File size: 4,261 Bytes
a28cff7
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
58debbc
 
 
a28cff7
 
4819352
58debbc
 
 
 
 
a28cff7
 
 
 
 
 
58debbc
 
 
 
 
 
a28cff7
 
 
 
 
 
 
58debbc
 
 
 
 
 
a28cff7
 
 
 
58debbc
 
 
 
 
 
 
 
 
a28cff7
 
 
 
 
 
 
58debbc
 
 
 
 
 
 
 
 
a28cff7
 
 
 
 
 
58debbc
 
 
 
 
 
 
 
 
a28cff7
 
 
 
 
 
58debbc
a28cff7
 
 
 
 
 
 
 
58debbc
 
 
 
 
 
 
a28cff7
 
 
 
 
 
 
58debbc
 
 
 
 
a28cff7
 
 
 
 
 
 
 
 
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
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()