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. Objectives
  2. 1. Why a queue
  3. 2. Celery on TicketFlow: the minimal honest configuration
  4. 3. The expiration task (28's logic, the new transport)
  5. 4. Idempotent tasks and on_commit
  6. 5. What goes in the queue and what doesn't
  7. Self-assessment

Stack: Django/DRF · Project: TicketFlow Status: Published — opening the async processing module Prerequisite: Lesson 28 — Refactoring without breaking anything


Objectives

  1. Migrate the expiration job to Celery: the logic already extracted (28) gains a transport with retries and supervision.
  2. Write tasks that are idempotent and transaction-safe: on_commit so no ghosts get published.
  3. Configure retries with backoff and decide which task is critical (retry forever) and which is disposable.

1. Why a queue

The synchronous checkout pays inside the request: gateway (2 s) + confirmation (50 ms) + email (300 ms, if SMTP is having a good day) + outgoing webhooks (1–3 s). The user waits for all of that, and if the email fails, does the checkout fail? No: the email is not part of the purchase contract. The queue separates "what the user waits for" from "what must happen afterwards": the request does transaction + outbox (25) and responds; the workers do the rest. And the periodic tasks (expire reservations every minute) stop being a human cron.

The cost (there always is one): tasks fail and repeat — your async code must be idempotent and tolerate weird ordering (30 goes deeper into the broker's guarantees).

2. Celery on TicketFlow: the minimal honest configuration

python
# core/celery.py
import os
from celery import Celery

os.environ.setdefault("DJANGO_SETTINGS_MODULE", "ticketflow.settings")
app = Celery("ticketflow")
app.config_from_object("django.conf:settings", namespace="CELERY")
app.autodiscover_tasks()

# settings.py (all 12-factor, 27)
CELERY_BROKER_URL = env("CELERY_BROKER_URL")            # redis:// or amqp://
CELERY_TASK_ACKS_LATE = True                            # ack after processing, not on receipt
CELERY_TASK_REJECT_ON_WORKER_LOST = True                # if the worker dies, the task returns to the queue
CELERY_TASK_TIME_LIMIT = 120                            # no task lives longer than 2 minutes
CELERY_WORKER_PREFETCH_MULTIPLIER = 1                   # fair share, no hoarding

acks_late + reject_on_worker_lost = real at-least-once (the task re-executes if the worker dies halfway). Immediate consequence: every task must tolerate running twice. prefetch_multiplier=1 with variable-duration tasks prevents one worker hoarding 20 slow tasks while another idles.

3. The expiration task (28's logic, the new transport)

python
# reservations/tasks.py
@shared_task(bind=True, max_retries=5, default_retry_delay=10, autoretry_for=(OperationalError,),
             retry_backoff=True, retry_jitter=True)
def expirar_reservas(self):
    return expirar_reservas_service(SystemClock())      # 28's stable logic

The logic didn't move: the task is transport. Backoff with jitter (10, 20, 40 s ± random) prevents 50 tasks failing due to a Postgres restart from all retrying at once and knocking it over as it comes back (12's stampede). autoretry_for for known transient failures; the unknown explodes and goes to failure visibility (46).

The periodic task with celery beat (the scheduler):

python
app.conf.beat_schedule = {
    "expirar-reservas-cada-minuto": {
        "task": "reservations.tasks.expirar_reservas",
        "schedule": crontab(minute="*/*1"),
    },
}

Two workers running but only ONE running beat (beat is single-writer); and the task itself carries 10/28's skip_locked: if two overlap due to lag, they don't step on each other.

4. Idempotent tasks and on_commit

The classic trap: enqueuing a task INSIDE a transaction:

python
with transaction.atomic():
    r.confirm()
    enviar_email_confirmacion.delay(r.public_ref)   # BUG: may run BEFORE the commit

If the transaction rolls back after enqueuing, the email goes out for a reservation that doesn't exist. The rule: every task is enqueued in transaction.on_commit (16 already used it for SSE):

python
with transaction.atomic():
    r.confirm()
    transaction.on_commit(lambda: enviar_email_confirmacion.delay(r.public_ref))

And idempotency on the task's side (the worker is at-least-once, 25/30): the email task stores email:{ref}:{kind} in Redis with cache.add (25) and skips if it already exists; the charge task uses intent.id as the gateway's idempotency key (14). Idempotency at two levels: enqueue-once (on_commit) and process-once (dedup by business key).

5. What goes in the queue and what doesn't

Queue admission criteria: (1) it is not part of the response the user waits for; (2) it can repeat without disaster (or has dedup); (3) it is latency-tolerant (an email 30 s late breaks nothing). Yes: emails, outgoing webhooks (17), expirations, PDF generation, settlements. No: reading availability (the user waits), charging (the user waits and demands an answer now — 32's saga orchestrates the rest), auth.

And minimum monitoring from day one: celery -A ticketflow inspect active, the task_failed_total{task} counter (46) and an alert if expirar_reservas hasn't run in 5 minutes (46's heartbeat). A queue without observability is a black-hole mailbox.


Self-assessment

  1. Which criterion separates "this goes in the request" from "this goes in the queue"? Give the 3 admission-criteria tests.
  2. What do acks_late and reject_on_worker_lost do, and what obligation do they impose on your task code?
  3. The bug of enqueuing inside the transaction: which sequence triggers it, and how does on_commit eliminate it?
  4. Why does backoff carry jitter, and which concrete disaster does it avoid on a Postgres restart?
  5. Beat runs on a single worker: which other two mechanisms protect the periodic task from overlapping runs?

Continue with the exercises. The solutions only after trying it yourself.