"""
DataHub v2 — CLI entry point
Uso:
python main.py weekly # tutti i Parquet weekly
python main.py daily # gemma + osservatorio
python main.py weekly --dry-run # logga conteggi, non scrive
python main.py prodotti # singolo domain (debug)
python main.py diritti # singolo domain (debug)
"""
import sys
import os
# Python embedded non include la dir corrente nel path — va aggiunto esplicitamente
sys.path.insert(0, os.path.dirname(os.path.abspath(__file__)))
import argparse
import duckdb
from loguru import logger
from config.settings import ensure_dirs, get_log_path, LOG_FILE
# Mappa singolo domain → modulo
DOMAIN_MAP = {
'reti_cluster': 'domains.reti_cluster',
'prodotti': 'domains.prodotti',
'cast': 'domains.cast',
'imdb': 'domains.imdb',
'veg': 'domains.veg',
'scelte_rete': 'domains.scelte_rete',
'boxoffice': 'domains.boxoffice',
'boxoffice_dettaglio': 'domains.boxoffice_dettaglio',
'osservatorio': 'domains.osservatorio',
'diritti': 'domains.diritti',
'diritti_enrich': 'domains.diritti_enrich',
'emesso': 'domains.emesso',
'emesso_enrich': 'domains.emesso_enrich',
'gemma': 'domains.gemma',
}
def _setup_logging(prefix: str) -> None:
log_path = LOG_FILE if LOG_FILE else get_log_path(prefix)
logger.remove()
logger.add(sys.stderr, level='INFO',
format='{time:HH:mm:ss} | {level:<8} | {message}')
logger.add(log_path, level='DEBUG', mode='a',
format='{time:YYYY-MM-DD HH:mm:ss} | {level:<8} | {message}',
encoding='utf-8')
def _run_single(domain: str, dry_run: bool) -> None:
import importlib
mod_path = DOMAIN_MAP.get(domain)
if not mod_path:
logger.error(f"Domain '{domain}' non riconosciuto. Disponibili: {list(DOMAIN_MAP)}")
sys.exit(1)
mod = importlib.import_module(mod_path)
with duckdb.connect() as con:
logger.info(f"Avvio domain: {domain}")
result = mod.run(con)
logger.success(f"{domain}: {result:,} righe")
def main() -> None:
parser = argparse.ArgumentParser(
description='DataHub v2 — pipeline CSV → Parquet'
)
parser.add_argument(
'target',
help=(
'Cosa eseguire: "weekly", "daily", '
'oppure nome di un singolo domain (es. "prodotti")'
)
)
parser.add_argument(
'--dry-run', action='store_true',
help='Esegue senza scrivere file (solo log conteggi)'
)
args = parser.parse_args()
ensure_dirs()
_setup_logging(args.target)
logger.info(f"DataHub v2 — target={args.target} dry_run={args.dry_run}")
if args.target == 'weekly':
from pipelines import weekly
results = weekly.run(dry_run=args.dry_run)
_print_summary(results)
elif args.target == 'daily':
from pipelines import daily
results = daily.run(dry_run=args.dry_run)
_print_summary(results)
elif args.target in DOMAIN_MAP:
_run_single(args.target, dry_run=args.dry_run)
else:
logger.error(
f"Target '{args.target}' non valido. "
f"Usa: weekly | daily | {' | '.join(DOMAIN_MAP)}"
)
sys.exit(1)
def _print_summary(results: dict) -> None:
logger.info("─" * 50)
logger.info("RIEPILOGO")
logger.info("─" * 50)
ok = all(v >= 0 for v in results.values())
for name, count in results.items():
if count < 0:
logger.error(f" {name:<25} ERRORE")
else:
logger.success(f" {name:<25} {count:>10,} righe")
logger.info("─" * 50)
if ok:
logger.success("Pipeline completata con successo.")
else:
logger.warning("Pipeline completata con errori (vedi log sopra).")
if __name__ == '__main__':
main()