🎯 Meta: el apéndice B te dio colas de tareas con Redis — perfectas para "procesa esto en segundo plano". Este apéndice cubre el escalón siguiente: cuando varios servicios independientes necesitan reaccionar al mismo evento, cuando el volumen es masivo, o cuando necesitas reproducir el historial de lo que pasó. Ahí es donde Redis se queda corto y entran Kafka o RabbitMQ.
Versiones: Apache Kafka 4.0 (KRaft, sin ZooKeeper) · RabbitMQ 4.1.
⚠️ Aviso de complejidad (como con Kubernetes, ap. G): Kafka y RabbitMQ son infraestructura con estado, operación y curva de aprendizaje reales. La mayoría de proyectos de este libro nunca necesitarán salir de las colas de Redis del apéndice B. Este apéndice existe para cuando sí las necesites — arquitecturas de microservicios, integraciones entre equipos, volúmenes de eventos que Redis no está diseñado para sostener.
S.1 · Tres herramientas, tres problemas distintos
Redis (ap. B) RabbitMQ Kafka
───────────── ──────── ─────
Cola simple Broker de mensajes Log de eventos distribuido
"procesa esto" con enrutado flexible "esto pasó, en este orden"
Un consumidor Colas, exchanges, routing Particiones, consumer groups
se lo lleva y listo keys — mensajería empresarial replay del historial completo
No hay historial Historial corto (colas) Historial completo (retención configurable)| Necesitas | Herramienta |
|---|---|
| Enviar un email en segundo plano | Redis (ap. B) — no compliques esto |
| Un pedido dispara: factura + email + actualizar inventario, cada uno un servicio distinto | RabbitMQ o Kafka |
| Varios equipos/servicios necesitan reaccionar al MISMO evento sin acoplarse | Kafka |
| Reconstruir el estado actual reproduciendo TODO el historial de eventos | Kafka (event sourcing) |
| Enrutado complejo (por prioridad, por tipo, con reintentos y dead-letter) | RabbitMQ |
| Streams de datos de altísimo volumen (analítica, IoT, logs) | Kafka |
🧠 La diferencia conceptual clave: en una cola tradicional (Redis, RabbitMQ), un mensaje se consume y desaparece — un solo consumidor se lo lleva. En Kafka, un evento se publica en un log y múltiples consumidores independientes lo leen a su propio ritmo, sin borrarlo. Es la diferencia entre "buzón de tareas" y "libro de registro público".
S.2 · RabbitMQ: el enrutador de mensajes flexible
El modelo mental: exchanges, colas y bindings
Productor ──▶ Exchange ──(routing key + binding)──▶ Cola(s) ──▶ Consumidor(es)Un exchange decide a qué cola(s) va cada mensaje. Tres tipos cubren el 95% de los casos:
direct: el mensaje va a la cola cuya routing key coincide EXACTA
"pedido.creado" → solo a quien escuche "pedido.creado"
topic: coincidencia con comodines (# = cualquier cosa, * = una palabra)
"pedido.*" escucha "pedido.creado" Y "pedido.cancelado"
fanout: el mensaje va a TODAS las colas conectadas (broadcast)
perfecto para "avisa a todo el que esté escuchando, sin filtrar"docker run -d -p 5672:5672 -p 15672:15672 rabbitmq:4.1-management # :15672 = panel web💡 Colas quorum, no clásicas, desde RabbitMQ 4. Las colas quorum (basadas en Raft, replicadas entre nodos) son el tipo recomendado por defecto — las colas clásicas espejadas (mirrored) están deprecadas. Para un solo nodo en desarrollo apenas notas la diferencia, pero declara tus colas como quorum desde el día 1 si algún día corres RabbitMQ en clúster:
arguments: { 'x-queue-type': 'quorum' }al crear la cola.
// Productor (NestJS + amqplib)
import amqp from 'amqplib';
const conn = await amqp.connect(process.env.RABBITMQ_URL!);
const canal = await conn.createChannel();
await canal.assertExchange('pedidos', 'topic', { durable: true });
function publicarEvento(routingKey: string, payload: object) {
canal.publish(
'pedidos',
routingKey, // ej. 'pedido.creado'
Buffer.from(JSON.stringify(payload)),
{ persistent: true }, // sobrevive a un reinicio del broker
);
}
// Al crear un pedido:
publicarEvento('pedido.creado', { pedidoId, userId, total });// Consumidor — servicio de FACTURACIÓN, independiente del que crea el pedido
const canal = await conn.createChannel();
await canal.assertExchange('pedidos', 'topic', { durable: true });
const { queue } = await canal.assertQueue('facturacion.pedidos', { durable: true });
await canal.bindQueue(queue, 'pedidos', 'pedido.creado'); // solo escucha ESTE evento
canal.consume(queue, async (msg) => {
if (!msg) return;
const evento = JSON.parse(msg.content.toString());
try {
await generarFactura(evento.pedidoId);
canal.ack(msg); // confirma que se procesó — se borra de la cola
} catch (err) {
canal.nack(msg, false, true); // falló: vuelve a la cola para reintentar
}
});⚠️
ackmanual, no automático. Sinackexplícito tras procesar con éxito, un consumidor que muere a mitad de proceso pierde el mensaje para siempre (con auto-ack) o lo reprocesa mal. El patrón correcto: procesa → si va bien,ack; si falla,nack(vuelve a la cola o va a una dead-letter queue tras N reintentos — igual que el patrón de reintentos del ap. B.4).
Dead-letter queue: qué hacer con lo que falla siempre
await canal.assertQueue('facturacion.pedidos', {
durable: true,
arguments: {
'x-dead-letter-exchange': 'pedidos.fallidos', // tras agotar reintentos, va aquí
'x-message-ttl': 30000, // 30s antes de considerarlo "atascado"
},
});💡 Sin dead-letter queue, un mensaje que siempre falla (bug en tu código, dato corrupto) reintenta para siempre, consumiendo recursos sin avanzar nunca. La DLQ es tu "cajón de revisión manual" — el mismo concepto que las tareas fallidas del ap. B.4.
S.3 · Kafka: el log de eventos distribuido
El modelo mental: topics, particiones y offsets
Topic "pedidos" dividido en 3 particiones (paralelismo + orden garantizado POR partición):
Partición 0: [evento1][evento2][evento5]...
Partición 1: [evento3][evento6]... ← cada evento tiene un OFFSET (posición)
Partición 2: [evento4][evento7]...
Los eventos del MISMO pedido van a la MISMA partición (por su "key") → orden garantizado
entre eventos de ese pedido; entre particiones distintas, NO hay orden garantizado.docker run -d -p 9092:9092 apache/kafka:4.0// Productor (kafkajs)
import { Kafka } from 'kafkajs';
const kafka = new Kafka({ clientId: 'tienda-api', brokers: ['localhost:9092'] });
const productor = kafka.producer();
await productor.connect();
await productor.send({
topic: 'pedidos',
messages: [{
key: pedidoId, // MISMA key → MISMA partición → orden
value: JSON.stringify({ tipo: 'creado', pedidoId, userId, total }),
}],
});// Consumidor — cada consumer GROUP recibe TODOS los eventos, independiente de otros grupos
const consumidor = kafka.consumer({ groupId: 'servicio-facturacion' });
await consumidor.connect();
await consumidor.subscribe({ topic: 'pedidos', fromBeginning: false });
await consumidor.run({
eachMessage: async ({ message, partition }) => {
const evento = JSON.parse(message.value!.toString());
await generarFactura(evento.pedidoId);
// el offset se confirma automáticamente al terminar (o gestiónalo tú para más control)
},
});🧠 Consumer groups son la clave de la escalabilidad de Kafka: dentro de un mismo
groupId, cada partición la procesa un solo consumidor (se reparten el trabajo — paralelismo real). Pero otrogroupId(ej.servicio-analitica) recibe los mismos eventos de nuevo, desde cero — cada grupo tiene su propio puntero (offset) independiente. Así es como 5 servicios distintos reaccionan al mismo evento sin pisarse ni coordinarse entre ellos — el desacoplamiento real que RabbitMQ no da tan naturalmente.
Replay: la superpotencia que Redis y RabbitMQ no tienen
// Kafka retiene los eventos (días, semanas, o para siempre — configurable por topic)
// Un consumidor NUEVO puede reprocesar TODO el historial desde el principio:
await consumidor.subscribe({ topic: 'pedidos', fromBeginning: true });💡 Casos reales donde esto importa: lanzas un nuevo servicio de analítica y necesita "ver" todos los pedidos de los últimos 2 años para construir su estado inicial; encuentras un bug en el servicio de facturación y necesitas reprocesar los últimos 3 días. Con Redis/RabbitMQ, ese historial ya no existe — se consumió y se borró.
S.4 · Exactly-once: cuando "al menos una vez" no es suficiente
Por defecto, Kafka (como la mayoría de sistemas de mensajería) da garantías "al menos una vez": si un productor no recibe confirmación de un envío, reintenta — y eso puede duplicar el mensaje. Para la mayoría de casos, un consumidor idempotente (S.2, S.5) basta. Pero cuando el propio conteo o suma de eventos importa (facturación, contabilidad, métricas exactas), Kafka ofrece exactly-once semantics (EOS) con dos piezas:
// 1. Productor IDEMPOTENTE — Kafka deduplica reintentos automáticamente por número de secuencia
const productor = kafka.producer({ idempotent: true }); // por defecto en kafkajs reciente
// 2. Transacciones — escritura atómica en VARIAS particiones/topics a la vez
const productor = kafka.producer({ transactionalId: 'servicio-facturacion-1', idempotent: true });
async function procesarPedidoPagado(evento: PedidoPagado) {
const transaccion = await productor.transaction();
try {
await transaccion.send({ topic: 'facturas', messages: [{ value: JSON.stringify(generarFactura(evento)) }] });
await transaccion.send({ topic: 'metricas-ingresos', messages: [{ value: JSON.stringify({ total: evento.total }) }] });
await transaccion.commit(); // AMBOS mensajes visibles a la vez, o NINGUNO
} catch (err) {
await transaccion.abort(); // si algo falla, ninguno de los dos queda visible
throw err;
}
}🧠
transactional.ides la clave para recuperarse de caídas. Si el proceso muere a mitad de una transacción y se reinicia con el MISMOtransactionalId, Kafka reconoce esa identidad, cierra ("fencea") cualquier transacción anterior a medio terminar de ese productor, y evita que un productor "zombi" (el proceso viejo que en realidad sigue vivo) escriba datos duplicados a la vez que el nuevo. SintransactionalId, no hay forma de detectar ese escenario.
⚠️ EOS no es gratis: añade latencia (2-5 ms por commit de transacción) y reduce el throughput (~10-20%) frente a "al menos una vez". Resérvalo para donde un duplicado o una pérdida sea un problema de negocio real (dinero, conteos auditables) — para logs de analítica o notificaciones, la idempotencia del consumidor (S.5, outbox) suele bastar y es más simple.
🔗 El consumidor también participa: para EOS de punta a punta, el consumidor debe leer con
isolation.level: read_committed(ignora mensajes de transacciones abortadas) y, en el patrón consume-transform-produce, confirmar sus offsets dentro de la misma transacción que produce sus resultados — así una relectura tras un fallo no reprocesa ni pierde nada.
S.5 · Patrones de arquitectura que habilitan
Event-driven entre microservicios (desacoplamiento real)
┌──────────────┐
Servicio │ │ Servicio de
de Pedidos ────▶│ Kafka: topic │◀──── Facturación (consumer group A)
(productor) │ "pedidos" │◀──── Notificaciones (consumer group B)
│ │◀──── Analítica (consumer group C)
└──────────────┘El servicio de Pedidos no sabe (ni le importa) cuántos servicios escuchan sus eventos. Añadir un consumidor nuevo (ej. "Fraude") no requiere tocar el servicio de Pedidos — solo suscribirse al topic. Es el mismo principio de acoplamiento del cap. 11, llevado a nivel de sistema completo.
Event sourcing (avanzado — úsalo con criterio)
En vez de guardar solo el estado actual ("el pedido está entregado"), guardas cada evento que llevó ahí (creado → pagado → en_preparacion → entregado) y el estado actual se deriva reproduciendo la secuencia. Potente para auditoría total, pero añade complejidad real — no lo adoptes "porque suena bien"; el CRUD normal (caps. 03-08) resuelve la mayoría de dominios sin esto.
Outbox pattern: el problema que resuelve la consistencia entre BD y broker
❌ Sin outbox: riesgo de inconsistencia
1. Guardas el pedido en PostgreSQL ✅
2. Publicas el evento en Kafka ❌ (el servidor cae justo aquí)
→ el pedido existe en tu BD pero NADIE se enteró — facturación nunca corre
✅ Con outbox: consistencia garantizada
1. En LA MISMA transacción de BD: guarda el pedido + guarda una fila en "outbox"
2. Un proceso aparte lee la tabla outbox y publica a Kafka, marcando como enviado
→ si el proceso 2 falla, reintenta — el evento nunca se pierde porque vive en tu BD-- Tabla outbox: parte de la MISMA transacción que crea el pedido (cap. 01, transacciones)
CREATE TABLE outbox (
id UUID PRIMARY KEY DEFAULT gen_random_uuid(),
topic TEXT NOT NULL,
payload JSONB NOT NULL,
publicado BOOLEAN DEFAULT FALSE,
creado_en TIMESTAMPTZ DEFAULT now()
);🧠 Es el mismo problema de "dos sistemas que deben ponerse de acuerdo" que resolvimos con idempotencia en el apéndice P (Stripe) y el cap. 20 (refresh tokens) — aquí la solución es apoyarte en la transacción ACID de PostgreSQL como ancla de verdad, y publicar al broker después, de forma reintentable.
S.6 · Observabilidad y operación (repaso del apéndice F)
□ Lag del consumidor: ¿cuántos mensajes sin procesar tiene cada consumer group?
(métrica CRÍTICA — un lag creciente = tu consumidor no da abasto, cap. F.4)
□ Dead-letter queue con alerta (ap. F.7): mensajes ahí = algo se está rompiendo en silencio
□ Retención del topic configurada a propósito (días/GB), no "para siempre" por defecto
□ Particiones suficientes para el paralelismo que necesitas (no se reducen después fácilmente)# Kafka: inspeccionar el lag de un consumer group
kafka-consumer-groups.sh --bootstrap-server localhost:9092 \
--group servicio-facturacion --describeS.7 · Buenas prácticas
- Redis (ap. B) por defecto. Sube a RabbitMQ/Kafka solo cuando el problema real (desacoplar servicios, replay, volumen masivo) lo justifique — misma regla de "no sobre-ingenierizar" que Kubernetes.
- RabbitMQ para enrutado flexible entre servicios; Kafka para streams de eventos con necesidad de replay o múltiples consumidores independientes.
ackmanual y dead-letter queue siempre — nunca dejes que un mensaje fallido reintente para siempre en silencio.- Outbox pattern cuando el evento debe ser consistente con un cambio en tu base de datos.
- Monitoriza el lag del consumidor como monitorizas la latencia de tu API (ap. F).
- Los eventos son un contrato entre servicios — versiónalos igual que versionas tu API REST (ap. D.7); cambiar su forma sin avisar rompe a todos los consumidores.
- No mezcles Kafka con tu cola de tareas simple. Enviar un email no necesita un log de eventos distribuido — necesita Redis (ap. B).
✅ Ejercicio del apéndice
1. RabbitMQ: al pagar un pedido en la Cantina (cap. 22), publica 'pedido.pagado'
(topic exchange) y crea DOS consumidores independientes: facturación y
notificaciones. Añade dead-letter queue con 3 reintentos.
2. Kafka: monta un topic 'eventos-cantina' con 3 particiones (key = pedidoId).
Publica creado/pagado/preparado/entregado. Crea 2 consumer groups distintos
(analítica y auditoría) y demuestra que cada uno procesa TODO el historial
de forma independiente.
3. Implementa el outbox pattern: la tabla outbox en la MISMA transacción que
crea el pedido, y un worker aparte que la vacía hacia Kafka con reintentos.
4. Provoca el fallo del outbox worker a mitad de proceso (mátalo) y demuestra
que al reiniciarlo, el evento pendiente SÍ se publica (no se perdió).
5. Métrica de lag del consumer group expuesta en Prometheus (ap. F) + alerta
si supera un umbral.
6. Productor transaccional: publica en 'facturas' y 'metricas-ingresos' en la
MISMA transacción de Kafka. Fuerza un error a mitad y comprueba que NINGUNO
de los dos mensajes queda visible (abort).
7. Reflexiona por escrito: ¿esta parte del proyecto REALMENTE necesitaba Kafka,
o las colas de Redis del apéndice B habrían bastado? ¿Necesitabas exactly-once
de verdad, o "al menos una vez" + idempotencia ya resolvía el problema?🧠 Autoevaluación
En una cola, un mensaje se consume y desaparece — un solo consumidor se lo lleva. En Kafka, el evento se publica en un log persistente que múltiples consumer groups pueden leer de forma independiente, cada uno con su propio offset, sin que unos afecten a otros ni el evento se borre al leerse.
Ack/nack garantizan que el consumidor procese el mensaje de forma fiable una vez publicado. El outbox resuelve el paso ANTERIOR: garantizar que el evento se publique si (y solo si) la transacción de base de datos que lo originó tuvo éxito — evitando el caso donde el pedido se guarda pero el evento nunca sale por una caída justo entre los dos pasos.
No. Cada groupId mantiene su propio offset independiente y recibe TODOS los mensajes del topic, sin importar lo que hagan otros grupos. Dentro de UN MISMO grupo sí se reparten las particiones para paralelizar el trabajo.
Cuando el problema es "procesa esto en segundo plano" con un solo consumidor y sin necesidad de que otros servicios reaccionen al mismo evento — ahí las colas de Redis (ap. B) resuelven todo con muchísima menos complejidad operativa.
Permite a Kafka reconocer al mismo productor tras un reinicio y "fencear" (invalidar) cualquier transacción anterior a medio terminar de ese mismo id — evitando que un proceso "zombi" (que en realidad sigue vivo tras perder la conexión) escriba datos duplicados a la vez que la instancia nueva. Sin transactionalId, Kafka no puede distinguir esos dos escritores.
EOS añade latencia (varios milisegundos por commit de transacción) y reduce el throughput frente a "al menos una vez" — un coste real que solo se justifica cuando un duplicado o una pérdida es un problema de negocio (dinero, conteos auditables). Para logs o notificaciones, un consumidor idempotente sobre "al menos una vez" resuelve lo mismo con menos complejidad operativa.
Volver al: README.md · Relacionado: B-redis-cache-colas.md, 11-arquitectura.md, F-observabilidad.md