learn.chetana.fr

Des jobs durables sans Celery

15 min de lectureAvancé

Le choix : pas de Celery, pas de broker

Une extraction dure 30 s : l'API répond 202 + job_id et lance un asyncio.create_task. Mais une task in-process meurt avec le pod (deploy, OOM, crash) — il faut la durabilité d'une queue sans en déployer une. La réponse du projet tient en trois pièces : un store de jobs Postgres, un payload rejouable, une boucle de réconciliation.

Pièce 1 : le JobsStore et son ContextVar

Un Protocol minimal (register_idempotent / get / get_result / update), deux implémentations : InMemoryJobsStore (dev/tests) et PgJobsStore (prod). Détail d'API élégant : la signature ne prend pas de tenant_id — le tenant voyage par ContextVar (job_tenant_scope(tenant_id)), que asyncio.create_task copie automatiquement dans la tâche enfant. Le Protocol reste stable et interchangeable, et le store Pg fail loud si le scope manque sur une écriture.

Pièce 2 : le payload rejouable (ADR-0018)

La règle : tout job persiste sa requête complète (ExtractionRequest.model_dump()) au moment de l'enregistrement, via un unique point d'entrée obligatoire (register_extraction_job). Conséquence magique : rejouer un job ne demande aucun join, aucune reconstruction — le payload EST la commande, prête à re-exécuter. États du job : pending → processing → completed | error (les deux derniers sont terminaux, intouchables).

Pièce 3 : le reconciler — la détection par staleness

Toutes les 60 s, la boucle de réconciliation cherche les orphelins. Comment sait-on qu'un job est mort ? Par l'immobilité de son updated_at dans un état non terminal :

processing depuis > 600 s sans update  → le pipeline a crashé quelque part
pending    depuis > 120 s              → la task n'a jamais démarré (pod tué juste après le 202)

Le claim réutilise exactement le pattern du polling email : FOR UPDATE SKIP LOCKED + bump immédiat, par tenant (session RLS), l'énumération cross-tenant passant par une session dédiée qui ne lit QUE les tenant_id. Job claimé → ExtractionRequest.model_validate(payload) → re-kick avec le même job_id. Payload absent (jobs d'avant la migration) ou invalide → marqué error avec warning, jamais de re-kick aveugle.

Les bornes de charge complètent le tableau : sémaphores 10 (matcher fast — protège le pool DB) et 5 (deep — protège les rate limits Anthropic/Cohere), batchs de claim à 10 — chaque ressource a sa borne, dimensionnée pour SON goulot.

🧩 Quiz1/4

Pourquoi persister le payload COMPLET du job (ADR-0018) ?

🃏 Flashcards1/4