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
- Comparar Redis, RabbitMQ, Kafka y SQS con las garantías reales (entrega, orden, retención, precio de operación).
- Manejar el at-least-once con deduplicación y el orden con claves de partición sin creer en exactly-once mágico.
- Diseñar dead-letter queues con política de re-procesado humano, no un cementerio que nadie mira.
1. Los cuatro candidatos, sin marketing
| Broker | Modelo | Entrega | Orden | Retención | Operación |
|---|---|---|---|---|---|
| Redis (listas/streams) | push en memoria | at-least-once con ACKs | por cola | baja (memoria; persistencia AOF/RDB opcional) | trivial, ya lo tienes (12) |
| RabbitMQ | push, exchanges + routing | at-least-once con acks | FIFO por cola (con cuidado: consumidores compiten) | hasta ACK, TTLs, DLX | moderada; gestión de colas |
| Kafka | pull, log particionado | at-least-once (transactions: EOS por app) | por partición, con clave | días/semanas replay | alta (ZK/KRaft, rebalances) |
| SQS | pull gestionado | at-least-once (FIFO: dedup+orden por grupo) | estándar: no; FIFO: por MessageGroupId | hasta 14 días | nula (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):
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étricadlq_depth_total(46) y alerta si crece. - Re-procesado con humano: un comando
manage.py dlq_replay --since 2h --dry-runque 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 5Reglas 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
- ¿Qué umbrales concretos justifican migrar de Redis a RabbitMQ/Kafka en TicketFlow y por qué no es fracaso?
- ¿Por qué "exactly-once" no existe entre sistemas y qué dos mecanismos lo sustituyen (con las lecciones donde ya los usaste)?
- En Kafka, ¿qué rompe el orden con múltiples consumidores y qué patrón lo salva por agregado?
- Diferencia entre comando (tarea) y evento: ¿qué canal lleva cada uno en TicketFlow y por qué no se mezclan?
- ¿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.