""" Catalogo Parquet per Linker — sostituisce CUSTOM_ANAGR_ORIZ e CUSTOM_DIRITTI da Access. Usa DuckDB in-memory per query sui Parquet DataHub_v2. """ import datetime import logging from pathlib import Path from typing import List, Optional import duckdb import pandas as pd logger = logging.getLogger(__name__) class ParquetCatalog: """ Legge CUSTOM_ANAGR_ORIZ e CUSTOM_DIRITTI dai Parquet DataHub_v2 invece che da MS Access. """ def __init__(self, parquet_path: str): self.parquet_path = Path(parquet_path) self._con: Optional[duckdb.DuckDBPyConnection] = None def _get_con(self) -> duckdb.DuckDBPyConnection: if self._con is None: self._con = duckdb.connect() self._setup_views() return self._con def _setup_views(self): p = self.parquet_path prod_f = str(p / 'prodotti.parquet').replace('\\', '/') cast_f = str(p / 'cast.parquet').replace('\\', '/') imdb_f = str(p / 'imdb.parquet').replace('\\', '/') dir_f = str(p / 'diritti.parquet').replace('\\', '/') self._con.execute(f""" CREATE VIEW v_custom_anagr AS WITH prod AS ( SELECT * FROM read_parquet('{prod_f}') QUALIFY ROW_NUMBER() OVER (PARTITION BY prodotto ORDER BY edizione) = 1 ), imdb AS ( SELECT codice, MIN(riferimento_imdb) AS riferimento_imdb FROM read_parquet('{imdb_f}') GROUP BY codice ), reg AS ( SELECT prodotto, cognome AS REGISTA_COGN, nome AS REGISTA_NOME FROM read_parquet('{cast_f}') WHERE ruolo = 'FRE' QUALIFY ROW_NUMBER() OVER (PARTITION BY prodotto ORDER BY progr_cast) = 1 ), a1 AS ( SELECT prodotto, TRIM(COALESCE(nome, '') || ' ' || COALESCE(cognome, '')) AS ATTORE1 FROM read_parquet('{cast_f}') WHERE ruolo = 'C001' AND progr_cast = 1 ), a2 AS ( SELECT prodotto, TRIM(COALESCE(nome, '') || ' ' || COALESCE(cognome, '')) AS ATTORE2 FROM read_parquet('{cast_f}') WHERE ruolo = 'C001' AND progr_cast = 2 ), a3 AS ( SELECT prodotto, TRIM(COALESCE(nome, '') || ' ' || COALESCE(cognome, '')) AS ATTORE3 FROM read_parquet('{cast_f}') WHERE ruolo = 'C001' AND progr_cast = 3 ), distr AS ( SELECT prod AS prodotto, MIN(ragsoc_distr) AS DF_RAGSOC_DISTR FROM read_parquet('{dir_f}') GROUP BY prod ) SELECT CAST(p.prodotto AS VARCHAR) AS RIFER, p.tipologia AS TIPOL, p.titolo_originale AS "TO", p.titolo_italiano AS TI, p.anno_produzione AS ANNO, p.paesi_produzione1 AS PAESE, p.genere1 AS GENERE1, p.num_episodi AS EPIS, p.durata AS DUR, p.superserie_descr AS SUPERSERIE, CASE p.veg WHEN 'EVER GREEN' THEN 'EV.GR' ELSE p.veg END AS VEG_FREE, i.riferimento_imdb AS IMDB_CODICE, r.REGISTA_COGN, r.REGISTA_NOME, a1.ATTORE1, a2.ATTORE2, a3.ATTORE3, d.DF_RAGSOC_DISTR FROM prod p LEFT JOIN imdb i ON i.codice = p.prodotto LEFT JOIN reg r ON r.prodotto = p.prodotto LEFT JOIN a1 ON a1.prodotto = p.prodotto LEFT JOIN a2 ON a2.prodotto = p.prodotto LEFT JOIN a3 ON a3.prodotto = p.prodotto LEFT JOIN distr d ON d.prodotto = p.prodotto """) self._con.execute(f""" CREATE VIEW v_diritti AS SELECT CAST(prod AS VARCHAR) AS PROD, tipologia AS TIPOLOGIA, scad AS SCAD, decr AS DECR, ragsoc_distr AS DF_RAGSOC_DISTR FROM read_parquet('{dir_f}') """) imdb_full_f = str(p / 'imdb_full.parquet').replace('\\', '/') if (p / 'imdb_full.parquet').exists(): self._con.execute(f""" CREATE VIEW v_imdb_full AS SELECT *, DUR AS DURATA FROM read_parquet('{imdb_full_f}') """) logger.info(f"ParquetCatalog: VIEW pronte su {p}") def get_all_imdb_rifer_map(self) -> dict: """Dict IMDB→RIFER da scansione completa di imdb.parquet (no IN clause).""" try: imdb_f = str(self.parquet_path / 'imdb.parquet').replace('\\', '/') df = self._get_con().execute(f""" SELECT riferimento_imdb AS IMDB_CODICE, CAST(codice AS VARCHAR) AS RIFER FROM read_parquet('{imdb_f}') WHERE riferimento_imdb IS NOT NULL AND riferimento_imdb != '' QUALIFY ROW_NUMBER() OVER (PARTITION BY riferimento_imdb ORDER BY codice) = 1 """).df() if df.empty: return {} return dict(zip(df['IMDB_CODICE'].astype(str), df['RIFER'].astype(str))) except Exception as e: logger.error(f"get_all_imdb_rifer_map: {e}") return {} def get_all_emesso_agg(self) -> pd.DataFrame: """Scan completo di emesso_agg.parquet (L3) se disponibile — per modalità completo.""" try: agg = self._emesso_agg_path() if not agg: logger.warning("get_all_emesso_agg: emesso_agg.parquet non disponibile") return pd.DataFrame() agg_f = str(agg).replace('\\', '/') return self._get_con().execute(f""" SELECT RIFER, tipo_rete, rete_last AS rete, data_last AS data_emissione, ora_last AS ora_inizio, fascia_last AS fascia, aud_last AS audience, sha_last AS share FROM read_parquet('{agg_f}') """).df() except Exception as e: logger.error(f"get_all_emesso_agg: {e}") return pd.DataFrame() def get_all_custom(self) -> pd.DataFrame: try: custom_agg = self.parquet_path.parent / "parquet_agg" / "custom_agg.parquet" if custom_agg.exists(): agg_f = str(custom_agg).replace('\\', '/') df = self._get_con().execute(f"SELECT * FROM read_parquet('{agg_f}')").df() else: df = self._get_con().execute("SELECT * FROM v_custom_anagr").df() logger.info(f"ParquetCatalog CUSTOM_ANAGR_ORIZ: {len(df)} righe") return df except Exception as e: logger.error(f"ParquetCatalog get_all_custom: {e}") return pd.DataFrame() def search_imdb(self, term: str, by_regista: bool = False) -> pd.DataFrame: try: tables = self._get_con().execute("SHOW TABLES").df()['name'].tolist() if 'v_imdb_full' not in tables: return pd.DataFrame() safe = term.replace("'", "''") if by_regista: sql = f"SELECT * FROM v_imdb_full WHERE REGISTA ILIKE '%{safe}%' LIMIT 999999" else: sql = f"SELECT * FROM v_imdb_full WHERE \"TO\" ILIKE '{safe}%' OR TI ILIKE '{safe}%' LIMIT 999999" return self._get_con().execute(sql).df() except Exception as e: logger.error(f"ParquetCatalog search_imdb: {e}") return pd.DataFrame() def search_custom(self, term: str, by_regista: bool = False) -> pd.DataFrame: try: safe = term.replace("'", "''") if by_regista: sql = ( f"SELECT RIFER, TIPOL, \"TO\", TI, ANNO, REGISTA_COGN, REGISTA_NOME " f"FROM v_custom_anagr WHERE REGISTA_COGN ILIKE '%{safe}%' LIMIT 999999" ) else: sql = ( f"SELECT RIFER, TIPOL, \"TO\", TI, ANNO, REGISTA_COGN, REGISTA_NOME " f"FROM v_custom_anagr WHERE \"TO\" ILIKE '{safe}%' OR TI ILIKE '{safe}%' LIMIT 999999" ) return self._get_con().execute(sql).df() except Exception as e: logger.error(f"ParquetCatalog search_custom: {e}") return pd.DataFrame() def get_custom_by_rifer(self, rifer: str) -> pd.DataFrame: try: safe = str(rifer).replace("'", "''") return self._get_con().execute( f"SELECT * FROM v_custom_anagr WHERE RIFER = '{safe}'" ).df() except Exception as e: logger.error(f"get_custom_by_rifer: {e}") return pd.DataFrame() def get_custom_by_imdb(self, imdb_code: str) -> pd.DataFrame: try: safe = str(imdb_code).replace("'", "''") return self._get_con().execute( f"SELECT * FROM v_custom_anagr WHERE IMDB_CODICE = '{safe}'" ).df() except Exception as e: logger.error(f"get_custom_by_imdb: {e}") return pd.DataFrame() def get_imdb_by_code(self, imdb_code: str) -> pd.DataFrame: try: con = self._get_con() tables = con.execute("SHOW TABLES").df()['name'].tolist() if 'v_imdb_full' not in tables: return pd.DataFrame() safe = str(imdb_code).replace("'", "''") return con.execute( f"SELECT * FROM v_imdb_full WHERE IMDB_CODICE = '{safe}'" ).df() except Exception as e: logger.error(f"get_imdb_by_code: {e}") return pd.DataFrame() def get_custom_by_rifers(self, rifers: list) -> pd.DataFrame: if not rifers: return pd.DataFrame() try: placeholders = ', '.join( f"'{str(r).replace(chr(39), chr(39)*2)}'" for r in rifers ) return self._get_con().execute( f"SELECT * FROM v_custom_anagr WHERE RIFER IN ({placeholders})" ).df() except Exception as e: logger.error(f"get_custom_by_rifers: {e}") return pd.DataFrame() def get_all_imdb(self) -> pd.DataFrame: try: con = self._get_con() # v_imdb_full esiste solo se il file è presente tables = con.execute("SHOW TABLES").df()['name'].tolist() if 'v_imdb_full' not in tables: return pd.DataFrame() df = con.execute("SELECT * FROM v_imdb_full").df() logger.info(f"ParquetCatalog IMDB: {len(df)} righe") return df except Exception as e: logger.error(f"ParquetCatalog get_all_imdb: {e}") return pd.DataFrame() def get_distributori(self) -> List[str]: try: dir_f = str(self.parquet_path / 'diritti.parquet').replace('\\', '/') df = self._get_con().execute(f""" SELECT DISTINCT ragsoc_distr FROM read_parquet('{dir_f}') WHERE ragsoc_distr IS NOT NULL ORDER BY ragsoc_distr """).df() return [str(v).strip() for v in df['ragsoc_distr'].tolist() if str(v).strip()] except Exception as e: logger.error(f"ParquetCatalog get_distributori: {e}") return [] def get_tipologie_diritti(self) -> List[str]: try: anno = datetime.date.today().year df = self._get_con().execute( "SELECT DISTINCT TIPOLOGIA FROM v_diritti " f"WHERE SCAD >= '{anno - 1}-01-01' AND SCAD <= '{anno + 5}-12-31' " "AND TIPOLOGIA IS NOT NULL ORDER BY TIPOLOGIA" ).df() return df['TIPOLOGIA'].dropna().tolist() except Exception as e: logger.error(f"ParquetCatalog get_tipologie_diritti: {e}") return [] def get_diritti_rifer(self, tipologie: Optional[List[str]] = None) -> List[str]: try: anno = datetime.date.today().year if tipologie: placeholders = ', '.join(f"'{t.replace(chr(39), chr(39)*2)}'" for t in tipologie) tipol_clause = f"AND TIPOLOGIA IN ({placeholders})" else: tipol_clause = ( "AND TIPOLOGIA IN " "('FILM','TV MOVIE','TELEFILM','MINISERIE','SIT COM','SOAP','TELENOVELAS')" ) sql = ( "SELECT DISTINCT PROD FROM v_diritti " f"WHERE SCAD >= '{anno - 1}-01-01' " f" AND SCAD <= '{anno + 5}-12-31' " f" {tipol_clause}" ) df = self._get_con().execute(sql).df() return [str(v).strip() for v in df['PROD'].dropna().tolist() if str(v).strip()] except Exception as e: logger.error(f"ParquetCatalog get_diritti_rifer: {e}") return [] def get_diritti_full_by_rifer(self, rifer: str) -> pd.DataFrame: try: safe = str(rifer).replace("'", "''") dir_f = str(self.parquet_path / 'diritti.parquet').replace('\\', '/') return self._get_con().execute(f""" SELECT CAST(prod AS VARCHAR) AS RIFER, decr, scad, ragsoc_distr, perc, pass_cons_tot, pass_eff_tot, causale FROM read_parquet('{dir_f}') WHERE CAST(prod AS VARCHAR) = '{safe}' """).df() except Exception as e: logger.error(f"get_diritti_full_by_rifer: {e}") return pd.DataFrame() def get_diritti_full_by_rifers(self, rifers: list) -> pd.DataFrame: if not rifers: return pd.DataFrame() try: placeholders = ', '.join( f"'{str(r).replace(chr(39), chr(39)*2)}'" for r in rifers ) dir_f = str(self.parquet_path / 'diritti.parquet').replace('\\', '/') return self._get_con().execute(f""" SELECT CAST(prod AS VARCHAR) AS RIFER, decr, scad, ragsoc_distr, perc, pass_cons_tot, pass_eff_tot, causale FROM read_parquet('{dir_f}') WHERE CAST(prod AS VARCHAR) IN ({placeholders}) """).df() except Exception as e: logger.error(f"get_diritti_full_by_rifers: {e}") return pd.DataFrame() _GENERALISTE = ('N1', 'N2', 'N3', 'C5', 'I1', 'MC', 'R4') def _emesso_agg_path(self) -> Path | None: """Ritorna il path di emesso_agg.parquet se disponibile, altrimenti None.""" p = self.parquet_path.parent / "parquet_agg" / "emesso_agg.parquet" return p if p.exists() else None def get_emesso_both_last_by_rifers(self, rifers: list) -> pd.DataFrame: """Ultima emissione generaliste E tematiche. Usa emesso_agg.parquet (L3) se disponibile — altrimenti scansiona emesso.parquet (L2).""" if not rifers: return pd.DataFrame() try: placeholders = ', '.join( f"'{str(r).replace(chr(39), chr(39)*2)}'" for r in rifers ) agg = self._emesso_agg_path() if agg: agg_f = str(agg).replace('\\', '/') return self._get_con().execute(f""" SELECT RIFER, tipo_rete, rete_last AS rete, data_last AS data_emissione, ora_last AS ora_inizio, fascia_last AS fascia, aud_last AS audience, sha_last AS share FROM read_parquet('{agg_f}') WHERE RIFER IN ({placeholders}) """).df() nets = ', '.join(f"'{n}'" for n in self._GENERALISTE) em_f = str(self.parquet_path / 'emesso.parquet').replace('\\', '/') return self._get_con().execute(f""" SELECT CAST(prodotto AS VARCHAR) AS RIFER, rete, data_emissione, ora_inizio, fascia, audience, share, CASE WHEN rete IN ({nets}) THEN 'G' ELSE 'T' END AS tipo_rete FROM read_parquet('{em_f}') WHERE CAST(prodotto AS VARCHAR) IN ({placeholders}) QUALIFY ROW_NUMBER() OVER ( PARTITION BY CAST(prodotto AS VARCHAR), CASE WHEN rete IN ({nets}) THEN 'G' ELSE 'T' END ORDER BY data_emissione DESC, ora_inizio DESC ) = 1 """).df() except Exception as e: logger.error(f"get_emesso_both_last_by_rifers: {e}") return pd.DataFrame() def get_emesso_both_first_by_rifers(self, rifers: list) -> pd.DataFrame: """Prima emissione generaliste E tematiche. Usa emesso_agg.parquet (L3) se disponibile — altrimenti scansiona emesso.parquet (L2).""" if not rifers: return pd.DataFrame() try: placeholders = ', '.join( f"'{str(r).replace(chr(39), chr(39)*2)}'" for r in rifers ) agg = self._emesso_agg_path() if agg: agg_f = str(agg).replace('\\', '/') return self._get_con().execute(f""" SELECT RIFER, tipo_rete, rete_first AS rete, data_first AS data_emissione, ora_first AS ora_inizio, fascia_first AS fascia, aud_first AS audience, sha_first AS share FROM read_parquet('{agg_f}') WHERE RIFER IN ({placeholders}) """).df() nets = ', '.join(f"'{n}'" for n in self._GENERALISTE) em_f = str(self.parquet_path / 'emesso.parquet').replace('\\', '/') return self._get_con().execute(f""" SELECT CAST(prodotto AS VARCHAR) AS RIFER, rete, data_emissione, ora_inizio, fascia, audience, share, CASE WHEN rete IN ({nets}) THEN 'G' ELSE 'T' END AS tipo_rete FROM read_parquet('{em_f}') WHERE CAST(prodotto AS VARCHAR) IN ({placeholders}) QUALIFY ROW_NUMBER() OVER ( PARTITION BY CAST(prodotto AS VARCHAR), CASE WHEN rete IN ({nets}) THEN 'G' ELSE 'T' END ORDER BY data_emissione ASC, ora_inizio ASC ) = 1 """).df() except Exception as e: logger.error(f"get_emesso_both_first_by_rifers: {e}") return pd.DataFrame() def get_emesso_last_by_rifers(self, rifers: list) -> pd.DataFrame: """Compatibilità: usa get_emesso_both_last_by_rifers filtrando su G.""" df = self.get_emesso_both_last_by_rifers(rifers) if df.empty: return df return df[df['tipo_rete'] == 'G'].drop(columns=['tipo_rete']) def get_emesso_tematiche_last_by_rifers(self, rifers: list) -> pd.DataFrame: """Compatibilità: usa get_emesso_both_last_by_rifers filtrando su T.""" df = self.get_emesso_both_last_by_rifers(rifers) if df.empty: return df return df[df['tipo_rete'] == 'T'].drop(columns=['tipo_rete']) def get_gemma_valutazioni_by_imdb_codes(self, imdb_codes: list) -> pd.DataFrame: """Restituisce una riga per IMDB code con i campi valutazioni da gemma_unified.db.""" if not imdb_codes: return pd.DataFrame() gemma_unified_db = self.parquet_path / 'gemma_unified.db' if gemma_unified_db.exists(): try: import sqlite3 as _sl placeholders = ','.join(['?'] * len(imdb_codes)) con = _sl.connect(str(gemma_unified_db)) con.row_factory = _sl.Row rows = con.execute( f"SELECT A_COD_IMDB, A_DISTRIBUTORE, V_RDA, " f"V_C5, V_I1, V_R4, V_LA5, V_I2, V_IRIS, " f"V_TOP, V_FOC, V_C20, V_CI34, V_C27, " f"A_SUPPORTO, A_DATA, GEM_SIN " f"FROM valutazioni_unified " f"WHERE A_COD_IMDB IN ({placeholders})", imdb_codes ).fetchall() con.close() return pd.DataFrame([dict(r) for r in rows]) except Exception as e: logger.error(f"get_gemma_valutazioni_by_imdb_codes (unified): {e}", exc_info=True) # Fallback: DuckDB su gemma.parquet (gemma_unified.db non ancora generato) gemma_f = str(self.parquet_path / 'gemma.parquet').replace('\\', '/') placeholders_sql = ', '.join( f"'{str(c).replace(chr(39), chr(39)*2)}'" for c in imdb_codes ) try: _win = "PARTITION BY A_COD_IMDB ORDER BY A_DATA DESC NULLS LAST ROWS BETWEEN UNBOUNDED PRECEDING AND UNBOUNDED FOLLOWING" df = self._get_con().execute(f""" SELECT DISTINCT A_COD_IMDB, FIRST_VALUE(A_DISTRIBUTORE IGNORE NULLS) OVER ({_win}) AS A_DISTRIBUTORE, FIRST_VALUE(V_RDA IGNORE NULLS) OVER ({_win}) AS V_RDA, FIRST_VALUE(V_C5 IGNORE NULLS) OVER ({_win}) AS V_C5, FIRST_VALUE(V_I1 IGNORE NULLS) OVER ({_win}) AS V_I1, FIRST_VALUE(V_R4 IGNORE NULLS) OVER ({_win}) AS V_R4, FIRST_VALUE(V_LA5 IGNORE NULLS) OVER ({_win}) AS V_LA5, FIRST_VALUE(V_I2 IGNORE NULLS) OVER ({_win}) AS V_I2, FIRST_VALUE(V_IRIS IGNORE NULLS) OVER ({_win}) AS V_IRIS, FIRST_VALUE(V_TOP IGNORE NULLS) OVER ({_win}) AS V_TOP, FIRST_VALUE(V_FOC IGNORE NULLS) OVER ({_win}) AS V_FOC, FIRST_VALUE(V_C20 IGNORE NULLS) OVER ({_win}) AS V_C20, FIRST_VALUE(V_CI34 IGNORE NULLS) OVER ({_win}) AS V_CI34, FIRST_VALUE(V_C27 IGNORE NULLS) OVER ({_win}) AS V_C27, FIRST_VALUE(A_SUPPORTO IGNORE NULLS) OVER ({_win}) AS A_SUPPORTO, FIRST_VALUE(A_DATA IGNORE NULLS) OVER ({_win}) AS A_DATA FROM read_parquet('{gemma_f}') WHERE A_COD_IMDB IN ({placeholders_sql}) """).df() df['GEM_SIN'] = 'G' return df except Exception as e: logger.error(f"get_gemma_valutazioni_by_imdb_codes: {e}", exc_info=True) return pd.DataFrame() def get_all_gemma(self) -> pd.DataFrame: """Restituisce tutti i record di gemma.parquet ordinati per A_DATA DESC.""" gemma_f = str(self.parquet_path / 'gemma.parquet').replace('\\', '/') try: return self._get_con().execute( f"SELECT * FROM read_parquet('{gemma_f}') ORDER BY A_DATA DESC NULLS LAST" ).df() except Exception as e: logger.error(f"get_all_gemma: {e}", exc_info=True) return pd.DataFrame() def get_rifer_by_imdb_codes(self, imdb_codes: list) -> pd.DataFrame: """Dato un elenco di codici IMDB (tt...) restituisce IMDB_CODICE → RIFER OnAir.""" if not imdb_codes: return pd.DataFrame() imdb_f = str(self.parquet_path / 'imdb.parquet').replace('\\', '/') placeholders = ', '.join( f"'{str(c).replace(chr(39), chr(39)*2)}'" for c in imdb_codes ) try: return self._get_con().execute(f""" SELECT riferimento_imdb AS IMDB_CODICE, CAST(codice AS VARCHAR) AS RIFER, COUNT(*) OVER (PARTITION BY riferimento_imdb) AS n_rifers FROM read_parquet('{imdb_f}') WHERE riferimento_imdb IN ({placeholders}) QUALIFY ROW_NUMBER() OVER ( PARTITION BY riferimento_imdb ORDER BY codice ) = 1 """).df() except Exception as e: logger.error(f"get_rifer_by_imdb_codes: {e}", exc_info=True) return pd.DataFrame() def get_imdb_data_by_codes(self, imdb_codes: list) -> pd.DataFrame: """Restituisce IMDB_CODICE, VOTO, VOTANTI per lista di codici IMDB.""" if not imdb_codes: return pd.DataFrame() imdb_f = str(self.parquet_path / 'imdb_full.parquet').replace('\\', '/') placeholders = ', '.join( f"'{str(c).replace(chr(39), chr(39)*2)}'" for c in imdb_codes ) try: return self._get_con().execute(f""" SELECT * FROM read_parquet('{imdb_f}') WHERE IMDB_CODICE IN ({placeholders}) """).df() except Exception as e: logger.error(f"get_imdb_data_by_codes: {e}", exc_info=True) return pd.DataFrame() def get_scelte_rete_by_rifers(self, rifers: list) -> pd.DataFrame: """Restituisce rete/slot aggregati per lista RIFER (una riga per RIFER).""" if not rifers: return pd.DataFrame() sr_path = str(self.parquet_path / 'scelte_rete.parquet').replace('\\', '/') placeholders = ', '.join( f"'{str(r).replace(chr(39), chr(39)*2)}'" for r in rifers ) try: df = self._get_con().execute(f""" WITH base AS ( SELECT DISTINCT CAST(prodotto AS VARCHAR) AS RIFER, codice_rete, COALESCE(slot, '') AS slot, CASE codice_rete WHEN 'C5' THEN 0 WHEN 'I1' THEN 1 WHEN 'R4' THEN 2 ELSE 3 END AS prio FROM read_parquet('{sr_path}') WHERE CAST(prodotto AS VARCHAR) IN ({placeholders}) ) SELECT RIFER, STRING_AGG(codice_rete || ' ' || slot, ', ' ORDER BY prio, codice_rete) AS rete_slot FROM base GROUP BY RIFER """).df() return df except Exception as e: logger.error(f"get_scelte_rete_by_rifers: {e}", exc_info=True) return pd.DataFrame() def get_boxoffice_by_rifers(self, rifers: list) -> pd.DataFrame: """Restituisce dati box office (debutto, incasso, spettatori, distributore) per lista RIFER.""" if not rifers: return pd.DataFrame() bo_path = str(self.parquet_path / 'boxoffice.parquet') placeholders = ', '.join( f"'{str(r).replace(chr(39), chr(39)*2)}'" for r in rifers ) try: con = self._get_con() df = con.execute(f""" SELECT CAST(prodotto AS VARCHAR) AS RIFER, data_debutto, incasso, spettatori, distributore FROM read_parquet('{bo_path}') WHERE CAST(prodotto AS VARCHAR) IN ({placeholders}) """).df() return df except Exception as e: logger.error(f"get_boxoffice_by_rifers: {e}", exc_info=True) return pd.DataFrame() def get_emissioni_generaliste_by_rifer(self, rifer: str) -> pd.DataFrame: """Restituisce tutte le emissioni generaliste di un prodotto ordinate per data.""" if not rifer: return pd.DataFrame() em_f = str(self.parquet_path / 'emesso.parquet').replace('\\', '/') safe = rifer.replace("'", "''") nets = "('N1','N2','N3','C5','I1','MC','R4')" try: return self._get_con().execute(f""" SELECT rete, SUBSTR(data_emissione, 9, 2) || '/' || SUBSTR(data_emissione, 6, 2) || '/' || SUBSTR(data_emissione, 1, 4) AS data, LPAD(CAST(ora_inizio // 10000 AS VARCHAR), 2, '0') || ':' || LPAD(CAST((ora_inizio % 10000) // 100 AS VARCHAR), 2, '0') AS hi, CASE WHEN prima_visione = 'S' THEN '1TV' ELSE '' END AS itv, fascia, audience AS aud, ROUND(share, 2) AS sha FROM read_parquet('{em_f}') WHERE CAST(prodotto AS VARCHAR) = '{safe}' AND rete IN {nets} ORDER BY data_emissione, ora_inizio """).df() except Exception as e: logger.error(f"get_emissioni_generaliste_by_rifer: {e}", exc_info=True) return pd.DataFrame() def get_emissioni_tematiche_by_rifer(self, rifer: str) -> pd.DataFrame: """Restituisce tutte le emissioni tematiche di un prodotto ordinate per data.""" if not rifer: return pd.DataFrame() em_f = str(self.parquet_path / 'emesso.parquet').replace('\\', '/') safe = rifer.replace("'", "''") nets = "('N1','N2','N3','C5','I1','MC','R4')" try: return self._get_con().execute(f""" SELECT rete, SUBSTR(data_emissione, 9, 2) || '/' || SUBSTR(data_emissione, 6, 2) || '/' || SUBSTR(data_emissione, 1, 4) AS data, LPAD(CAST(ora_inizio // 10000 AS VARCHAR), 2, '0') || ':' || LPAD(CAST((ora_inizio % 10000) // 100 AS VARCHAR), 2, '0') AS hi, CASE WHEN prima_visione = 'S' THEN '1TV' ELSE '' END AS itv, fascia, audience AS aud, ROUND(share, 2) AS sha FROM read_parquet('{em_f}') WHERE CAST(prodotto AS VARCHAR) = '{safe}' AND rete NOT IN {nets} ORDER BY data_emissione, ora_inizio """).df() except Exception as e: logger.error(f"get_emissioni_tematiche_by_rifer: {e}", exc_info=True) return pd.DataFrame() def get_edizioni_by_rifer(self, rifer: str) -> pd.DataFrame: """Restituisce tutte le edizioni di un prodotto (ED, TIPOL, DUR, VM).""" if not rifer: return pd.DataFrame() prod_f = str(self.parquet_path / 'prodotti.parquet').replace('\\', '/') safe = rifer.replace("'", "''") try: return self._get_con().execute(f""" SELECT edizione AS ED, tipologia AS TIPOL, durata AS DUR, vm AS VM FROM read_parquet('{prod_f}') WHERE CAST(prodotto AS VARCHAR) = '{safe}' ORDER BY edizione """).df() except Exception as e: logger.error(f"get_edizioni_by_rifer: {e}", exc_info=True) return pd.DataFrame() def get_prodotti_by_superserie(self, superserie_name: str) -> pd.DataFrame: """Restituisce tutti i prodotti di una superserie (VEG, TIPOL, TI).""" if not superserie_name: return pd.DataFrame() prod_f = str(self.parquet_path / 'prodotti.parquet').replace('\\', '/') safe = superserie_name.replace("'", "''") try: return self._get_con().execute(f""" SELECT CASE veg WHEN 'EVER GREEN' THEN 'EV.GR' ELSE veg END AS VEG, tipologia AS TIPOL, titolo_italiano AS TI FROM read_parquet('{prod_f}') WHERE superserie_descr = '{safe}' QUALIFY ROW_NUMBER() OVER (PARTITION BY prodotto ORDER BY edizione) = 1 ORDER BY tipologia, titolo_italiano """).df() except Exception as e: logger.error(f"get_prodotti_by_superserie: {e}", exc_info=True) return pd.DataFrame() # ── OTT (Parquet) + OTT_EXT (SQLite in-memory cache) ──────────────────── _SVOD_PROVIDERS = ('netflix', 'amazon prime video', 'apple tv plus', 'apple tv+') @property def ott_parquet_path(self) -> Path: return self.parquet_path / 'ott.parquet' def register_ott_ext(self, sqlite_con) -> int: """Carica OTT_EXT da SQLite in una tabella DuckDB in-memory. Chiamato al boot del Linker e dopo ogni scrittura su OTT_EXT.""" try: df = pd.read_sql( "SELECT mediaset_id, provider, data_creazione, " "segnalazioni_automatiche, imdb_override FROM OTT_EXT", sqlite_con, ) con = self._get_con() con.register('_ott_ext_tmp', df) con.execute("CREATE OR REPLACE TABLE ott_ext AS SELECT * FROM _ott_ext_tmp") con.unregister('_ott_ext_tmp') logger.info(f"ParquetCatalog: ott_ext caricato ({len(df)} righe)") return len(df) except Exception as e: logger.error(f"register_ott_ext: {e}") return 0 def refresh_ott_ext(self, sqlite_con) -> int: """Ricarica ott_ext dopo scritture su OTT_EXT (imdb_override, segnalazioni).""" return self.register_ott_ext(sqlite_con) def get_ott_joined_for_mode(self, mode: str, anno_min: int = 0, cutoff: str = None) -> pd.DataFrame: """JOIN DuckDB tra ott.parquet e ott_ext (in-memory) per la modalità richiesta.""" ott_f = str(self.ott_parquet_path).replace('\\', '/') if not self.ott_parquet_path.exists(): logger.warning("get_ott_joined_for_mode: ott.parquet non trovato") return pd.DataFrame() tipo = 'movie' if 'movie' in mode else 'show' p = ', '.join(f"'{v}'" for v in self._SVOD_PROVIDERS) p = f"({p})" def _date_filter(cutoff): if not cutoff: return '' if isinstance(cutoff, list): quoted = ', '.join(f"'{d}'" for d in cutoff if d) return f"AND e.data_creazione IN ({quoted})" if quoted else '' return f"AND e.data_creazione >= '{cutoff}'" if mode == 'movie_filtrato': where = (f"LOWER(o.tipo)='{tipo}' AND UPPER(o.tipo_finestra)='SVOD' " f"AND LOWER(o.provider) IN {p} " f"AND TRY_CAST(NULLIF(TRIM(o.anno), '') AS INTEGER) >= {anno_min}") join_type = 'INNER' extra_join = _date_filter(cutoff) order_by = 'e.data_creazione DESC, CAST(o.mediaset_id AS VARCHAR) ASC' limit_sql = '' elif mode == 'show_filtrato': where = (f"LOWER(o.tipo)='{tipo}' AND UPPER(o.tipo_finestra)='SVOD' " f"AND LOWER(o.provider) IN {p} " "AND o.data_inizio >= '2019-01-01' " "AND (o.paesi LIKE '%IT%' OR o.paesi='-' OR o.paesi='[XX]')") join_type = 'INNER' extra_join = _date_filter(cutoff) order_by = 'e.data_creazione DESC, CAST(o.mediaset_id AS VARCHAR) ASC' limit_sql = '' else: # completo where = f"LOWER(o.tipo)='{tipo}' AND UPPER(o.tipo_finestra)='SVOD'" join_type = 'LEFT' extra_join = '' order_by = 'o.data_inizio DESC' limit_sql = '' # Per filtrato: deduplica per (mediaset_id, provider) tenendo la finestra più recente dedup = ("QUALIFY ROW_NUMBER() OVER " "(PARTITION BY CAST(o.mediaset_id AS VARCHAR), LOWER(CAST(o.provider AS VARCHAR)) " "ORDER BY o.data_inizio DESC) = 1" ) if mode.endswith('_filtrato') else '' sql = f""" SELECT CAST(o.mediaset_id AS VARCHAR) AS mediaset_id, CAST(o.titolo AS VARCHAR) AS titolo, CAST(o.tipo AS VARCHAR) AS tipo, CASE WHEN TRIM(COALESCE(CAST(o.imdb_id AS VARCHAR), '')) = '' THEN COALESCE(CAST(e.imdb_override AS VARCHAR), '') ELSE CAST(o.imdb_id AS VARCHAR) END AS imdb_id, CAST(o.anno AS VARCHAR) AS anno, CAST(o.nr_stagioni AS VARCHAR) AS nr_stagioni, CAST(o.tot_episodi AS VARCHAR) AS tot_episodi, CAST(o.provider AS VARCHAR) AS provider, CAST(o.tipo_finestra AS VARCHAR) AS tipo_finestra, CAST(o.monetization_type AS VARCHAR) AS monetization_type, CAST(o.data_inizio AS VARCHAR) AS data_inizio, CAST(o.data_fine AS VARCHAR) AS data_fine, CAST(o.regista AS VARCHAR) AS regista, CAST(o.generi AS VARCHAR) AS generi, CAST(o.paesi AS VARCHAR) AS paesi, COALESCE(CAST(e.data_creazione AS VARCHAR), '') AS data_creazione, COALESCE(CAST(e.segnalazioni_automatiche AS VARCHAR), '') AS segnalazioni_automatiche, COALESCE(CAST(e.imdb_override AS VARCHAR), '') AS imdb_override FROM read_parquet('{ott_f}') o {join_type} JOIN ott_ext e ON CAST(o.mediaset_id AS VARCHAR) = e.mediaset_id AND LOWER(CAST(o.provider AS VARCHAR)) = LOWER(e.provider) {extra_join} WHERE {where} {dedup} ORDER BY {order_by} {limit_sql} """ try: return self._get_con().execute(sql).df() except Exception as e: logger.error(f"get_ott_joined_for_mode: {e}") return pd.DataFrame() def close(self): if self._con: try: self._con.close() except Exception: pass self._con = None