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. Objetivos
  2. 1. Los cuatro candidatos, sin marketing
  3. 2. Las garantías: nadie da exactly-once
  4. 3. Redis Streams: el medio campo que ya tienes
  5. 4. Dead-letter queues con política
  6. 5. El diseño de mensajería de TicketFlow
  7. Autoevaluación

Stack: Redis / RabbitMQ / Kafka / SQS · Proyecto: TicketFlow Estado: Publicada — las garantías de cada broker, sin marketing Prerrequisito: Lección 29 — Colas y jobs en segundo plano


Objetivos

  1. Comparar Redis, RabbitMQ, Kafka y SQS con las garantías reales (entrega, orden, retención, precio de operación).
  2. Manejar el at-least-once con deduplicación y el orden con claves de partición sin creer en exactly-once mágico.
  3. Diseñar dead-letter queues con política de re-procesado humano, no un cementerio que nadie mira.

1. Los cuatro candidatos, sin marketing

BrokerModeloEntregaOrdenRetenciónOperación
Redis (listas/streams)push en memoriaat-least-once con ACKspor colabaja (memoria; persistencia AOF/RDB opcional)trivial, ya lo tienes (12)
RabbitMQpush, exchanges + routingat-least-once con acksFIFO por cola (con cuidado: consumidores compiten)hasta ACK, TTLs, DLXmoderada; gestión de colas
Kafkapull, log particionadoat-least-once (transactions: EOS por app)por partición, con clavedías/semanas replayalta (ZK/KRaft, rebalances)
SQSpull gestionadoat-least-once (FIFO: dedup+orden por grupo)estándar: no; FIFO: por MessageGroupIdhasta 14 díasnula (servicio), por mensaje

TicketFlow usa Redis como broker de Celery (29) porque ya corre para caché (12) y el volumen lo permite: ~50k mensajes/día. Los umbrales para migrar: cuando el AOF empiece a doler, o necesites replay (volver a consumir un día entero para reparar datos), o el orden por clave sea requisito (la saga de la 32 con eventos de la MISMA reserva debe procesarse en orden), Kafka/RabbitMQ entran en la conversación. Migrar broker no es fracasar: es el refactor del transporte (28) con la lógica en tareas estables.

2. Las garantías: nadie da exactly-once

Todos los brokers reales entregan at-least-once: si garantizan no-perder, a veces repiten (el ACK se perdió por red aunque el procesamiento fue exitoso). El exactly-once NO existe entre sistemas: existe exactamente-una-VEZ-EFECTO y lo construyes tú con deduplicación por clave de negocio (25: event:{event_id} con cache.add) o claves de idempotencia en el receptor (14: Idempotency-Key). Kafka ofrece "exactly-once semantics" — léase: transacciones Kafka-a-Kafka en una misma app; tu consumidor Django-Postgres sigue necesitando su dedup.

El orden: Kafka/RabbitMQ (con un solo consumidor por cola) prometen orden de entrega, pero con consumidores competindo (prefetch, varias replicas) el orden de PROCESAMIENTO se rompe igual. El patrón que funciona: el mensaje lleva sequence o updated_at; el consumidor descarta el viejo (UPDATE... WHERE updated_at < mensaje.updated_at — el optimistic concurrency de la 10) y las particiones/agrupaciones van por AGREGADO (reserva), no global (global = rendimiento 1, el cuello de la 55).

3. Redis Streams: el medio campo que ya tienes

Entre pub/sub (fire-and-forget, no sirve para esto) y Kafka hay Redis Streams — XADD/XREADGROUP con consumer groups, ACK y PEL (pending entries list):

bash
XADD ticketflow:events * event_id 01H... type ReservationConfirmed ref RF-8A21
XGROUP CREATE ticketflow:events notificaciones $ 0
XREADGROUP GROUP notificaciones worker-1 COUNT 10 BLOCK 5000 STREAMS ticketflow:events >
XACK ticketflow:events notificaciones <id>

Consumer group: cada mensaje lo procesa UN consumidor del grupo (fan-out entre servicios = grupos distintos), el ACK explícito y el PEL permite reclamar (XAUTOCLAIM) los pendientes de un consumidor muerto. Es Kafka en miniatura con la operación de Redis; para TicketFlow hoy: el outbox poller (25) publica a un stream y los consumidores (email, analytics) son grupos. Cuando el replay o el volumen justifiquen, Kafka toma el relevo con el mismo contrato.

4. Dead-letter queues con política

Toda cola necesita destino de errores: tras N reintentos fallidos, el mensaje a la DLQ — no se descarta en silencio ni se reintenta hasta el infinito atascando la cola. Celery: el recauchutado del custom state o la cola con task_reject_on_worker_lost y una tarea procesar_dlq que inspecciona. Lo que distingue una DLQ operada de un cementerio:

  • Clasificación: el mensaje a la DLQ se registra con causa (task, excepción, payload) — el evento entra a la métrica dlq_depth_total (46) y alerta si crece.
  • Re-procesado con humano: un comando manage.py dlq_replay --since 2h --dry-run que lista, luego aplica. El re-procesado NO automático: si falló 5 veces con backoff, la sexta solita no va a arreglarlo — alguien arregla el bug o decide descartar.
  • Causa raíz: la DLQ con profundidad > 0 es un incidente de la 47: el postmortem la apaga de verdad.

5. El diseño de mensajería de TicketFlow

Postgres (outbox)  →  poller  →  Redis Stream "events"
                                    ├── grupo "notificaciones"  → emails (29)
                                    ├── grupo "analytics"       → dashboards
                                    └── grupo "sagas"           → orquestador de la 32
Celery broker (Redis)  →  workers: tareas retry con backoff + DLQ tras 5

Reglas del proyecto: los eventos de dominio (25) salen del outbox al stream (garantía transaccional primero, broker después); las TAREAS (trabajo a hacer) van por el broker de Celery; los EVENTOS (hechos consumados) van por el stream. No mezclar: una tarea es un comando ("haz esto"), un evento es un hecho ("pasó esto") — la 32 monta la saga sobre eventos para desacoplar el orquestador del detalle de quién hace qué.


Autoevaluación

  1. ¿Qué umbrales concretos justifican migrar de Redis a RabbitMQ/Kafka en TicketFlow y por qué no es fracaso?
  2. ¿Por qué "exactly-once" no existe entre sistemas y qué dos mecanismos lo sustituyen (con las lecciones donde ya los usaste)?
  3. En Kafka, ¿qué rompe el orden con múltiples consumidores y qué patrón lo salva por agregado?
  4. Diferencia entre comando (tarea) y evento: ¿qué canal lleva cada uno en TicketFlow y por qué no se mezclan?
  5. ¿Qué distingue una DLQ operada de un cementerio? ¿Quién decide el re-procesado y por qué no es automático?

Continúa con los ejercicios. Las solutions.md solo tras intentarlo.