""" ImdbTransformer - DuckDB-based transformation (v2) Sostituisce database.py + processor.py. Legge TSV direttamente via DuckDB, produce imdb.parquet. """ import logging from pathlib import Path import duckdb logger = logging.getLogger(__name__) class ImdbTransformer: """ Trasformatore DuckDB per i dataset IMDb. Carica TSV filtrati in-memory, esegue join per step, scrive parquet. """ def __init__(self, temp_dir: Path, parquet_output: Path, max_directors: int = 5, max_actors: int = 5, duckdb_threads: int = 0, duckdb_memory_limit: str = '4GB'): self.temp_dir = Path(temp_dir) self.parquet_output = Path(parquet_output) self.max_directors = max_directors self.max_actors = max_actors self.duckdb_threads = duckdb_threads self.duckdb_memory_limit = duckdb_memory_limit self._con = None def __enter__(self): self._con = duckdb.connect() # Thread singolo: riduce picco di memoria parallela su macchine con 4GB RAM self._con.execute("SET threads TO 1") self._con.execute(f"SET memory_limit='{self.duckdb_memory_limit}'") self._con.execute("SET preserve_insertion_order=false") # Spill-to-disk: evita OOM su join grandi (~31M × ~15M righe) temp_dir_str = str(self.temp_dir).replace('\\', '/') self._con.execute(f"SET temp_directory='{temp_dir_str}'") return self def __exit__(self, *args): if self._con: self._con.close() self._con = None def transform(self, files: dict) -> int: """ Entry point. Carica tabelle, trasforma, scrive parquet. Ritorna row count. """ self._load_title_basics(files['title_basics']) self._load_title_akas(files['title_akas']) self._load_title_ratings(files['title_ratings']) self._load_title_principals(files['title_principals']) self._load_name_basics(files['name_basics']) self._build_new_imdb() return self._write_parquet() # ------------------------------------------------------------------ # Caricamento tabelle sorgente # ------------------------------------------------------------------ def _load_title_basics(self, path: Path): logger.info(f"Caricamento title_basics: {path.name}") path_str = str(path).replace('\\', '/') self._con.execute(f""" CREATE OR REPLACE TABLE title_basics AS SELECT tconst, titleType, SUBSTR(primaryTitle, 1, 100) AS primaryTitle, runtimeMinutes, startYear, genres FROM read_csv('{path_str}', delim='\t', header=true, nullstr='\\N', ignore_errors=true) WHERE titleType IN ('movie', 'tvMiniSeries', 'tvMovie', 'tvSeries') """) count = self._con.execute("SELECT COUNT(*) FROM title_basics").fetchone()[0] logger.info(f" title_basics: {count:,} righe (filtrate per titleType)") def _load_title_akas(self, path: Path): logger.info(f"Caricamento title_akas: {path.name}") path_str = str(path).replace('\\', '/') self._con.execute(f""" CREATE OR REPLACE TABLE title_akas AS SELECT titleId, title FROM read_csv('{path_str}', delim='\t', header=true, nullstr='\\N', ignore_errors=true) WHERE region = 'IT' """) count = self._con.execute("SELECT COUNT(*) FROM title_akas").fetchone()[0] logger.info(f" title_akas: {count:,} righe (solo region=IT)") def _load_title_ratings(self, path: Path): logger.info(f"Caricamento title_ratings: {path.name}") path_str = str(path).replace('\\', '/') self._con.execute(f""" CREATE OR REPLACE TABLE title_ratings AS SELECT tconst, averageRating, numVotes FROM read_csv('{path_str}', delim='\t', header=true, nullstr='\\N', ignore_errors=true) """) count = self._con.execute("SELECT COUNT(*) FROM title_ratings").fetchone()[0] logger.info(f" title_ratings: {count:,} righe") def _load_title_principals(self, path: Path): logger.info(f"Caricamento title_principals: {path.name}") path_str = str(path).replace('\\', '/') self._con.execute(f""" CREATE OR REPLACE TABLE title_principals AS SELECT tconst, nconst, ordering, category FROM read_csv('{path_str}', delim='\t', header=true, nullstr='\\N', ignore_errors=true) WHERE category IN ('director', 'actor') """) count = self._con.execute("SELECT COUNT(*) FROM title_principals").fetchone()[0] logger.info(f" title_principals: {count:,} righe (director+actor)") def _load_name_basics(self, path: Path): logger.info(f"Caricamento name_basics: {path.name}") path_str = str(path).replace('\\', '/') self._con.execute(f""" CREATE OR REPLACE TABLE name_basics AS SELECT nconst, primaryName FROM read_csv('{path_str}', delim='\t', header=true, nullstr='\\N', ignore_errors=true) """) count = self._con.execute("SELECT COUNT(*) FROM name_basics").fetchone()[0] logger.info(f" name_basics: {count:,} righe") # ------------------------------------------------------------------ # Trasformazione principale # ------------------------------------------------------------------ def _build_new_imdb(self): """ Costruisce new_imdb in step separati per limitare il picco di memoria. Ogni tabella intermedia viene materializzata prima del join finale. """ logger.info("Costruzione new_imdb (step separati)...") # Step 1: titoli italiani logger.info(" Step 1/4: titoli italiani") self._con.execute(""" CREATE OR REPLACE TABLE italian_titles AS SELECT titleId, ANY_VALUE(title) AS ti FROM title_akas GROUP BY titleId """) # Step 2: registi (top N per ordering) # Pre-filtro su title_basics riduce title_principals da 31M a ~pochi M logger.info(" Step 2/4: registi") self._con.execute(f""" CREATE OR REPLACE TABLE directors AS SELECT p.tconst, STRING_AGG(n.primaryName, ', ' ORDER BY p.ordering) AS regista FROM ( SELECT tp.tconst, tp.nconst, tp.ordering, ROW_NUMBER() OVER (PARTITION BY tp.tconst ORDER BY tp.ordering) AS rn FROM title_principals tp WHERE tp.category = 'director' AND tp.tconst IN (SELECT tconst FROM title_basics) ) p JOIN name_basics n ON p.nconst = n.nconst WHERE p.rn <= {self.max_directors} GROUP BY p.tconst """) # Step 3: attori (top N per ordering) logger.info(" Step 3/4: attori") self._con.execute(f""" CREATE OR REPLACE TABLE actors AS SELECT p.tconst, STRING_AGG(n.primaryName, ', ' ORDER BY p.ordering) AS cast_str FROM ( SELECT tp.tconst, tp.nconst, tp.ordering, ROW_NUMBER() OVER (PARTITION BY tp.tconst ORDER BY tp.ordering) AS rn FROM title_principals tp WHERE tp.category = 'actor' AND tp.tconst IN (SELECT tconst FROM title_basics) ) p JOIN name_basics n ON p.nconst = n.nconst WHERE p.rn <= {self.max_actors} GROUP BY p.tconst """) # Step 4: join finale logger.info(" Step 4/4: join finale") self._con.execute(""" CREATE OR REPLACE TABLE new_imdb AS SELECT b.tconst AS IMDB_CODICE, b.titleType AS TIPOL, b.primaryTitle AS "TO", it.ti AS TI, TRY_CAST(b.runtimeMinutes AS INTEGER) AS DUR, TRY_CAST(b.startYear AS INTEGER) AS ANNO, b.genres AS GENERE, r.averageRating AS VOTO, r.numVotes AS VOTANTI, d.regista AS REGISTA, a.cast_str AS "CAST" FROM title_basics b LEFT JOIN italian_titles it ON b.tconst = it.titleId LEFT JOIN title_ratings r ON b.tconst = r.tconst LEFT JOIN directors d ON b.tconst = d.tconst LEFT JOIN actors a ON b.tconst = a.tconst """) count = self._con.execute("SELECT COUNT(*) FROM new_imdb").fetchone()[0] logger.info(f" new_imdb: {count:,} righe") def _write_parquet(self) -> int: """Scrive new_imdb su file parquet. Ritorna row count.""" self.parquet_output.parent.mkdir(parents=True, exist_ok=True) parquet_str = str(self.parquet_output).replace('\\', '/') logger.info(f"Scrittura parquet: {self.parquet_output}") self._con.execute(f""" COPY (SELECT * FROM new_imdb) TO '{parquet_str}' (FORMAT PARQUET, COMPRESSION ZSTD, ROW_GROUP_SIZE 100000) """) count = self._con.execute(f"SELECT COUNT(*) FROM '{parquet_str}'").fetchone()[0] size_mb = self.parquet_output.stat().st_size / (1024 * 1024) logger.info(f" Parquet scritto: {count:,} righe, {size_mb:.1f} MB") return count