Ejercicio 1 — El worker local
- 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] succeededSin 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
- La protección de triple capa contra solapamiento: beat único (un solo contenedor con
celery beat),skip_lockeden 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- Eager mode: el test corre la tarea EN el proceso del test (sin Redis ni worker) — ideal para CI:
@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
- El test rojo:
@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- 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_commitse ejecuta al salir del bloque y el mailbox marca 1. En tests no-atómicos (TestCaseenvuelve cada test en transacción) on_commit NO dispara: usaself.captureOnCommitCallbacks(execute=True)(Django ≥ 3.2) oTransactionTestCase— el detalle que hace perder una tarde.
- El grep final vacío es la regla del proyecto:
.delay()solo dentro deon_commit(o fuera de toda transacción, raro y documentado).
Ejercicio 4 — Idempotencia
- Dedup:
@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 listaTres 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.
- 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
| Tarea | Crítica | Latencia tolerada | Idempotencia | Reintentos |
|---|---|---|---|---|
| Email confirmación | media (el usuario la espera, pero hay reenvío) | 1 min | dedup Redis por ref | 5 + backoff |
| Email recordatorio 24h | baja | 1 h | dedup por ref+kind+fecha | 3 |
| Webhooks salientes (17) | alta (contrato con terceros) | 5 min | dedup por event_id | backoff 1m→6h + DLQ |
| Expiración | alta (libera inventario) | 1 min | skip_locked + lock por minuto | infinito + alerta |
| Liquidación mensual | alta (dinero) | 24 h | lote idempotente (31) | infinito + alerta |
| Disponibilidad | NUNCA en cola | 0 (el usuario espera) | — | — |
El retry infinito con alerta:
@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.