""" ImdbUpdate v2 - Pipeline DuckDB Scarica dataset IMDb, trasforma via DuckDB in-memory, produce imdb.parquet e db_split_IMDB.accdb. """ import logging import shutil from pathlib import Path from datetime import datetime import sys sys.path.insert(0, str(Path(__file__).parent)) from config import ( BASE_DIR, TEMP_DIR, EMPTY_DIR, ACCESS_DB_TEMPLATE, ACCESS_DB_FINAL, IMDB_URLS, EXPORT_PATHS, GEMMA_LOADER_PATH, GEMMA_LOADER_PASSWORD, LOG_DIR, MAX_DIRECTORS, MAX_ACTORS, DUCKDB_THREADS, DUCKDB_MEMORY_LIMIT, PARQUET_OUTPUT_PATH, PARQUET_PUBLISH_PATHS, ) from src.downloader import ImdbDownloader from src.transformer import ImdbTransformer from src.access_exporter import AccessExporter logger = logging.getLogger(__name__) def setup_logging(log_file: Path = None): log_format = '%(asctime)s - %(name)s - %(levelname)s - %(message)s' date_format = '%Y-%m-%d %H:%M:%S' handlers = [logging.StreamHandler()] if log_file: log_file.parent.mkdir(parents=True, exist_ok=True) handlers.append(logging.FileHandler(log_file, encoding='utf-8')) logging.basicConfig( level=logging.INFO, format=log_format, datefmt=date_format, handlers=handlers ) def prepare_temp_directory(): logger.info("Preparazione directory temporanea") if TEMP_DIR.exists(): logger.info(f"Eliminazione {TEMP_DIR}") shutil.rmtree(TEMP_DIR) TEMP_DIR.mkdir(parents=True) logger.info(f"Directory creata: {TEMP_DIR}") def download_imdb_datasets(): logger.info("="*60) logger.info("FASE 2: Download dataset IMDb") logger.info("="*60) downloader = ImdbDownloader(TEMP_DIR) files = downloader.download_and_extract_all(IMDB_URLS) logger.info(f"Download completato: {len(files)} file") return files def transform_and_export_parquet(files: dict) -> bool: """ Step 3: DuckDB transform → imdb.parquet + copia in publish paths. """ logger.info("="*60) logger.info("FASE 3: Transform DuckDB → imdb.parquet") logger.info("="*60) with ImdbTransformer( temp_dir=TEMP_DIR, parquet_output=PARQUET_OUTPUT_PATH, max_directors=MAX_DIRECTORS, max_actors=MAX_ACTORS, duckdb_threads=DUCKDB_THREADS, duckdb_memory_limit=DUCKDB_MEMORY_LIMIT, ) as transformer: row_count = transformer.transform(files) logger.info(f"Transform completato: {row_count:,} righe in parquet") # Copia parquet nelle destinazioni di rete for dest_dir in PARQUET_PUBLISH_PATHS: dest = Path(dest_dir) / PARQUET_OUTPUT_PATH.name try: Path(dest_dir).mkdir(parents=True, exist_ok=True) shutil.copy2(PARQUET_OUTPUT_PATH, dest) logger.info(f" Parquet copiato in: {dest}") except Exception as e: logger.warning(f" Copia parquet fallita in {dest_dir}: {e}") return True def export_to_access(): """Step 4: Parquet → Access (.accdb) + copia destinazioni.""" logger.info("="*60) logger.info("FASE 4: Export Parquet → MS Access") logger.info("="*60) exporter = AccessExporter() template = EMPTY_DIR / ACCESS_DB_TEMPLATE if not template.exists(): logger.error(f"Template non trovato: {template}") logger.error("Crea cartella Empty/ e inserisci db_split_IMDB.accdb vuoto") return False exporter.create_access_db(template, ACCESS_DB_FINAL) exporter.export_parquet_to_access(PARQUET_OUTPUT_PATH, ACCESS_DB_FINAL, "IMDB") exporter.enrich_ott_in_access(ACCESS_DB_FINAL) exporter.compact_access_db(ACCESS_DB_FINAL) exporter.copy_to_destinations(ACCESS_DB_FINAL, EXPORT_PATHS) return True def import_to_gemma_loader(): """Step 5: IMDB → GemmaLoader via tabella linkata (opzionale).""" logger.info("="*60) logger.info("FASE 5: Copia IMDB in GemmaLoader (opzionale)") logger.info("="*60) try: exporter = AccessExporter() exporter.copy_to_gemma_loader(PARQUET_OUTPUT_PATH, GEMMA_LOADER_PATH, GEMMA_LOADER_PASSWORD) except Exception as e: logger.warning(f"Copia GemmaLoader fallita (opzionale): {e}") def write_log_file(): logger.info("="*60) logger.info("FASE 6: Scrittura log finale") logger.info("="*60) try: LOG_DIR.mkdir(parents=True, exist_ok=True) for old_log in LOG_DIR.glob("UPDATE_IMDB *.txt"): old_log.unlink() log_filename = f"UPDATE_IMDB al {datetime.now().strftime('%d-%m-%y')}.txt" log_file = LOG_DIR / log_filename log_file.touch() logger.info(f"Log creato: {log_file}") except Exception as e: logger.warning(f"Errore scrittura log: {e}") def cleanup_temp_files(): """Pulizia completa file temporanei (tsv.gz, txt, accdb).""" logger.info("="*60) logger.info("FASE 7: Pulizia file temporanei") logger.info("="*60) if not TEMP_DIR.exists(): logger.info("Directory Temp/ non esiste, skip pulizia") return total_freed = 0 for pattern in ["*.tsv.gz", "*.txt", "*.accdb", "*.duckdb", "*.duckdb.wal"]: for f in TEMP_DIR.glob(pattern): size = f.stat().st_size f.unlink() total_freed += size logger.info(f" Eliminato: {f.name} ({size / (1024**3):.2f} GB)") logger.info(f"Spazio liberato totale: {total_freed / (1024**3):.2f} GB") def cleanup_on_startup(): """Pulizia file residui da esecuzioni precedenti interrotte.""" if not TEMP_DIR.exists(): return logger.info("Pulizia file residui all'avvio...") files_to_clean = ( list(TEMP_DIR.glob("*.tsv.gz")) + list(TEMP_DIR.glob("*.txt")) + list(TEMP_DIR.glob("*.accdb")) + list(TEMP_DIR.glob("*.duckdb")) + list(TEMP_DIR.glob("*.duckdb.wal")) ) if not files_to_clean: logger.info("Nessun file residuo trovato") return total_size = sum(f.stat().st_size for f in files_to_clean) logger.info(f"Trovati {len(files_to_clean)} file residui ({total_size / (1024**3):.2f} GB)") for f in files_to_clean: f.unlink() logger.info(f" Eliminato: {f.name}") logger.info("Pulizia completata") def main(skip_download: bool = False, keep_temp: bool = False): """ Pipeline v2: download → DuckDB transform → parquet → Access → GemmaLoader Args: skip_download: Se True, salta download (usa file esistenti in Temp/) keep_temp: Se True, mantiene file temporanei """ start_time = datetime.now() log_file = BASE_DIR / "logs" / f"imdb_update_{start_time.strftime('%Y%m%d_%H%M%S')}.log" setup_logging(log_file) logger.info("="*60) logger.info("ImdbUpdate v2 - Avvio (DuckDB pipeline)") logger.info("="*60) logger.info(f"Base directory: {BASE_DIR}") logger.info(f"Temp directory: {TEMP_DIR}") logger.info(f"Parquet output: {PARQUET_OUTPUT_PATH}") try: # 0. Pulizia residui all'avvio if not skip_download: cleanup_on_startup() # 1. Prepara directory if not skip_download: prepare_temp_directory() # 2. Download dataset if skip_download: logger.warning("Download saltato - uso file esistenti in Temp/") files = { 'name_basics': TEMP_DIR / 'name_basics.txt', 'title_akas': TEMP_DIR / 'title_akas.txt', 'title_basics': TEMP_DIR / 'title_basics.txt', 'title_ratings': TEMP_DIR / 'title_ratings.txt', 'title_principals': TEMP_DIR / 'title_principals.txt', } else: files = download_imdb_datasets() # 3. Transform DuckDB → parquet if not transform_and_export_parquet(files): logger.error("Transform fallito") return False # 4. Export Access if not export_to_access(): logger.error("Export Access fallito") return False # 5. Import GemmaLoader import_to_gemma_loader() # 6. Log finale write_log_file() # 7. Cleanup if not keep_temp: cleanup_temp_files() else: logger.info("File temporanei preservati (keep_temp=True)") elapsed = datetime.now() - start_time logger.info("="*60) logger.info(f"Completato con successo in {elapsed}") logger.info("="*60) return True except KeyboardInterrupt: logger.warning("\nInterrotto dall'utente") return False except Exception as e: logger.exception(f"Errore fatale: {e}") return False if __name__ == '__main__': import argparse parser = argparse.ArgumentParser(description='ImdbUpdate v2 - DuckDB pipeline') parser.add_argument( '--skip-download', action='store_true', help='Salta download, usa file esistenti in Temp/' ) parser.add_argument( '--keep-temp', action='store_true', help='Mantiene file temporanei in Temp/' ) args = parser.parse_args() success = main( skip_download=args.skip_download, keep_temp=args.keep_temp ) sys.exit(0 if success else 1)