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. Ejercicio 1 — El worker local
  2. Ejercicio 2 — La expiración
  3. Ejercicio 3 — on_commit
  4. Ejercicio 4 — Idempotencia
  5. Ejercicio 5 — Clasificación
  6. Resumen del profesor

Ejercicio 1 — El worker local

  1. Con acks_late=True, la evidencia del log:
[Worker-1] Task reservations.tasks.lenta[job-42] received
(kill -9 del worker a los 5 s)
[Worker-2] Task reservations.tasks.lenta[job-42] received     <- vuelve a la cola
[Worker-2] Task reservations.tasks.lenta[job-42] succeeded

Sin acks_late, el mensaje se consume al recibirlo y la muerte del worker lo pierde silenciosamente — por eso el par con reject_on_worker_lost. Nota: la tarea debe ser idempotente (la segunda ejecución repite el efecto de la primera incompleta).

Ejercicio 2 — La expiración

  1. La protección de triple capa contra solapamiento: beat único (un solo contenedor con celery beat), skip_locked en la query (el segundo worker no ve las filas que el primero tiene lockeadas) y la dedup opcional por minuto con Redis (lock:expirar:{minute}, cache.add, TTL 60 s — la 12). Evidencia esperada en logs:
[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: el test corre la tarea EN el proceso del test (sin Redis ni worker) — ideal para 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()          # llama SystemClock: para eager, monkeypatch o inyección
    self.assertEqual(refs, ["RF-8A21"])

Ojo con EAGER_PROPAGATES=True: sin él, las excepciones de la tarea se tragan y el test pasa en falso — el error del eager debe romper el test, como rompería el worker.

Ejercicio 3 — on_commit

  1. El test rojo:
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)  # duplicado: rollback
            transaction.on_commit(lambda: None)  # (antes: .delay() directo aquí)
    self.assertEqual(len(self.mailbox), 1)   # ROJO: el email salió de una reserva inexistente
  1. Con on_commit, el rollback no dispara el callback (se registra pero la transacción abortó) y el mailbox queda en 0; en el test del commit feliz, transaction.on_commit se ejecuta al salir del bloque y el mailbox marca 1. En tests no-atómicos (TestCase envuelve cada test en transacción) on_commit NO dispara: usa self.captureOnCommitCallbacks(execute=True) (Django ≥ 3.2) o TransactionTestCase — el detalle que hace perder una tarde.
  1. El grep final vacío es la regla del proyecto: .delay() solo dentro de on_commit (o fuera de toda transacción, raro y documentado).

Ejercicio 4 — Idempotencia

  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(...)          # el fake de SMTP del test apunta a una lista

Tres lanzamientos → un email; el retorno "dedup" es la trazabilidad en el log del worker. El TTL de 7 días cubre reintentos razonables; después, un reenvío sería un caso de soporte, no un bug del worker.

  1. El intent.id como Idempotency-Key hace el trabajo SUCIO en la pasarela: la repetición devuelve el cargo original (200, mismo body) sin cobrar dos veces — tu tarea puede ser invocada N veces, el dinero se mueve una. La combinación dedup-local + idempotency-key remota es la postura estándar: la local es optimización, la remota es la garantía.

Ejercicio 5 — Clasificación

TareaCríticaLatencia toleradaIdempotenciaReintentos
Email confirmaciónmedia (el usuario la espera, pero hay reenvío)1 mindedup Redis por ref5 + backoff
Email recordatorio 24hbaja1 hdedup por ref+kind+fecha3
Webhooks salientes (17)alta (contrato con terceros)5 mindedup por event_idbackoff 1m→6h + DLQ
Expiraciónalta (libera inventario)1 minskip_locked + lock por minutoinfinito + alerta
Liquidación mensualalta (dinero)24 hlote idempotente (31)infinito + alerta
DisponibilidadNUNCA en cola0 (el usuario espera)——

El retry infinito con alerta:

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 alerta si > 0
    ...

El criterio: desechable = su pérdida molesta; crítica = su pérdida rompe contrato o dinero — y las críticas se vigilan por ausencia (heartbeat: "no corrió en 5 min") tanto como por fallo.


Resumen del profesor

  • La cola separa lo esperado de lo consecuente; el criterio de admisión (latencia-tolerante, repetible, no-en-la-respuesta) decide.
  • acks_late + reject_on_worker_lost = at-least-once real, y por eso cada tarea nace idempotente (dedup local + Idempotency-Key remota).
  • on_commit siempre; eager con EAGER_PROPAGATES en tests; backoff con jitter para no tumbarnos a nosotros mismos al volver.