Stack: Django/DRF · Proyecto: TicketFlow Estado: Publicada — apertura del módulo de procesamiento asíncrono Prerrequisito: Lección 28 — Refactorizar sin romper nada
Objetivos
- Migrar el job de expiración a Celery: la lógica ya extraída (28) gana un transporte con reintentos y supervisión.
- Escribir tareas idempotentes y seguras frente a la transacción:
on_commitpara no publicar fantasmas. - 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
# 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 acaparamientoacks_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)
# 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 28La 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):
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:
with transaction.atomic():
r.confirm()
enviar_email_confirmacion.delay(r.public_ref) # BUG: puede ejecutarse ANTES del commitSi 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):
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
- ¿Qué criterio separa "esto va en el request" de "esto va en la cola"? Da los 3 tests del criterio de admisión.
- ¿Qué hacen
acks_lateyreject_on_worker_lost, y qué obligación imponen sobre el código de tus tareas? - El bug de encolar dentro de la transacción: ¿qué secuencia lo dispara y cómo lo elimina
on_commit? - ¿Por qué el backoff lleva jitter y qué desastre concreto evita en un reinicio de Postgres?
- 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.