""" emesso_enrich.py — produce emesso_enrich.parquet (sidecar) Join key con emesso.parquet: (rete, data_emissione, ora_inizio) Colonne derivate: - fornitore_diritto / fornitore_cluster : join emesso → diritti Free/Analogico → fornitori_cluster_mapping - is_primetime_main : programma con maggiore copertura 21:10-24:00 per rete/data - prima_visione_recalc : logica a cascata FILM/anno/ncparpgm/prima emissione assoluta - timestamp_normalizzato : datetime normalizzato da data_emissione + ora_inizio Dipende da: emesso.parquet, diritti.parquet (già prodotto), prodotti.parquet """ import duckdb from config.settings import out_parquet, FORNITORE_CLUSTER_FILE def run(con: duckdb.DuckDBPyConnection) -> int: forn_file = FORNITORE_CLUSTER_FILE.replace('\\', '/') em_pq = out_parquet('emesso').replace('\\', '/') dir_pq = out_parquet('diritti').replace('\\', '/') prod_pq = out_parquet('prodotti').replace('\\', '/') out = out_parquet('emesso_enrich').replace('\\', '/') # ------------------------------------------------------------------ # Step 1: diritti filtrati (Free/Analogico) per il join fornitore # Legge diritti.parquet — già nel bidone, colonne già tipizzate # ------------------------------------------------------------------ con.execute(f""" CREATE OR REPLACE TABLE _dir_forn AS SELECT prod, TRY_STRPTIME(decr, '%Y-%m-%d')::DATE AS decr_date, TRY_STRPTIME(scad, '%Y-%m-%d')::DATE AS scad_date, ragsoc_distr, contratto, riga, situazione FROM '{dir_pq}' WHERE UPPER(tipo_diritto) = 'FREE' AND UPPER(piattaforma) = 'ANALOGICO' AND prod IS NOT NULL AND decr IS NOT NULL AND scad IS NOT NULL """) # ------------------------------------------------------------------ # Step 2: fornitore_cluster mapping # ------------------------------------------------------------------ con.execute(f""" CREATE OR REPLACE TABLE _forn_map AS SELECT UPPER(TRIM(fornitore_diritto)) AS fornitore_key, fornitore_cluster FROM read_csv('{forn_file}', delim=';', encoding='cp1252', header=true) """) # ------------------------------------------------------------------ # Step 3: anno_produzione da prodotti # ------------------------------------------------------------------ con.execute(f""" CREATE OR REPLACE TABLE _prod_anno AS SELECT prodotto, MAX(anno_produzione) AS anno_produzione FROM '{prod_pq}' WHERE prodotto IS NOT NULL GROUP BY prodotto """) # ------------------------------------------------------------------ # Step 4: join emesso → diritti → fornitore_cluster # ------------------------------------------------------------------ con.execute(f""" CREATE OR REPLACE TABLE _em_forn AS WITH joined AS ( SELECT e.rete, e.data_emissione, e.ora_inizio, e.prodotto, e.tipologia, e.ncparpgm, e.prima_visione, e.fascia, e.ora_fine, TRY_CAST(e.data_emissione AS DATE) AS data_emissione_date, d.ragsoc_distr AS fornitore_diritto, ROW_NUMBER() OVER ( PARTITION BY e.rete, e.data_emissione, e.ora_inizio ORDER BY d.contratto NULLS LAST, d.riga NULLS LAST, d.situazione NULLS LAST ) AS _rn FROM '{em_pq}' e LEFT JOIN _dir_forn d ON e.prodotto = d.prod AND TRY_CAST(e.data_emissione AS DATE) BETWEEN d.decr_date AND d.scad_date ) SELECT j.rete, j.data_emissione, j.ora_inizio, j.prodotto, j.tipologia, j.ncparpgm, j.prima_visione, j.fascia, j.ora_fine, j.data_emissione_date, j.fornitore_diritto, CASE WHEN j.fornitore_diritto IS NOT NULL THEN COALESCE(fm.fornitore_cluster, 'altro') ELSE NULL END AS fornitore_cluster FROM joined j LEFT JOIN _forn_map fm ON UPPER(TRIM(j.fornitore_diritto)) = fm.fornitore_key WHERE j._rn = 1 """) # ------------------------------------------------------------------ # Step 5: is_primetime_main # Primetime = 21:10–24:00 → interi 211000–240000 # ------------------------------------------------------------------ con.execute(""" CREATE OR REPLACE TABLE _em_pt AS WITH coverage AS ( SELECT *, CASE WHEN fascia = 'PR' THEN GREATEST(0, LEAST( CASE WHEN ora_fine < ora_inizio THEN ora_fine + 240000 ELSE ora_fine END, 240000 ) - GREATEST(ora_inizio, 211000) ) ELSE NULL END AS _pt_cov FROM _em_forn ), ranked AS ( SELECT *, CASE WHEN fascia = 'PR' THEN ROW_NUMBER() OVER ( PARTITION BY rete, data_emissione ORDER BY _pt_cov DESC, ora_inizio ASC ) ELSE NULL END AS _pt_rank FROM coverage ) SELECT * EXCLUDE (_pt_cov, _pt_rank), CASE WHEN fascia != 'PR' THEN 0 WHEN _pt_rank = 1 THEN 1 ELSE 0 END AS is_primetime_main FROM ranked """) # ------------------------------------------------------------------ # Step 6: prima_visione_recalc + timestamp_normalizzato # ------------------------------------------------------------------ con.execute(""" CREATE OR REPLACE TABLE _em_pv AS WITH em_joined AS ( SELECT e.*, p.anno_produzione, (TRY_STRPTIME(e.data_emissione, '%Y-%m-%d')::TIMESTAMP + (e.ora_inizio / 10000) * INTERVAL '1' HOUR + ((e.ora_inizio % 10000) / 100) * INTERVAL '1' MINUTE + (e.ora_inizio % 100) * INTERVAL '1' SECOND ) AS _ts FROM _em_pt e LEFT JOIN _prod_anno p USING (prodotto) ), with_rn AS ( SELECT *, ROW_NUMBER() OVER ( PARTITION BY prodotto ORDER BY _ts ASC NULLS LAST ) AS _rn_ts FROM em_joined ) SELECT rete, data_emissione, ora_inizio, fornitore_diritto, fornitore_cluster, is_primetime_main, STRFTIME(_ts, '%Y-%m-%d %H:%M:%S') AS timestamp_normalizzato, CASE WHEN tipologia != 'FILM' THEN prima_visione WHEN tipologia = 'FILM' AND anno_produzione < 1984 THEN prima_visione WHEN tipologia = 'FILM' AND (anno_produzione >= 1984 OR anno_produzione IS NULL) AND ncparpgm > 1 THEN prima_visione WHEN tipologia = 'FILM' AND (anno_produzione >= 1984 OR anno_produzione IS NULL) AND (ncparpgm IS NULL OR ncparpgm <= 1) AND _rn_ts = 1 THEN 'S' WHEN tipologia = 'FILM' AND (anno_produzione >= 1984 OR anno_produzione IS NULL) AND (ncparpgm IS NULL OR ncparpgm <= 1) THEN 'N' ELSE COALESCE(prima_visione, 'N') END AS prima_visione_recalc FROM with_rn """) # ------------------------------------------------------------------ # Step 7: scrivi Parquet # ------------------------------------------------------------------ con.execute(f""" COPY ( SELECT rete, data_emissione, ora_inizio, fornitore_diritto, fornitore_cluster, is_primetime_main, timestamp_normalizzato, prima_visione_recalc FROM _em_pv ) TO '{out}' (FORMAT PARQUET, COMPRESSION ZSTD) """) for t in ('_dir_forn', '_forn_map', '_prod_anno', '_em_forn', '_em_pt', '_em_pv'): con.execute(f"DROP TABLE IF EXISTS {t}") return con.execute(f"SELECT COUNT(*) FROM '{out}'").fetchone()[0]