""" Parquet_layer1ToParquet_layer2 — trasformazione Layer 1 -> Layer 2 Legge i parquet da Layer 1 e produce i parquet Layer 2. Per la maggior parte dei domini: copia diretta. Per domini con logica di business: trasformazione specifica. Transform con logica: - gemma : flag '!' su campi V_* per titoli IN CORSO - emesso : filtro 34 reti + dedup QUALIFY (max audience) - emesso_enrich: fornitore_cluster, is_primetime_main, prima_visione_recalc, timestamp_normalizzato - diritti_enrich: fornitore_cluster da ragsoc_distr Uso: python main.py daily # gemma + osservatorio python main.py weekly # tutti i parquet """ import sys import shutil import datetime import duckdb from pathlib import Path # Fornitore cluster mapping CSV (nella config di CsvToParquet_layer1) _SCRIPT_DIR = Path(__file__).parent _FORNITORE_CLUSTER_FILE = str( _SCRIPT_DIR.parent / "CsvToParquet_layer1" / "config" / "fornitori_cluster_mapping.csv" ).replace('\\', '/') # 36 reti del filtro emesso Layer 2 _RETI_FILTRO = [ '7C', '7D', 'B5', 'B6', 'BP', 'C5', 'CL', 'CM', 'DD', 'DE', 'DG', 'DI', 'FH', 'FI', 'FT', 'FU', 'GI', 'I1', 'I2', 'KA', 'KB', 'KI', 'KQ', 'LA', 'LB', 'LT', 'MC', 'N1', 'N2', 'N3', 'N4', 'NM', 'NP', 'R4', 'TS', 'WB' ] # Layer 1 — sorgente L1_DIR = Path(r"\\mediaset.it\share\Indirizzo_controllo_risorse\SOFTWARE\PYTHON_SRV\parquet_layer1") # Layer 2 — destinazione (consumer-ready) L2_DIR = Path(r"\\mediaset.it\share\Indirizzo_controllo_risorse\SOFTWARE\PYTHON_SRV\PYTHON_LOCAL\MyICR_Suite\local_db\parquet") # Parquet prodotti dal daily DAILY_PARQUETS = [ "gemma.parquet", "osservatorio.parquet", ] # Tutti i parquet (weekly) ALL_PARQUETS = [ "boxoffice.parquet", "cast.parquet", "diritti.parquet", "diritti_enrich.parquet", "emesso.parquet", "emesso_enrich.parquet", "gemma.parquet", "imdb.parquet", "osservatorio.parquet", "prodotti.parquet", "reti_cluster.parquet", "scelte_rete.parquet", ] def _transform_gemma(src: Path, dst: Path) -> None: """gemma Layer 1 → Layer 2: applica flag '!' sui campi V_*. Per ogni coppia (A_X, V_X): se V_STATO = 'IN CORSO' AND A_X != '0' AND V_X è vuoto/null → V_X = '!' """ src_s = str(src).replace('\\', '/') dst_s = str(dst).replace('\\', '/') with duckdb.connect() as con: con.execute(f"CREATE TABLE _gemma AS SELECT * FROM '{src_s}'") cols = [row[0] for row in con.execute("DESCRIBE _gemma").fetchall() if row[0].upper() != 'COLUMN58'] v_stato = next((c for c in cols if c.upper() == 'V_STATO'), None) parts = [] for col in cols: col_up = col.upper() if col_up.startswith('V_') and col_up not in ('V_STATO', 'V_RDA') and v_stato: a_col = next((c for c in cols if c.upper() == f'A_{col_up[2:]}'), None) if a_col: parts.append(f""" CASE WHEN UPPER(TRIM("{v_stato}")) = 'IN CORSO' AND TRIM("{a_col}") <> '0' AND "{a_col}" IS NOT NULL AND (TRIM("{col}") IS NULL OR TRIM("{col}") = '') THEN '!' ELSE "{col}" END AS "{col}" """) continue parts.append(f'"{col}"') con.execute(f""" COPY (SELECT {', '.join(parts)} FROM _gemma) TO '{dst_s}' (FORMAT PARQUET, COMPRESSION ZSTD) """) def _transform_diritti(src: Path, dst: Path) -> None: """diritti Layer 1 → Layer 2: solo Free+Analogico + Free+DVB-T non contenuto in Analogico. Inclusi: Free+Analogico (tutti) + Free+DVB-T non contenuto (decr/scad) in alcun Analogico per lo stesso prodotto. Esclusi: TVOD, IPTV, DVB-H, Internet, DVB-T il cui periodo è coperto da un Analogico (analog.decr<=dvbt.decr AND analog.scad>=dvbt.scad). Colonne tipo_diritto e piattaforma rimosse (ridondanti dopo il filtro). Date normalizzate: 0001-01-01 → 1001-01-01; tutte le date convertite a DATE. DECR_WIN*/SCAD_WIN* convertite da DD/MM/YYYY a DATE. """ src_s = str(src).replace('\\', '/') dst_s = str(dst).replace('\\', '/') win_replace = '\n'.join( f" TRY_CAST(TRY_STRPTIME(d.DECR_WIN_{i}, '%d/%m/%Y') AS DATE) AS DECR_WIN_{i}," f"\n TRY_CAST(TRY_STRPTIME(d.SCAD_WIN_{i}, '%d/%m/%Y') AS DATE) AS SCAD_WIN_{i}," for i in range(1, 10) ).rstrip(',') with duckdb.connect() as con: con.execute(f"CREATE TABLE _dir AS SELECT * FROM '{src_s}'") con.execute(f""" COPY ( SELECT * EXCLUDE (tipo_diritto, piattaforma) REPLACE ( CASE WHEN d.decr = '0001-01-01' THEN DATE '1001-01-01' ELSE TRY_CAST(d.decr AS DATE) END AS decr, CASE WHEN d.scad = '0001-01-01' THEN DATE '1001-01-01' ELSE TRY_CAST(d.scad AS DATE) END AS scad, TRY_CAST(d.decr_inib AS DATE) AS decr_inib, TRY_CAST(d.scad_inib AS DATE) AS scad_inib, {win_replace} ) FROM _dir d WHERE (UPPER(d.tipo_diritto) = 'FREE' AND UPPER(d.piattaforma) = 'ANALOGICO') OR (UPPER(d.tipo_diritto) = 'FREE' AND UPPER(d.piattaforma) = 'DVB-T' AND NOT EXISTS ( SELECT 1 FROM _dir a WHERE UPPER(a.tipo_diritto) = 'FREE' AND UPPER(a.piattaforma) = 'ANALOGICO' AND a.prod = d.prod AND a.decr <= d.decr AND a.scad >= d.scad )) ) TO '{dst_s}' (FORMAT PARQUET, COMPRESSION ZSTD) """) def _transform_emesso(src: Path, dst: Path) -> None: """emesso Layer 1 → Layer 2: filtro 36 reti + dedup per (rete, data_emissione, ora_inizio).""" src_s = str(src).replace('\\', '/') dst_s = str(dst).replace('\\', '/') reti_sql = ', '.join(f"'{r}'" for r in _RETI_FILTRO) with duckdb.connect() as con: con.execute(f""" COPY ( SELECT * FROM '{src_s}' WHERE rete IN ({reti_sql}) QUALIFY ROW_NUMBER() OVER ( PARTITION BY rete, data_emissione, ora_inizio ORDER BY audience DESC NULLS LAST ) = 1 ) TO '{dst_s}' (FORMAT PARQUET, COMPRESSION ZSTD) """) def _transform_emesso_enrich(dst: Path) -> None: """emesso_enrich Layer 2: fornitore_cluster, is_primetime_main, prima_visione_recalc, timestamp_normalizzato. Legge da: - L2_DIR/emesso.parquet (già filtrato per 34 reti + dedup) - L1_DIR/diritti.parquet - L1_DIR/prodotti.parquet - _FORNITORE_CLUSTER_FILE """ em_pq = str(L2_DIR / 'emesso.parquet').replace('\\', '/') dir_pq = str(L1_DIR / 'diritti.parquet').replace('\\', '/') prod_pq = str(L1_DIR / 'prodotti.parquet').replace('\\', '/') dst_s = str(dst).replace('\\', '/') with duckdb.connect() as con: # diritti filtrati Free/Analogico + fallback DVB-T per prodotti senza Analogico con.execute(f""" CREATE TABLE _has_analogico_em AS SELECT DISTINCT prod FROM '{dir_pq}' WHERE UPPER(tipo_diritto) = 'FREE' AND UPPER(piattaforma) = 'ANALOGICO' AND prod IS NOT NULL """) con.execute(f""" CREATE 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' OR (UPPER(piattaforma) = 'DVB-T' AND prod NOT IN (SELECT prod FROM _has_analogico_em)) ) AND prod IS NOT NULL AND decr IS NOT NULL AND scad IS NOT NULL """) con.execute(f""" CREATE TABLE _forn_map AS SELECT UPPER(TRIM(fornitore_diritto)) AS fornitore_key, fornitore_cluster FROM read_csv('{_FORNITORE_CLUSTER_FILE}', delim=';', encoding='cp1252', header=true) """) con.execute(f""" CREATE TABLE _prod_anno AS SELECT prodotto, MAX(anno_produzione) AS anno_produzione FROM '{prod_pq}' WHERE prodotto IS NOT NULL GROUP BY prodotto """) # join emesso → diritti → fornitore_cluster con.execute(f""" CREATE 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 """) # is_primetime_main (fascia PR, copertura 21:10–24:00) con.execute(""" CREATE 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 """) # prima_visione_recalc + timestamp_normalizzato con.execute(""" CREATE 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 """) 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 '{dst_s}' (FORMAT PARQUET, COMPRESSION ZSTD) """) def _transform_diritti_enrich(dst: Path) -> None: """diritti_enrich Layer 2: fornitore_cluster + ESCLUSIVO_TEM da ragsoc_distr. Copre solo le righe presenti in L2 diritti.parquet (stessa logica di _transform_diritti): - Free+Analogico (tutti) → ESCLUSIVO_TEM=FALSE - Free+DVB-T non contenuto in alcun Analogico → ESCLUSIVO_TEM=TRUE "Contenuto" = esiste un Analogico per lo stesso prodotto con decr<=dvbt.decr AND scad>=dvbt.scad. Legge da: - L1_DIR/diritti.parquet - _FORNITORE_CLUSTER_FILE """ dir_pq = str(L1_DIR / 'diritti.parquet').replace('\\', '/') dst_s = str(dst).replace('\\', '/') with duckdb.connect() as con: con.execute(f""" CREATE TABLE _forn_map AS SELECT UPPER(TRIM(fornitore_diritto)) AS fornitore_key, fornitore_cluster FROM read_csv('{_FORNITORE_CLUSTER_FILE}', delim=';', encoding='cp1252', header=true) """) con.execute(f"CREATE TABLE _dir AS SELECT * FROM '{dir_pq}' WHERE UPPER(tipo_diritto) = 'FREE'") con.execute(f""" COPY ( SELECT d.id_diritto, CASE WHEN d.ragsoc_distr IS NOT NULL THEN COALESCE(fm.fornitore_cluster, 'altro') ELSE NULL END AS fornitore_cluster, (UPPER(d.piattaforma) = 'DVB-T') AS ESCLUSIVO_TEM FROM _dir d LEFT JOIN _forn_map fm ON UPPER(TRIM(d.ragsoc_distr)) = fm.fornitore_key WHERE (UPPER(d.piattaforma) = 'ANALOGICO') OR (UPPER(d.piattaforma) = 'DVB-T' AND NOT EXISTS ( SELECT 1 FROM _dir a WHERE UPPER(a.piattaforma) = 'ANALOGICO' AND a.prod = d.prod AND a.decr <= d.decr AND a.scad >= d.scad )) ) TO '{dst_s}' (FORMAT PARQUET, COMPRESSION ZSTD) """) def _transform_boxoffice(src: Path, dst: Path) -> None: """boxoffice Layer 1 → Layer 2: una riga per prodotto. Solo prodotti con almeno una riga 'D' (prima uscita cinema). data_debutto = MIN(data_debutto) dalle righe 'D' → tipo DATE. stagione, distributore, sale dalla prima riga 'D' (ORDER BY data_debutto ASC). incasso e spettatori = SUM su tutte le righe del prodotto. Le 168K righe senza tipo_programmazione (solo codice prodotto) vengono escluse. """ src_s = str(src).replace('\\', '/') dst_s = str(dst).replace('\\', '/') with duckdb.connect() as con: con.execute(f""" COPY ( WITH _d AS ( SELECT prodotto, MIN(TRY_CAST(data_debutto AS DATE)) AS data_debutto, FIRST(stagione ORDER BY TRY_CAST(data_debutto AS DATE) ASC) AS stagione, FIRST(distributore ORDER BY TRY_CAST(data_debutto AS DATE) ASC) AS distributore FROM '{src_s}' WHERE tipo_programmazione = 'D' AND prodotto IS NOT NULL GROUP BY prodotto ), _agg AS ( SELECT prodotto, SUM(incasso) AS incasso, CAST(SUM(spettatori) AS BIGINT) AS spettatori FROM '{src_s}' WHERE tipo_programmazione IS NOT NULL AND prodotto IS NOT NULL GROUP BY prodotto ) SELECT d.prodotto, d.stagione, d.data_debutto, d.distributore, a.incasso, a.spettatori FROM _d d LEFT JOIN _agg a ON a.prodotto = d.prodotto ORDER BY d.prodotto ) TO '{dst_s}' (FORMAT PARQUET, COMPRESSION ZSTD) """) def _transform_scelte_rete(src: Path, dst: Path) -> None: """scelte_rete Layer 1 → Layer 2: solo stagione corrente, rete_estesa → codice_rete, slot abbreviato.""" src_s = str(src).replace('\\', '/') dst_s = str(dst).replace('\\', '/') reti_s = str(L1_DIR / 'reti_cluster.parquet').replace('\\', '/') today = datetime.date.today() if today > datetime.date(today.year, 9, 10): current_season = f"{today.year}/{today.year + 1}" else: current_season = f"{today.year - 1}/{today.year}" with duckdb.connect() as con: con.execute(f""" COPY ( SELECT DISTINCT s.prodotto, COALESCE(r.codice_rete, s.rete) AS codice_rete, CASE s.slot WHEN 'NOTTE' THEN 'NO' WHEN 'PRIMETIME GARANZIA' THEN 'PT G' WHEN 'MATTINA' THEN 'MA' WHEN 'POMERIGGIO' THEN 'PO' WHEN 'SECONDA SERATA GARANZIA' THEN 'SS G' WHEN 'PRIMETIME ESTATE' THEN 'PT E' WHEN 'SECONDA SERATA ESTATE' THEN 'SS E' WHEN 'ACCESS PRIMETIME GARANZIA' THEN 'AP G' WHEN 'PRIMETIME STRENNE' THEN 'PT S' WHEN 'SECONDA SERATA STRENNE' THEN 'SS S' WHEN 'ACCESS PRIMETIME STRENNE' THEN 'AP S' WHEN 'ACCESS PRIMETIME ESTATE' THEN 'AP E' END AS slot FROM '{src_s}' s LEFT JOIN '{reti_s}' r ON s.rete = r.rete_estesa WHERE s.stagione = '{current_season}' AND s.rete IS NOT NULL AND s.rete != 'None' AND s.slot IS NOT NULL AND s.slot != 'None' ) TO '{dst_s}' (FORMAT PARQUET, COMPRESSION ZSTD) """) def _transform_imdb(src: Path, dst: Path) -> None: """imdb Layer 1 → Layer 2: drop edizione (sempre 1, zero valore informativo).""" src_s = str(src).replace('\\', '/') dst_s = str(dst).replace('\\', '/') with duckdb.connect() as con: con.execute(f""" COPY ( SELECT codice, riferimento_imdb FROM '{src_s}' ) TO '{dst_s}' (FORMAT PARQUET, COMPRESSION ZSTD) """) def transform(name: str, src: Path, dst: Path) -> None: """Dispatcha alla trasformazione specifica del dominio.""" if name == 'gemma.parquet': _transform_gemma(src, dst) elif name == 'emesso.parquet': _transform_emesso(src, dst) elif name == 'boxoffice.parquet': _transform_boxoffice(src, dst) elif name == 'diritti.parquet': _transform_diritti(src, dst) elif name == 'emesso_enrich.parquet': _transform_emesso_enrich(dst) elif name == 'diritti_enrich.parquet': _transform_diritti_enrich(dst) elif name == 'scelte_rete.parquet': _transform_scelte_rete(src, dst) elif name == 'imdb.parquet': _transform_imdb(src, dst) else: shutil.copy2(src, dst) # Parquet calcolati da sorgenti multiple — non esistono in L1_DIR _COMPUTED = {'emesso_enrich.parquet', 'diritti_enrich.parquet'} def run(mode: str) -> None: if mode == "daily": files = DAILY_PARQUETS elif mode == "weekly": files = ALL_PARQUETS else: print(f"ERRORE: mode '{mode}' non valido. Usa: daily | weekly") sys.exit(1) L2_DIR.mkdir(parents=True, exist_ok=True) ok, failed = 0, [] for name in files: src = L1_DIR / name dst = L2_DIR / name # I parquet _COMPUTED vengono calcolati dalle funzioni dedicate, # non esistono in Layer 1 quindi non si controlla src.exists() if name not in _COMPUTED and not src.exists(): print(f" SKIP {name} — non trovato in Layer 1") continue try: transform(name, src, dst) print(f" OK {name}") ok += 1 except Exception as e: print(f" ERR {name} — {e}") failed.append(name) print(f"\nlayer1->layer2 {mode}: {ok} prodotti, {len(failed)} errori.") if failed: sys.exit(1) if __name__ == "__main__": if len(sys.argv) != 2: print("Uso: python main.py daily | weekly") sys.exit(1) run(sys.argv[1])