Ejercicio 1 — Consumer groups
- La evidencia:
> XADD ticketflow:events * event_id e1 type ReservationConfirmed ref RF-8A21
1695800000000-0
> XGROUP CREATE ticketflow:events notificaciones $ 0
OK
(worker-1) XREADGROUP GROUP notificaciones worker-1 COUNT 10 STREAMS ticketflow:events ">"
→ recibe e1, e3 (worker-2 recibió e2: reparto, no duplicación)- Sin ACK, el mensaje queda en el PEL del consumidor muerto;
XAUTOCLAIM ticketflow:events notificaciones worker-2 60000 0-0lo transfiere (idle > 60 s) y el nuevo consumidor lo procesa y acka. El PEL es la "cola de vergüenza" que evita la pérdida silenciosa.
- El poller con skip_locked:
with transaction.atomic():
batch = list(OutboxEvent.objects.select_for_update(skip_locked=True)
.filter(publicado_at__isnull=True).order_by("id")[:50])
for ev in batch:
xadd("ticketflow:events", {"event_id": str(ev.id), ...})
ev.publicado_at = timezone.now()Dos pollers en paralelo: skip_locked hace que el segundo no vea las filas del primero (10) — sin duplicar ni pisar. Es el mismo lock del job de expiración (28): el patrón "batch con lock + marca" es la base de todo procesamiento en fondo sobre Postgres.
Ejercicio 2 — Dedup y orden
- El test (el mismo de la 25, ahora en el consumidor del stream):
def test_duplicado_un_efecto(self):
publicar(ev) # email fake: mailbox 1
publicar(ev) # at-least-once del broker: llega otra vez
self.assertEqual(len(mailbox), 1)- El guard de orden:
def aplicar(msg, estado):
if msg.updated_at < estado.updated_at:
return "obsoleto" # llegó tarde: el estado ya es más nuevo
estado.update(msg)Con eventos 10:00 → 10:05 → 10:02: el 10:02 se descarta (obsoleto) aunque llegó último. Es optimistic concurrency (10) aplicado a mensajería: el dato lleva su versión.
- Por reserva y no global: el orden SOLO importa entre eventos de la misma reserva; particionar por
refparalleliza (cada partición con su consumidor), particionar global serializa todo el sistema a 1 consumidor. Con 1M de reservas concurrentes y partición global: throughput = 1 consumidor = el cuello de la 55; por ref: throughput = nº de particiones.
Ejercicio 3 — DLQ operada
- La política en el worker (Celery):
@shared_task(bind=True, max_retries=5, retry_backoff=True, retry_jitter=True)
def procesar(self, payload):
try:
...
except PermanentError:
dlq_push({"task": "procesar", "payload": payload,
"cause": repr(self.request.last_exception or exc),
"trace_id": get_trace_id(), "attempts": self.request.retries})El test: tarea siempre fallida → el mock de dlq_push recibe 1 llamada con attempts=5; la cola principal queda vacía (no se atasca). El detalle: la excepción PERMANENTE (payload inválido) va directo a DLQ sin quemar 5 reintentos — el backoff es para fallos TRANSITORIOS (red, BD), no para "esto nunca va a funcionar".
- El runbook de alerta > 50:
- Congelar: no re-encolar masivamente (repite el mismo fallo × 50).
- Agrupar por
cause: la DLQ con 50 mensajes casi nunca son 50 problemas — son 1 problema × 50 mensajes. - Identificar el deploy/cambio de hace ~2 h (el mensaje más viejo data la causa).
- Fix + test de regresión; deploy.
dlq_replay --since 2h(sin dry-run) y vigilardlq_depth → 0ytask_failed_total(46).
El dry-run primero SIEMPRE: el re-procesado masivo sin fix es la definición de tsunami repetido.
Ejercicio 4 — Matriz de decisión
Volúmenes estimados de TicketFlow (~2 eventos de dominio por reserva, 50k reservas/día pico → 150k eventos/día; emails ~60k; webhooks salientes ~40k):
| Broker | Garantía | Orden | Replay | Operación | Coste mensual estimado |
|---|---|---|---|---|---|
| Redis listas | at-least-once | FIFO por lista | no | trivial | ~0 (ya lo pagas) |
| Redis streams | at-least-once + PEL | por stream | sí (con XREAD desde id) | trivial-moderada | ~0 |
| RabbitMQ | at-least-once | FIFO por cola | no (dead-letter sí) | moderada | VPS ~20 € |
| Kafka | at-least-once (EOS intra-app) | por partición | sí (retención configurable) | alta | 3 nodos ~90-150 € |
| SQS | at-least-once (FIFO opc.) | por MessageGroupId (FIFO) | no | nula | ~5-20 € a este volumen |
El ADR (extracto):
ADR-0011 — Broker de eventos: Redis Streams. Estado: aceptado. Contexto: ~150k eventos/día, equipo backend de 1, Redis ya operado para caché. Decisión: outbox → stream con consumer groups; Celery-Redis para tareas. Razones: replay básico, PEL para reclamos, cero infra nueva. Alternativa descartada: Kafka (replay superior, pero operación 3× y el volumen no lo exige). Revisión: si el volumen ×10 o la saga (32) exige retención semanal para reparaciones, reevaluar.
El patrón del ADR: decisión + razones + alternativa + umbral de revisión — el umbral convierte el ADR en documento vivo, no en lápida.
Ejercicio 5 — Comando vs evento
- Broker de tareas (comandos):
enviar_email,cobrar_intent,expirar_reservas. Stream de eventos:ReservationConfirmed,PaymentFailed,EventPublished.
- Lo que se rompe conceptualmente:
enviar_emailcomo evento dice "pasó que hay que enviar un email" — un hecho que no pasó, emitido por quien NO es la fuente del hecho: nadie puede auditar ni re-procesar con confianza (¿el email ya salió? ¿quién lo ordenó?). YReservationConfirmedcomo tarea asigna el conocimiento del consumidor al productor: el servicio de reservas tendría que saber que existe "el que procesa confirmaciones" — el acoplamiento que los eventos eliminan (la 32 lo explota). El comando tiene UN destinatario que espera el efecto; el evento tiene N suscriptores que reaccionan sin permiso del emisor.
Resumen del profesor
- El broker se elige por garantías + operación, no por README: streams de Redis cubren el volumen de TicketFlow con la operación que ya pagas.
- At-least-once + dedup por clave de negocio = exactly-once-efecto; el orden por agregado con guard de versión.
- DLQ operada: causa + trace_id + dry-run + runbook; una DLQ creciente es un incidente con un solo problema multiplicado.