""" ParquetToAccess — convertitore centralizzato Parquet -> MS Access Flusso: 1. Copia DB_DAY_master.accdb -> output_db locale 2. INSERT di tutti i job (Parquet -> tabelle Access locali) 3. Pubblica output_db su destinazioni rete Uso: python main.py # esegue tutti i job python main.py --job nome # esegue solo il job specificato python main.py --no-publish # salta la copia su rete (debug) """ import json import os import sys import shutil import argparse from pathlib import Path from datetime import datetime class _Tee: """Scrive su stdout e su file log contemporaneamente.""" def __init__(self, log_path: Path, mode: str = 'a'): log_path.parent.mkdir(parents=True, exist_ok=True) self._file = open(log_path, mode, encoding='utf-8') self._stdout = sys.stdout def write(self, data): self._stdout.write(data) self._file.write(data) def flush(self): self._stdout.flush() self._file.flush() def close(self): sys.stdout = self._stdout self._file.close() BASE_DIR = Path(__file__).parent EXPORTS_FILE = BASE_DIR / 'config' / 'exports.json' SETTINGS_FILE = BASE_DIR / 'config' / 'settings.json' def ts(): return datetime.now().strftime('%H:%M:%S') def load_json(path): with open(path, encoding='utf-8') as f: return json.load(f) def step_copy_master(settings: dict): master = BASE_DIR / settings['master_db'] output = Path(settings['output_db']) if not master.exists(): raise FileNotFoundError(f"Master DB non trovato: {master}") output.parent.mkdir(parents=True, exist_ok=True) # Pulizia: rimuovi DB e lock file residui da run precedenti lock = output.with_suffix('.laccdb') for f in (output, lock): if f.exists(): try: f.unlink() except PermissionError: raise PermissionError(f"File in uso, chiudi Access prima di procedere: {f}") shutil.copy2(master, output) print(f"[{ts()}] Master copiato -> {output}") def run_job(job: dict) -> bool: name = job['name'] print(f"[{ts()}] {name} — avvio") try: from core.exporter import export total = export(job) print(f"[{ts()}] {name} — OK ({total:,} righe -> [{job['table']}])") return True except Exception as e: print(f"[{ts()}] {name} — ERRORE: {e}") return False def step_publish(settings: dict): output = Path(settings['output_db']) for entry in settings.get('publish', []): dest = Path(entry['dest']) dest.parent.mkdir(parents=True, exist_ok=True) shutil.copy2(output, dest) print(f"[{ts()}] Pubblicato -> {dest}") def main(): parser = argparse.ArgumentParser(description='ParquetToAccess') parser.add_argument('--job', help='Esegui solo questo job (per nome)') parser.add_argument('--no-publish', action='store_true', help='Salta pubblicazione su rete') parser.add_argument('--no-enrich', action='store_true', help='Salta enrichment (solo populate)') args = parser.parse_args() settings = load_json(SETTINGS_FILE) exports = load_json(EXPORTS_FILE) # Setup log file tee = None pipeline_log = os.environ.get('PIPELINE_LOG_FILE') if pipeline_log: log_path = Path(pipeline_log) else: log_dir = settings.get('enrichment', {}).get('pta_log_dir') log_path = Path(log_dir) / f"pta_{datetime.now().strftime('%Y%m%d_%H%M%S')}.log" if log_dir else None if log_path: tee = _Tee(log_path, mode='a') sys.stdout = tee if args.job: jobs = [j for j in exports if j['name'] == args.job] if not jobs: print(f"ERRORE: job '{args.job}' non trovato in exports.json") sys.exit(1) else: jobs = exports print(f"ParquetToAccess — {len(jobs)} job da eseguire") print() # Step 1 — copia master in locale step_copy_master(settings) print() # Step 2 — INSERT tutti i job failed = [] for job in jobs: # inietta target dal settings se non specificato nel job if 'target' not in job: job['target'] = settings['output_db'] ok = run_job(job) if not ok: failed.append(job['name']) print() if failed: print(f"ERRORE: {len(failed)} job falliti: {', '.join(failed)}") if tee: tee.close() sys.exit(1) # Step 3 — enrichment enrich_cfg = settings.get('enrichment', {}) if enrich_cfg.get('enabled', False) and not args.no_enrich: print(f"[{ts()}] Enrichment — avvio") from core.enrichment import run_all run_all(settings['output_db'], enrich_cfg) print() # Step 4 — pubblica su rete if not args.no_publish: step_publish(settings) print() print(f"Completato — {len(jobs)} job eseguiti con successo.") if tee: tee.close() sys.exit(0) if __name__ == '__main__': main()