# app/database/connection.py """ Gestore connessione dati — DuckDB in-memory con VIEW sui Parquet DataHub. db_user.db (SQLite) mantenuto separatamente per scritture utente. """ import sqlite3 import threading import duckdb import pandas as pd from pathlib import Path from loguru import logger from app.config_loader import get_app_root, load_config class DatabaseManager: _instance = None _lock = threading.Lock() def __new__(cls): if cls._instance is None: with cls._lock: if cls._instance is None: cls._instance = super().__new__(cls) cls._instance._initialized = False return cls._instance def __init__(self): if self._initialized: return config = load_config() app_root = get_app_root() parquet_dir = app_root / config.get('parquet_dir', 'local_db/parquet') user_db_path = app_root / config['databases']['user'] # Connessione DuckDB in-memory — VIEW su ogni Parquet self.conn = duckdb.connect() self._register_parquet_views(parquet_dir) # Connessione SQLite separata per db_user (scritture) user_db_path.parent.mkdir(parents=True, exist_ok=True) self._user_conn = sqlite3.connect(str(user_db_path), check_same_thread=False) logger.info(f"DatabaseManager: {len(list(parquet_dir.glob('*.parquet')))} Parquet registrati da {parquet_dir}") self._initialized = True def _register_parquet_views(self, parquet_dir: Path): """Crea una VIEW DuckDB per ogni file .parquet nella directory.""" if not parquet_dir.exists(): logger.warning(f"Parquet dir non trovata: {parquet_dir}") return for f in sorted(parquet_dir.glob("*.parquet")): view_name = f.stem self.conn.execute(f'CREATE VIEW "{view_name}" AS SELECT * FROM \'{f}\'') logger.debug(f" VIEW: {view_name} → {f.name}") def execute_query(self, query: str, params=None) -> pd.DataFrame: """Esegue una query SELECT su DuckDB e restituisce un DataFrame.""" try: if params: return self.conn.execute(query, params).df() return self.conn.execute(query).df() except Exception as e: logger.error(f"Errore query DuckDB: {e}\nQuery: {query}") raise def execute_write(self, query: str, params: tuple = None): """Esegue una scrittura su db_user.db (SQLite).""" try: cursor = self._user_conn.cursor() cursor.execute(query, params or ()) self._user_conn.commit() return cursor.lastrowid except Exception as e: self._user_conn.rollback() logger.error(f"Errore write SQLite: {e}\nQuery: {query}") raise def close(self): if self.conn: self.conn.close() self.conn = None if self._user_conn: self._user_conn.close() self._user_conn = None logger.info("DatabaseManager chiuso.") @classmethod def get_instance(cls): return cls() try: db_manager = DatabaseManager() except Exception as e: db_manager = None logger.critical(f"DatabaseManager non inizializzato: {e}")