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

Lección 30 — Brokers de mensajería

RabbitMQ, Kafka y SQS: entrega at-least-once, orden y dead-letter queues.

Publicada
En esta lección
  1. Ejercicio 1 — Consumer groups
  2. Ejercicio 2 — Dedup y orden
  3. Ejercicio 3 — DLQ operada
  4. Ejercicio 4 — Matriz de decisión
  5. Ejercicio 5 — Comando vs evento
  6. Resumen del profesor

Ejercicio 1 — Consumer groups

  1. 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)
  1. Sin ACK, el mensaje queda en el PEL del consumidor muerto; XAUTOCLAIM ticketflow:events notificaciones worker-2 60000 0-0 lo 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.
  1. El poller con skip_locked:
python
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

  1. El test (el mismo de la 25, ahora en el consumidor del stream):
python
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)
  1. El guard de orden:
python
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.

  1. Por reserva y no global: el orden SOLO importa entre eventos de la misma reserva; particionar por ref paralleliza (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

  1. La política en el worker (Celery):
python
@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".

  1. El runbook de alerta > 50:
  1. Congelar: no re-encolar masivamente (repite el mismo fallo × 50).
  2. Agrupar por cause: la DLQ con 50 mensajes casi nunca son 50 problemas — son 1 problema × 50 mensajes.
  3. Identificar el deploy/cambio de hace ~2 h (el mensaje más viejo data la causa).
  4. Fix + test de regresión; deploy.
  5. dlq_replay --since 2h (sin dry-run) y vigilar dlq_depth → 0 y task_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):

BrokerGarantíaOrdenReplayOperaciónCoste mensual estimado
Redis listasat-least-onceFIFO por listanotrivial~0 (ya lo pagas)
Redis streamsat-least-once + PELpor streamsí (con XREAD desde id)trivial-moderada~0
RabbitMQat-least-onceFIFO por colano (dead-letter sí)moderadaVPS ~20 €
Kafkaat-least-once (EOS intra-app)por particiónsí (retención configurable)alta3 nodos ~90-150 €
SQSat-least-once (FIFO opc.)por MessageGroupId (FIFO)nonula~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

  1. Broker de tareas (comandos): enviar_email, cobrar_intent, expirar_reservas. Stream de eventos: ReservationConfirmed, PaymentFailed, EventPublished.
  1. Lo que se rompe conceptualmente: enviar_email como 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ó?). Y ReservationConfirmed como 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.