Skip to content

Latest commit

 

History

History
131 lines (109 loc) · 6.46 KB

File metadata and controls

131 lines (109 loc) · 6.46 KB

TODO: Migrazione Celery + Redis

Obiettivo: sostituire l'esecuzione custom del worker con Celery + Redis, mantenendo Postgres come fonte ufficiale per stato, audit e tracciabilita dei job.

Fonti operative:

Decisione architetturale

  • Usare Redis solo come broker Celery.
  • Valutare se usare Redis anche come result backend; preferenza iniziale: no, perche lo stato business resta in parse_jobs.
  • Tenere Postgres come fonte ufficiale per:
    • parse_jobs.status
    • parse_jobs.attempts
    • parse_jobs.last_error
    • audit_logs
    • flight_logs.parse_status
  • Rendere i task Celery idempotenti: una riesecuzione non deve duplicare track point, eventi o audit critici.

Dipendenze e configurazione

  • Aggiungere celery[redis] o celery + redis in api/requirements.txt.
  • Aggiungere variabili in api/.env.example:
    • CELERY_BROKER_URL=redis://redis:6379/0
    • CELERY_RESULT_BACKEND=
    • CELERY_TASK_TIME_LIMIT=300
    • CELERY_TASK_SOFT_TIME_LIMIT=240
    • CELERY_REDIS_VISIBILITY_TIMEOUT=900
  • Aggiungere servizio redis in docker-compose.yml.
  • Aggiornare worker in docker-compose.yml per eseguire Celery:
    • celery -A app.celery_app worker --loglevel=INFO --concurrency=2
  • Valutare un container separato per celery beat solo se servono job periodici.

Implementazione backend

  • Creare api/app/celery_app.py.
  • Configurare Celery con:
    • broker Redis
    • task routing dedicato per parsing DJI
    • task_acks_late=True
    • worker_prefetch_multiplier=1
    • task_reject_on_worker_lost=True
    • visibility timeout Redis coerente con il tempo massimo previsto del parser
  • Creare api/app/tasks/parse.py.
  • Implementare task parse_flight_log_task(parse_job_id: int).
  • Il task deve:
    • leggere ParseJob da Postgres
    • verificare che il job sia ancora valido
    • marcare running
    • chiamare parse_flight_log_or_raise
    • marcare succeeded o failed
    • scrivere audit_logs
    • gestire retry con backoff
  • Modificare enqueue_parse_job per chiamare parse_flight_log_task.delay(job.id).
  • Tenere una modalita fallback senza Celery solo se utile per sviluppo locale.

Idempotenza e concorrenza

  • Prima del parsing, bloccare il job con lock DB o transizione atomica queued -> running.
  • Evitare doppia esecuzione contemporanea dello stesso flight_log_id.
  • Garantire che _save_track_points continui a cancellare e ricreare i punti del volo.
  • Verificare che _extract_events non accumuli duplicati a ogni retry.
  • Audit: evitare spam su retry ravvicinati o distinguere chiaramente started, retry_scheduled, failed, succeeded.
  • Gestire task redelivered da Redis visibility timeout.

API e UI

  • Lasciare invariata la risposta upload: flight_id, log_id, job_id, parse_status.
  • Esporre lo stato Celery solo se utile per debug; lo stato utente deve restare quello DB.
  • Aggiornare polling frontend solo se cambia il formato di /logs.
  • Aggiungere messaggio chiaro quando un job e in retry.

Osservabilita

  • Log strutturati nel worker Celery.
  • Metriche minime:
    • job queued/running/failed/succeeded
    • durata parsing
    • numero retry
    • errori DJI API/subprocess
  • Valutare Flower solo per ambiente admin/dev, non esposto pubblicamente.

Scalabilita Postgres

  • Monitorare crescita di flights, flight_logs, flight_track_points, flight_events, parse_jobs per organizzazione.
  • Aggiungere indici composti mirati prima del partizionamento:
    • flights (organization_id, started_at DESC)
    • flight_logs (flight_id, parse_status)
    • flight_track_points (flight_id, id)
    • flight_events (flight_id, severity, flight_time_s)
    • parse_jobs (organization_id, status, next_run_at)
  • Valutare partizionamento Postgres quando una tabella supera milioni di righe o query/report per org diventano lente.
  • Opzione preferita per multi-tenant: partizionare per organization_id dove il pattern query e sempre tenant-scoped.
  • Valutare partizionamento ibrido per tabelle molto grandi:
    • flights: partition/hash o list per organization_id
    • flight_track_points: partition per flight_id derivato o per organization_id denormalizzato
    • flight_events: partition per organization_id denormalizzato o per periodo se report storici pesano
  • Considerare una colonna organization_id denormalizzata su flight_track_points e flight_events per evitare join pesanti quando si partiziona per tenant.
  • Prima di partizionare, misurare con EXPLAIN ANALYZE le query reali di dashboard, report, mappe e PDF.
  • Se pochi tenant diventano enormi, valutare sharding/logical database per org enterprise invece di partizionare tutto subito.
  • Documentare strategia retention/archiviazione:
    • mantenere track point completi per N mesi
    • mantenere downsample storico per mappe/report
    • conservare PDF/audit/hash log secondo policy cliente
  • Verificare che Alembic supporti bene la strategia scelta prima di migrare dati produzione.

Test

  • Unit test per enqueue_parse_job.
  • Unit test task Celery in eager mode.
  • Test idempotenza: eseguire due volte lo stesso job e verificare no duplicati.
  • Test retry su errore temporaneo del parser.
  • Test failure finale dopo max_attempts.
  • Test Docker Compose: upload log -> job queued -> worker parsed -> UI vede parsed.

Cutover

  • Prima fase: aggiungere Celery mantenendo il worker custom disabilitabile.
  • Seconda fase: usare Celery come path default in Docker Compose.
  • Terza fase: rimuovere api/app/worker.py se non serve piu.
  • Aggiornare README e documentazione deploy.
  • Verificare docker compose up --build -d da ambiente pulito.

Rischi aperti

  • Redis non e fonte durevole come Postgres: non salvare solo su Redis stati importanti.
  • Visibility timeout troppo basso puo causare riesecuzioni.
  • Visibility timeout troppo alto ritarda il recupero di task persi.
  • Parsing DJI deve restare idempotente per tollerare redelivery e retry.
  • Celery aumenta complessita operativa: healthcheck, log, tuning concurrency, memory leak del subprocess.