hpce-dev / data /executor.py
์ด๋™ํ˜„
[DOCS] ๋ฐฑ์—”๋“œ ์ „๋ฐ˜ docstring ๋ณด๊ฐ• (Google ์Šคํƒ€์ผ Args/Returns/Raises)
58debbc
Raw
History Blame Contribute Delete
4.26 kB
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()