| from __future__ import annotations |
| """ |
| DuckDB Executor (์ฑ๊ธํค) |
| """ |
| import logging |
| import threading |
| from pathlib import Path |
| from typing import Any, Iterable |
|
|
| import duckdb |
| import pandas as pd |
|
|
| from config import settings |
|
|
| logger = logging.getLogger(__name__) |
|
|
| _SCHEMA_PATH = Path(__file__).parent / "schema.sql" |
|
|
|
|
| class DuckDBExecutor: |
| """DuckDB ๋จ์ผ connection์ ๋ณดํธํ๋ wrapper. |
| |
| ๋ฉํฐ์ค๋ ๋ ํ๊ฒฝ์์ RLock์ผ๋ก ์ ๊ทผ์ ์ง๋ ฌํํ์ฌ thread-safeํ๊ฒ ๋์ํ๋ค. |
| """ |
|
|
| def __init__(self, db_path: str) -> None: |
| """DuckDB ์ฐ๊ฒฐ์ ์ด๊ณ lock์ ์ด๊ธฐํํ๋ค. |
| |
| Args: |
| db_path: ์ฐ๊ฒฐํ DuckDB ํ์ผ ๊ฒฝ๋ก. |
| """ |
| self._db_path = db_path |
| self._conn = duckdb.connect(db_path) |
| self._lock = threading.RLock() |
| logger.info(f"DuckDB connected: {db_path}") |
|
|
| def execute(self, sql: str, params: Iterable[Any] | None = None) -> None: |
| """SQL 1๊ฑด์ ์คํํ๋ค (lock ๋ณดํธ). |
| |
| Args: |
| sql: ์คํํ SQL ๋ฌธ. |
| params: ๋ฐ์ธ๋ฉ ํ๋ผ๋ฏธํฐ (์์ผ๋ฉด None). |
| """ |
| with self._lock: |
| if params is None: |
| self._conn.execute(sql) |
| else: |
| self._conn.execute(sql, list(params)) |
|
|
| def executemany(self, sql: str, rows: list[list[Any]]) -> None: |
| """๋ค์ค ํ์ ์ผ๊ด ์คํํ๋ค (lock ๋ณดํธ). |
| |
| Args: |
| sql: ์คํํ SQL ๋ฌธ. |
| rows: ๊ฐ ํ์ ๋ฐ์ธ๋ฉ ํ๋ผ๋ฏธํฐ ๋ชฉ๋ก. |
| """ |
| with self._lock: |
| self._conn.executemany(sql, rows) |
|
|
| def to_pandas(self, sql: str, params: Iterable[Any] | None = None) -> pd.DataFrame: |
| """SELECT ๊ฒฐ๊ณผ๋ฅผ DataFrame์ผ๋ก ๋ฐํํ๋ค. |
| |
| Args: |
| sql: ์คํํ SELECT ๋ฌธ. |
| params: ๋ฐ์ธ๋ฉ ํ๋ผ๋ฏธํฐ (์์ผ๋ฉด None). |
| |
| Returns: |
| ์กฐํ ๊ฒฐ๊ณผ DataFrame. |
| """ |
| with self._lock: |
| if params is None: |
| return self._conn.execute(sql).df() |
| else: |
| return self._conn.execute(sql, list(params)).df() |
|
|
| def fetchone(self, sql: str, params: Iterable[Any] | None = None) -> tuple | None: |
| """์ฒซ ํ 1๊ฑด์ ๋ฐํํ๋ค. |
| |
| Args: |
| sql: ์คํํ SELECT ๋ฌธ. |
| params: ๋ฐ์ธ๋ฉ ํ๋ผ๋ฏธํฐ (์์ผ๋ฉด None). |
| |
| Returns: |
| ์ฒซ ํ tuple, ๊ฒฐ๊ณผ๊ฐ ์์ผ๋ฉด None. |
| """ |
| with self._lock: |
| if params is None: |
| return self._conn.execute(sql).fetchone() |
| return self._conn.execute(sql, list(params)).fetchone() |
|
|
| def fetchall(self, sql: str, params: Iterable[Any] | None = None) -> list[tuple]: |
| """์ ์ฒด ํ์ ๋ฐํํ๋ค. |
| |
| Args: |
| sql: ์คํํ SELECT ๋ฌธ. |
| params: ๋ฐ์ธ๋ฉ ํ๋ผ๋ฏธํฐ (์์ผ๋ฉด None). |
| |
| Returns: |
| ์ ์ฒด ํ tuple์ ๋ชฉ๋ก. |
| """ |
| with self._lock: |
| if params is None: |
| return self._conn.execute(sql).fetchall() |
| return self._conn.execute(sql, list(params)).fetchall() |
|
|
| def close(self) -> None: |
| """connection์ ์ข
๋ฃํ๋ค.""" |
| with self._lock: |
| self._conn.close() |
|
|
|
|
| _executor: DuckDBExecutor | None = None |
|
|
|
|
| def get_executor() -> DuckDBExecutor: |
| """DuckDB Executor ์ฑ๊ธํค์ ๋ฐํํ๋ค. |
| |
| ์ต์ด ํธ์ถ ์ settings.DB_PATH๋ก ์ธ์คํด์ค๋ฅผ ์์ฑํ์ฌ ์บ์ํ๋ค. |
| |
| Returns: |
| ํ๋ก์ธ์ค ์ ์ญ DuckDBExecutor ์ฑ๊ธํค. |
| """ |
| global _executor |
| if _executor is None: |
| _executor = DuckDBExecutor(settings.DB_PATH) |
| return _executor |
|
|
|
|
| def init_db() -> None: |
| """์คํค๋ง๋ฅผ ์์ฑํ๊ณ ์นดํ๋ก๊ทธ๋ฅผ ์๋ํ๋ค. |
| |
| schema.sql์ ์ฝ์ด statement ๋จ์๋ก ์คํํ ๋ค ์นดํ๋ก๊ทธ ์๋๋ฅผ ์ํํ๋ค. |
| DuckDB๋ multi-statement๋ฅผ ํ ๋ฒ์ ๋ฐ์ง ๋ชปํ๋ฏ๋ก ';' ๊ธฐ์ค์ผ๋ก ๋ถ๋ฆฌ ์คํํ๋ค. |
| """ |
| ex = get_executor() |
| with open(_SCHEMA_PATH, encoding="utf-8") as f: |
| schema_sql = f.read() |
| for stmt in [s.strip() for s in schema_sql.split(";") if s.strip()]: |
| ex.execute(stmt) |
|
|
| |
| from data.seed import seed_catalogs |
| seed_catalogs() |
|
|