Módulo 6 · Procesamiento asíncrono y mensajería

Lección 29 — Colas y jobs en segundo plano

Celery en TicketFlow: la expiración de reservas deja de ser un cron humano.

Publicada
En esta lección
  1. Objetivos
  2. 1. Por qué una cola
  3. 2. Celery en TicketFlow: la configuración mínima honesta
  4. 3. La tarea de expiración (la lógica de la 28, el transporte nuevo)
  5. 4. Tareas idempotentes y on_commit
  6. 5. Qué va en la cola y qué no
  7. Autoevaluación

Stack: Django/DRF · Proyecto: TicketFlow Estado: Publicada — apertura del módulo de procesamiento asíncrono Prerrequisito: Lección 28 — Refactorizar sin romper nada


Objetivos

  1. Migrar el job de expiración a Celery: la lógica ya extraída (28) gana un transporte con reintentos y supervisión.
  2. Escribir tareas idempotentes y seguras frente a la transacción: on_commit para no publicar fantasmas.
  3. Configurar reintentos con backoff y decidir qué tarea es crítica (retry hasta el infinito) y cuál es desechable.

1. Por qué una cola

El checkout sincrónico paga en el request: gateway (2 s) + confirmación (50 ms) + email (300 ms, si SMTP va bien) + webhooks salientes (1–3 s). El usuario espera todo eso, y si el email falla, ¿el checkout falla? No: el email no es parte del contrato de la compra. La cola separa "lo que el usuario espera" de "lo que debe pasar después": el request hace transacción + outbox (25) y responde; los workers hacen el resto. Y las tareas periódicas (expirar reservas cada minuto) dejan de ser un cron humano.

El coste (siempre lo hay): las tareas fallan y se repiten — tu código asíncrono debe ser idempotente y tolerar orden raro (la 30 profundiza en las garantías del broker).

2. Celery en TicketFlow: la configuración mínima honesta

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 (todo 12-factor, la 27)
CELERY_BROKER_URL = env("CELERY_BROKER_URL")            # redis:// o amqp://
CELERY_TASK_ACKS_LATE = True                            # ack tras procesar, no al recibir
CELERY_TASK_REJECT_ON_WORKER_LOST = True                # si el worker muere, la tarea vuelve a la cola
CELERY_TASK_TIME_LIMIT = 120                            # ninguna tarea vive más de 2 minutos
CELERY_WORKER_PREFETCH_MULTIPLIER = 1                   # reparto justo, no acaparamiento

acks_late + reject_on_worker_lost = at-least-once de verdad (la tarea se re-ejecuta si el worker muere a mitad). Consecuencia inmediata: toda tarea debe tolerar ejecutarse dos veces. prefetch_multiplier=1 con tareas de duración variable evita que un worker acapare 20 tareas lentas mientras otro se aburre.

3. La tarea de expiración (la lógica de la 28, el transporte nuevo)

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())      # la lógica estable de la 28

La lógica no se movió: el task es transporte. El backoff con jitter (10, 20, 40 s ± aleatorio) evita que 50 tareas fallidas por un reinicio de Postgres reintenten todas a la vez y lo tumben al volver (el stampede de la 12). autoretry_for para los fallos transitorios conocidos; lo desconocido explota y va a la visibilidad del fallo (46).

La tarea periódica con celery beat (el scheduler):

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

Dos workers corriendo el beat solo UNO (beat es single-writer); y la tarea en sí lleva el skip_locked de la 10/28: si dos se solapan por un lag, no se pisan.

4. Tareas idempotentes y on_commit

La trampa clásica: lanzar una tarea DENTRO de una transacción:

python
with transaction.atomic():
    r.confirm()
    enviar_email_confirmacion.delay(r.public_ref)   # BUG: puede ejecutarse ANTES del commit

Si la transacción hace rollback después de encolar, el email sale de una reserva que no existe. La regla: toda tarea se encola en transaction.on_commit (la 16 ya lo usó para SSE):

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

Y la idempotencia del lado de la tarea (el worker es at-least-once, la 25/30): la tarea de email guarda email:{ref}:{kind} en Redis con cache.add (25) y salta si ya existe; la de cobro usa el intent.id como clave de idempotencia de la pasarela (14). Idempotencia en dos niveles: encolar-once (on_commit) y procesar-once (dedup por clave de negocio).

5. Qué va en la cola y qué no

Criterio de admisión en la cola: (1) no es parte de la respuesta que el usuario espera; (2) puede repetirse sin desastre (o tiene dedup); (3) es latencia-tolerante (un email 30 s tarde no rompe nada). Sí: emails, webhooks salientes (17), expiraciones, generación de PDFs, liquidaciones. No: leer disponibilidad (el usuario espera), cobrar (el usuario espera y exige respuesta en el momento — la saga de la 32 orquesta lo demás), auth.

Y el monitoreo mínimo desde el día uno: celery -A ticketflow inspect active, el contador task_failed_total{task} (46) y una alerta si expirar_reservas no corre en 5 minutos (el heartbeat de la 46). Una cola sin observabilidad es un buzón negro.


Autoevaluación

  1. ¿Qué criterio separa "esto va en el request" de "esto va en la cola"? Da los 3 tests del criterio de admisión.
  2. ¿Qué hacen acks_late y reject_on_worker_lost, y qué obligación imponen sobre el código de tus tareas?
  3. El bug de encolar dentro de la transacción: ¿qué secuencia lo dispara y cómo lo elimina on_commit?
  4. ¿Por qué el backoff lleva jitter y qué desastre concreto evita en un reinicio de Postgres?
  5. El beat corre en un solo worker: ¿qué otros dos mecanismos protegen la tarea periódica de ejecuciones solapadas?

Continúa con los ejercicios. Las solutions.md solo tras intentarlo.