Module 6 · Async processing and messaging

Lesson 29 — Queues and background jobs

Celery in TicketFlow: reservation expiry stops being a human cron.

Published
In this lesson
  1. Exercise 1 — The local worker
  2. Exercise 2 — The expiration
  3. Exercise 3 — on_commit
  4. Exercise 4 — Idempotency
  5. Exercise 5 — Classification
  6. Professor's summary

Exercise 1 — The local worker

  1. 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] succeeded

Without 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

  1. The triple-layer protection against overlap: single beat (one container with celery beat), skip_locked in 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
  1. Eager mode: the test runs the task IN the test process (no Redis, no worker) — ideal for CI:
python
@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

  1. The red test:
python
@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
  1. 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_commit executes on leaving the block and the mailbox shows 1. In non-atomic tests (TestCase wraps each test in a transaction) on_commit does NOT fire: use self.captureOnCommitCallbacks(execute=True) (Django ≥ 3.2) or TransactionTestCase — the detail that costs you an afternoon.
  1. The empty final grep is the project's rule: .delay() only inside on_commit (or outside any transaction, rare and documented).

Exercise 4 — Idempotency

  1. Dedup:
python
@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 list

Three 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.

  1. 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

TaskCriticalTolerated latencyIdempotencyRetries
Confirmation emailmedium (the user waits for it, but there is resending)1 minRedis dedup by ref5 + backoff
24h reminder emaillow1 hdedup by ref+kind+date3
Outgoing webhooks (17)high (third-party contract)5 mindedup by event_idbackoff 1m→6h + DLQ
Expirationhigh (frees inventory)1 minskip_locked + per-minute lockinfinite + alert
Monthly settlementhigh (money)24 hidempotent batch (31)infinite + alert
AvailabilityNEVER queued0 (the user waits)——

Infinite retry with alert:

python
@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.