Exercise 1 — The local worker
- With
acks_late=True, the log evidence:
[Worker-1] Task reservations.tasks.lenta[job-42] received
(kill -9 of the worker at 5 s)
[Worker-2] Task reservations.tasks.lenta[job-42] received <- back in the queue
[Worker-2] Task reservations.tasks.lenta[job-42] succeededWithout acks_late, the message is consumed on receipt and the worker's death loses it silently — hence the pair with reject_on_worker_lost. Note: the task must be idempotent (the second run repeats the effect of the first, incomplete one).
Exercise 2 — The expiration
- The triple-layer protection against overlap: single beat (one container with
celery beat),skip_lockedin the query (the second worker doesn't see the rows the first has locked) and optional per-minute dedup with Redis (lock:expirar:{minute},cache.add, 60 s TTL — 12). Expected log evidence:
[beat] Scheduler: sending reservations.tasks.expirar_reservas
[worker] expirar_reservas[]: 1 reservas expiradas (refs: RF-8A21)
[worker] expirar_reservas[]: 0 reservas expiradas
[worker] expirar_reservas[]: 0 reservas expiradas- Eager mode: the test runs the task IN the test process (no Redis, no worker) — ideal for CI:
@override_settings(CELERY_TASK_ALWAYS_EAGER=True, CELERY_TASK_EAGER_PROPAGATES=True)
def test_expira_solo_la_vencida(self):
clock = FakeClock(datetime(2027, 3, 2, 10, 0, tzinfo=dt.UTC))
refs = expirar_reservas() # calls SystemClock: for eager, monkeypatch or inject
self.assertEqual(refs, ["RF-8A21"])Watch out for EAGER_PROPAGATES=True: without it, the task's exceptions get swallowed and the test passes falsely — the eager error must break the test, just as it would break the worker.
Exercise 3 — on_commit
- The red test:
@override_settings(CELERY_TASK_ALWAYS_EAGER=True)
def test_email_fantasma(self):
with self.assertRaises(DomainError):
with transaction.atomic():
reservar(self.user, self.event, ["A1", "A1"], clock=self.clock) # duplicate: rollback
transaction.on_commit(lambda: None) # (before: a direct .delay() here)
self.assertEqual(len(self.mailbox), 1) # RED: the email went out for a nonexistent reservation- With on_commit, the rollback doesn't fire the callback (it gets registered but the transaction aborted) and the mailbox stays at 0; in the happy-commit test,
transaction.on_commitexecutes on leaving the block and the mailbox shows 1. In non-atomic tests (TestCasewraps each test in a transaction) on_commit does NOT fire: useself.captureOnCommitCallbacks(execute=True)(Django ≥ 3.2) orTransactionTestCase— the detail that costs you an afternoon.
- The empty final grep is the project's rule:
.delay()only insideon_commit(or outside any transaction, rare and documented).
Exercise 4 — Idempotency
- Dedup:
@shared_task
def enviar_email_confirmacion(ref: str):
if not cache.add(f"email:{ref}:confirmed", "1", 60 * 60 * 24 * 7):
return f"dedup: {ref}"
send_mail(...) # the test's SMTP fake points at a listThree launches → one email; the "dedup" return is the traceability in the worker's log. The 7-day TTL covers reasonable retries; after that, a resend would be a support case, not a worker bug.
- The intent.id as Idempotency-Key makes the gateway do the DIRTY work: the repeat returns the original charge (200, same body) without charging twice — your task may be invoked N times, the money moves once. The local-dedup + remote-idempotency-key combination is the standard posture: the local one is an optimization, the remote one is the guarantee.
Exercise 5 — Classification
| Task | Critical | Tolerated latency | Idempotency | Retries |
|---|---|---|---|---|
| Confirmation email | medium (the user waits for it, but there is resending) | 1 min | Redis dedup by ref | 5 + backoff |
| 24h reminder email | low | 1 h | dedup by ref+kind+date | 3 |
| Outgoing webhooks (17) | high (third-party contract) | 5 min | dedup by event_id | backoff 1m→6h + DLQ |
| Expiration | high (frees inventory) | 1 min | skip_locked + per-minute lock | infinite + alert |
| Monthly settlement | high (money) | 24 h | idempotent batch (31) | infinite + alert |
| Availability | NEVER queued | 0 (the user waits) | — | — |
Infinite retry with alert:
@shared_task(bind=True, max_retries=None, retry_backoff=True, retry_jitter=True)
def expirar_reservas(self):
if self.request.retries >= 10:
metrics.incr("task_retry_total", tags={"task": "expirar_reservas"}) # 46 alerts if > 0
...The criterion: disposable = losing it is annoying; critical = losing it breaks a contract or money — and critical ones are watched by absence (heartbeat: "didn't run in 5 min") as much as by failure.
Professor's summary
- The queue separates the awaited from the consequential; the admission criterion (latency-tolerant, repeatable, not-in-the-response) decides.
- acks_late + reject_on_worker_lost = real at-least-once, and that is why every task is born idempotent (local dedup + remote Idempotency-Key).
- on_commit always; eager with EAGER_PROPAGATES in tests; backoff with jitter so we don't knock ourselves over on the way back up.