En la lección anterior fijamos los contratos REST de TechCorp, pero dejamos claro que las llamadas síncronas serán pocas: Pedidos → Catálogo y Pedidos → Clientes. El resto del flujo "un cliente hace un pedido" (reservar stock, cobrar, confirmar, notificar) lo diseñamos en 02-05 como una saga por coreografía en la que los servicios no se llaman entre sí, sino que publican y consumen eventos. Hasta ahora esos eventos (pedido.creado, stock.reservado, pago.confirmado, pedido.confirmado, pedido.cancelado) han sido nombres en un diagrama. Esta lección los convierte en mensajes reales que viajan por RabbitMQ.
Veremos primero cuándo conviene la comunicación asíncrona y qué acoplamiento elimina; después el vocabulario imprescindible (mensaje, evento y comando, cola y pub-sub, broker, ack, reentrega, dead-letter queue, orden y garantías de entrega); a continuación RabbitMQ en detalle (exchanges, routing keys, colas, bindings) con la topología concreta de TechCorp; el código con amqplib para publicar el sobre de evento estándar y para consumir con prefetch, ack/nack y reenvío a la DLQ; el enganche con el outbox y la idempotencia de 02-05; y una comparación de RabbitMQ con Kafka y las colas gestionadas de la nube que justifica la elección de TechCorp. La integración de punta a punta de un servicio real (con su outbox, sus consumidores arrancando junto a Express) es de 04-04, y desplegar el broker es del módulo 5.
Contenido
- Síncrono frente a asíncrono: qué acoplamiento eliminamos
- Vocabulario de la mensajería
- Garantías de entrega: at-most-once, at-least-once, exactly-once
- RabbitMQ: exchanges, colas, routing keys y bindings
- La topología de TechCorp
- Publicar un evento con
amqplib - Consumir eventos:
prefetch,ack,nacky DLQ - Enganche con outbox e idempotencia
- RabbitMQ frente a Kafka y colas gestionadas
- Síncrono frente a asíncrono: qué acoplamiento eliminamos
Cuando Pedidos llama a GET /productos?ids= y espera la respuesta, Pedidos y Catálogo tienen que estar vivos en el mismo instante: es el acoplamiento temporal de 02-01. Con seis servicios encadenados en el flujo de pedido, ese acoplamiento sería fatal: la disponibilidad del conjunto sería el producto de las seis, y una caída de Notificaciones impediría crear pedidos. La mensajería asíncrona rompe esa cadena: Pedidos deja el evento en el broker y sigue; Inventario lo recogerá cuando pueda, aunque sea dentro de un minuto.
| Criterio | Síncrono (REST, gRPC) | Asíncrono (mensajes) |
|---|---|---|
| El llamante espera la respuesta | Sí | No; sigue su trabajo |
| Acoplamiento temporal | Sí: ambos deben estar disponibles | No: el broker guarda el mensaje |
| Quién conoce a quién | El llamante conoce al receptor (URL) | El productor no sabe quién consume |
| Fallo del receptor | Error inmediato al llamante | El mensaje espera; se procesa al volver |
| Picos de carga | El receptor los sufre en tiempo real | La cola los amortigua |
| Modelo mental | Pregunta-respuesta | Hecho ocurrido / orden dada |
| Depuración | Sencilla: una traza lineal | Más difícil: causa y efecto separados en el tiempo (06-02) |
| Cuándo la usa TechCorp | Cuando necesitamos la respuesta para continuar: consultar precio y nombre antes de guardar el pedido | Cuando el resultado no hace falta ahora: reservar, cobrar, notificar |
La regla práctica que aplicará TechCorp: por defecto asíncrono entre servicios; síncrono solo cuando la respuesta es imprescindible para responder al usuario. Es la misma conclusión de 02-02, ahora con la justificación técnica.
- Vocabulario de la mensajería
- Mensaje. Unidad de datos que viaja por el broker: unas cabeceras y un cuerpo (en TechCorp, el sobre JSON de 02-05).
- Evento frente a comando. Un evento describe algo que ya ha ocurrido (
pedido.creado); quien lo publica no espera nada de nadie, y puede haber cero o muchos consumidores. Un comando es una orden dirigida a un receptor concreto (reservar-stock): quien lo envía espera que alguien lo ejecute. La saga por coreografía de 02-05 usa solo eventos; una saga por orquestación usaría comandos. La distinción importa porque cambia quién decide: con eventos, decide el consumidor si le interesa; con comandos, decide el emisor a quién manda. - Productor (publisher) y consumidor (subscriber). Quien envía y quien recibe. Un servicio suele ser ambas cosas: Inventario consume
pedido.creadoy producestock.reservado. - Broker. El intermediario que recibe, almacena y entrega mensajes: RabbitMQ, Kafka, SQS. Es la "tubería tonta" de 02-01: enruta, no decide.
- Cola (queue). Un búfer FIFO del que uno o varios consumidores compiten por los mensajes: cada mensaje lo procesa una instancia. Ideal para repartir trabajo entre réplicas del mismo servicio.
- Tópico / pub-sub. Un mensaje se copia a todos los suscriptores interesados. Ideal para eventos:
pedido.confirmadointeresa a Notificaciones y a Inventario, y ambos deben recibirlo. En RabbitMQ, la combinación "un exchange que reparte a varias colas, cada cola con las réplicas de un servicio compitiendo" da los dos comportamientos a la vez. - Ack / nack. El consumidor confirma (acknowledge) que ha procesado el mensaje; hasta entonces el broker lo considera pendiente. Si el consumidor lo rechaza (negative ack) o muere sin confirmar, el broker lo reentrega (redelivery), a la misma o a otra instancia.
- Dead-letter queue (DLQ). Cola a la que va un mensaje que no se ha podido procesar (rechazado sin reencolar, caducado o expulsado por límite). Evita el "mensaje venenoso" que se reentrega infinitamente y permite inspeccionarlo y reprocesarlo a mano.
- Orden. RabbitMQ conserva el orden dentro de una cola para un solo consumidor; con varias réplicas consumiendo en paralelo, el orden global no está garantizado. Por eso el diseño de 02-05 hace que cada evento sea autocontenido y que la máquina de estados del pedido rechace transiciones fuera de orden.
- Garantías de entrega: at-most-once, at-least-once, exactly-once
| Garantía | Cómo se consigue | Qué puede pasar | Cuándo vale |
|---|---|---|---|
| At-most-once (como mucho una vez) | El broker entrega y olvida; el consumidor hace ack antes de procesar (o no hay ack) | Se pierden mensajes si el consumidor falla a mitad | Métricas, logs no críticos |
| At-least-once (al menos una vez) | Mensajes persistentes; el consumidor hace ack después de procesar; si falla, reentrega | Se duplican mensajes: si el consumidor procesa y muere antes del ack, lo recibe otra vez | Todo lo de negocio: pedidos, pagos |
| Exactly-once (exactamente una vez) | Requiere que broker y consumidor compartan una transacción, o que el consumidor deduplique | En la práctica no existe de extremo a extremo entre sistemas distintos (el broker no puede saber si tu UPDATE en PostgreSQL se confirmó) |
Solo dentro de un mismo sistema (Kafka Streams entre tópicos de Kafka) |
La conclusión realista, que ya adelantó 02-05: elegimos at-least-once + idempotencia. El broker garantiza que ningún evento se pierde (aunque alguno llegue dos veces), y el consumidor garantiza que procesar dos veces el mismo eventoId no tiene efecto (procesarUnaVez() con la tabla eventos_procesados). Es la única combinación que funciona con una base de datos por servicio.
- RabbitMQ: exchanges, colas, routing keys y bindings
RabbitMQ implementa el protocolo AMQP 0-9-1, cuyo modelo tiene cuatro piezas:
- Exchange. El productor nunca publica directamente en una cola: publica en un exchange con una routing key (una cadena como
pedido.creado). El exchange decide a qué colas copiar el mensaje según su tipo:direct: a las colas cuyo binding coincida exactamente con la routing key.fanout: a todas las colas enlazadas, ignorando la routing key.topic: a las colas cuyo patrón de binding encaje con la routing key, con comodines:*sustituye exactamente una palabra (separadas por puntos) y#cero o más.pedido.*encaja conpedido.creadoypedido.cancelado;#encaja con todo.headers: enruta por cabeceras; poco usado.
- Cola. Donde esperan los mensajes. Puede ser duradera (sobrevive a un reinicio del broker) y contener mensajes persistentes (escritos a disco). Para at-least-once hacen falta las dos cosas.
- Binding. La regla que une un exchange con una cola: "los mensajes con routing key que encaje con
pedido.creadovan a la colainventario.pedidos". - Canal (channel). Conexión lógica multiplexada dentro de una conexión TCP; cada hilo o consumidor usa el suyo.
Un mensaje llega a un exchange con su routing key, el exchange lo copia a cada cola cuyo binding encaja (así se consigue el pub-sub), y dentro de cada cola las réplicas del servicio consumidor compiten (así se consigue el reparto de carga). Un mismo mensaje puede acabar en tres colas y ser procesado una vez en cada una.
- La topología de TechCorp
Decisiones:
- Un único exchange de tipo
topic, llamadotechcorp.eventos, duradero. Todos los servicios publican en él. - La routing key es igual al tipo de evento:
pedido.creado,stock.reservado,pago.confirmado, etc. Así el nombre del evento y su enrutamiento son la misma cosa y no hay que mantener dos vocabularios. - Una cola por servicio consumidor, con nombre
<consumidor>.<tema>, con bindings a los eventos que le interesan. Las réplicas de un servicio comparten su cola (reparto de carga); servicios distintos tienen colas distintas (cada uno recibe su copia). - Cada cola tiene una DLQ asociada (
<cola>.dlq) a través de un exchangetechcorp.eventos.dlxde tipodirect.
| Cola | Servicio | Bindings (routing keys) | Qué hace con ellos |
|---|---|---|---|
inventario.pedidos |
Inventario | pedido.creado, pedido.confirmado, pedido.cancelado |
Reserva stock; consume la reserva; libera la reserva |
pagos.stock |
Pagos | stock.reservado, pedido.cancelado |
Cobra cuando hay reserva; reembolsa si el pedido se cancela tras cobrar |
notificaciones.pedidos |
Notificaciones | pedido.confirmado, pedido.cancelado |
Envía correo de confirmación o de cancelación |
pedidos.saga |
Pedidos | stock.reservado, stock.rechazado, pago.confirmado, pago.rechazado |
Avanza la máquina de estados del pedido y publica pedido.confirmado / pedido.cancelado |
pedidos.clientes |
Pedidos | cliente.actualizado |
Mantiene la réplica clientes_ref de 02-04 |
flowchart LR
subgraph Productores
P[servicio-pedidos]
I[servicio-inventario]
G[servicio-pagos]
C[servicio-clientes]
end
X{{"exchange techcorp.eventos (topic)"}}
P -- "pedido.creado / pedido.confirmado / pedido.cancelado" --> X
I -- "stock.reservado / stock.rechazado / stock.liberado" --> X
G -- "pago.confirmado / pago.rechazado / pago.reembolsado" --> X
C -- "cliente.actualizado" --> X
X -- "pedido.creado, pedido.confirmado, pedido.cancelado" --> Q1[(inventario.pedidos)]
X -- "stock.reservado, pedido.cancelado" --> Q2[(pagos.stock)]
X -- "pedido.confirmado, pedido.cancelado" --> Q3[(notificaciones.pedidos)]
X -- "stock.*, pago.*" --> Q4[(pedidos.saga)]
X -- "cliente.actualizado" --> Q5[(pedidos.clientes)]
Q1 --> CI[Inventario x N réplicas]
Q2 --> CG[Pagos x N]
Q3 --> CN[Notificaciones x N]
Q4 --> CP[Pedidos x N]
Q5 --> CP
Q1 -. "nack sin requeue" .-> DLX{{"techcorp.eventos.dlx"}}
Q2 -.-> DLX
Q3 -.-> DLX
Q4 -.-> DLX
DLX --> D1[(inventario.pedidos.dlq)]
DLX --> D2[(pagos.stock.dlq)]
DLX --> D3[(notificaciones.pedidos.dlq)]
DLX --> D4[(pedidos.saga.dlq)]
Fíjate en que Pedidos publica pedido.confirmado al recibir pago.confirmado, y ese evento lo consumen dos colas distintas (Notificaciones e Inventario): la copia la hace el exchange, no Pedidos. Pedidos no sabe ni le importa cuántos consumidores hay: eso es la "conformidad" de Notificaciones y la asimetría del mapa de contextos de 02-03 hechas topología.
Sobre pedidos.saga con el binding stock.* y pago.*: es cómodo, pero también recibiría stock.liberado y pago.reembolsado, que Pedidos no necesita. En la práctica se declaran los cuatro bindings explícitos de la tabla; el comodín aparece en el diagrama solo por brevedad.
- Publicar un evento con
amqplib
amqplibamqplib es la librería estándar de Node.js para AMQP 0-9-1. Primero, la declaración de la topología. Cada servicio declara al arrancar lo que usa (el exchange y sus propias colas); assert* es idempotente: si ya existe con los mismos parámetros, no hace nada.
// mensajeria/topologia.js (compartido vía @techcorp/comun-http o copiado en cada servicio)
const amqp = require('amqplib');
const EXCHANGE = 'techcorp.eventos';
const EXCHANGE_DLX = 'techcorp.eventos.dlx';
async function conectar(url = process.env.RABBITMQ_URL ?? 'amqp://localhost:5672') {
// 1. Una conexión TCP por proceso...
const conexion = await amqp.connect(url);
// 2. ...y un canal por uso (aquí uno para publicar; los consumidores abrirán el suyo)
const canal = await conexion.createChannel();
// 3. Declarar el exchange principal (topic, duradero) y el de dead-letter (direct, duradero)
await canal.assertExchange(EXCHANGE, 'topic', { durable: true });
await canal.assertExchange(EXCHANGE_DLX, 'direct', { durable: true });
return { conexion, canal };
}
// Declara una cola de consumidor con su DLQ y sus bindings
async function declararColaConsumidor(canal, nombreCola, routingKeys) {
// 4. La DLQ: cola duradera enlazada al exchange DLX con la routing key = nombre de la cola original
await canal.assertQueue(`${nombreCola}.dlq`, { durable: true });
await canal.bindQueue(`${nombreCola}.dlq`, EXCHANGE_DLX, nombreCola);
// 5. La cola principal: duradera y con dead-lettering configurado
await canal.assertQueue(nombreCola, {
durable: true,
arguments: {
'x-dead-letter-exchange': EXCHANGE_DLX, // a dónde van los mensajes rechazados
'x-dead-letter-routing-key': nombreCola // con qué routing key (→ su .dlq)
}
});
// 6. Un binding por cada tipo de evento que interesa a este consumidor
for (const rk of routingKeys) {
await canal.bindQueue(nombreCola, EXCHANGE, rk);
}
}
module.exports = { conectar, declararColaConsumidor, EXCHANGE };Ahora la publicación. La función recibe el sobre de evento estándar de 02-05 ya construido (eventoId, tipo, version, ocurridoEn, carga) y lo envía:
// mensajeria/publicador.js
const { randomUUID } = require('node:crypto');
const { EXCHANGE } = require('./topologia');
function construirSobre(tipo, carga, { version = 1 } = {}) {
return {
eventoId: `evt-${randomUUID()}`, // id único: lo usa procesarUnaVez() del consumidor
tipo, // 'pedido.creado'
version, // versión del esquema de la carga (03-06)
ocurridoEn: new Date().toISOString(), // cuándo pasó, en UTC
carga // el JSON de negocio
};
}
function publicarEvento(canal, sobre) {
const cuerpo = Buffer.from(JSON.stringify(sobre));
// publish devuelve false si el búfer interno está lleno (backpressure); lo tratamos en 06-04
return canal.publish(
EXCHANGE, // exchange destino
sobre.tipo, // routing key = tipo de evento
cuerpo,
{
persistent: true, // se escribe a disco: sobrevive a un reinicio del broker
contentType: 'application/json',
messageId: sobre.eventoId, // duplicamos el id en la cabecera AMQP para herramientas
type: sobre.tipo,
timestamp: Math.floor(Date.now() / 1000),
headers: { 'x-version': sobre.version }
}
);
}
module.exports = { construirSobre, publicarEvento };Y su uso desde el relay del outbox de Pedidos, con el evento pedido.creado del pedido ped-88213 (el JSON de carga es el que fijamos en 02-05):
const sobre = construirSobre('pedido.creado', {
pedidoId: 'ped-88213',
clienteId: 'c-1024',
cliente: { email: '[email protected]', nombre: 'Ana Ruiz' },
direccionEnvio: { calle: 'Gran Vía 12', codigoPostal: '28013', ciudad: 'Madrid', pais: 'ES' },
lineas: [
{ productoId: 'p-501', nombre: 'Auriculares BT X200', cantidad: 1, precioUnitario: 59.90 },
{ productoId: 'p-777', nombre: 'Cable USB-C 2 m', cantidad: 2, precioUnitario: 9.90 }
],
total: 79.70
});
publicarEvento(canal, sobre);Tres detalles que marcan la diferencia entre "funciona en mi máquina" y "no pierde pedidos":
persistent: truey colasdurable: truevan juntos. Un mensaje persistente en una cola no duradera se pierde igual al reiniciar; una cola duradera con mensajes no persistentes, también.- Confirmaciones del productor. Con
createChannel(),publishes "dispara y olvida": si RabbitMQ cae justo entonces, el mensaje se pierde sin error. ConcreateConfirmChannel()el broker confirma cada publicación y podemos esperar conawait canal.waitForConfirms()antes de marcar la fila del outbox como enviada. Es lo que usará el relay en 04-04. - El
eventoIdviaja en el cuerpo y enmessageId. En el cuerpo porque es parte del contrato del evento; en la cabecera porque la consola de RabbitMQ y las herramientas de DLQ lo muestran sin abrir el JSON.
- Consumir eventos:
prefetch, ack, nack y DLQ
prefetch, ack, nack y DLQEl consumidor de Inventario para la cola inventario.pedidos:
// mensajeria/consumidorInventario.js (servicio-inventario)
const { conectar, declararColaConsumidor } = require('./topologia');
async function iniciarConsumidor({ manejadores, procesarUnaVez }) {
const { conexion, canal } = await conectar();
const COLA = 'inventario.pedidos';
await declararColaConsumidor(canal, COLA, ['pedido.creado', 'pedido.confirmado', 'pedido.cancelado']);
// 1. prefetch: cuántos mensajes sin ack puede tener esta instancia a la vez.
// Sin esto, RabbitMQ volcaría toda la cola en la primera réplica que se conecte.
await canal.prefetch(10);
await canal.consume(COLA, async (msg) => {
if (msg === null) return; // el canal se ha cerrado
let sobre;
try {
sobre = JSON.parse(msg.content.toString());
} catch (err) {
// 2. Mensaje que ni siquiera es JSON: no tiene sentido reintentar → a la DLQ
// nack(msg, allUpTo=false, requeue=false) → RabbitMQ lo manda al DLX configurado
console.error('Mensaje ilegible, enviado a DLQ', { messageId: msg.properties.messageId });
return canal.nack(msg, false, false);
}
const manejador = manejadores[sobre.tipo];
if (!manejador) {
// 3. Evento con binding pero sin manejador (p. ej. una versión desplegada a medias): DLQ, no perderlo
console.warn('Sin manejador para el tipo', { tipo: sobre.tipo, eventoId: sobre.eventoId });
return canal.nack(msg, false, false);
}
try {
// 4. procesarUnaVez (02-05): si eventoId ya está en eventos_procesados para 'inventario', no hace nada
await procesarUnaVez(sobre.eventoId, 'inventario', () => manejador(sobre));
// 5. Todo bien: ack. Solo ahora RabbitMQ borra el mensaje de la cola
canal.ack(msg);
} catch (err) {
// 6. Error al procesar. ¿Es transitorio (BD caída) o permanente (datos imposibles)?
const esPrimeraVez = !msg.fields.redelivered;
if (err.transitorio && esPrimeraVez) {
// Reencolar UNA vez: RabbitMQ lo volverá a entregar (a esta u otra réplica)
console.warn('Error transitorio, reencolando', { eventoId: sobre.eventoId, error: err.message });
canal.nack(msg, false, true);
} else {
// Segundo fallo o error permanente: a la DLQ para inspección humana
console.error('Evento enviado a DLQ', { eventoId: sobre.eventoId, tipo: sobre.tipo, error: err.message });
canal.nack(msg, false, false);
}
}
}, { noAck: false }); // 7. noAck: false = modo at-least-once (ack manual). Es el valor por defecto, pero se explicita
// 8. Cierre ordenado: al parar el proceso, cerrar canal y conexión para que los mensajes sin ack se reentreguen ya
process.on('SIGTERM', async () => { await canal.close(); await conexion.close(); });
}
module.exports = { iniciarConsumidor };Y los manejadores que Inventario registra (solo la firma; la lógica de reserva es del módulo 4):
iniciarConsumidor({
procesarUnaVez,
manejadores: {
'pedido.creado': (sobre) => reservarStockParaPedido(sobre.carga), // → publica stock.reservado o stock.rechazado
'pedido.confirmado': (sobre) => consumirReserva(sobre.carga.pedidoId), // ACTIVA → CONSUMIDA
'pedido.cancelado': (sobre) => liberarReserva(sobre.carga.pedidoId) // ACTIVA → LIBERADA, publica stock.liberado
}
});Puntos que conviene entender bien:
prefetch(10)limita el trabajo en vuelo por réplica. Un valor bajo reparte mejor entre réplicas y evita que un proceso que muere arrastre cientos de mensajes a reentrega; un valor alto aumenta el rendimiento. Diez es un punto de partida razonable para manejadores que tocan la base de datos.ackdespués de procesar es lo que da at-least-once. Si el proceso muere entremanejador()yack, RabbitMQ reentrega yprocesarUnaVezevita el efecto doble. Nunca hagasackal principio "para que no se atasque": eso es at-most-once con nombre bonito.nack(msg, false, requeue): el segundo argumento (allUpTo) rechaza también todos los anteriores sin ack; casi siemprefalse. El tercero decide entre reencolar (true) y descartar/dead-letter (false).- Reencolar sin límite es un bucle infinito. Un mensaje que falla siempre volvería a la cabeza de la cola y bloquearía a los demás. Por eso se reencola como mucho una vez (
msg.fields.redelivereddice si ya lo fue) y después va a la DLQ. Los reintentos con espera creciente se ven en 06-03; RabbitMQ no los da de serie. - La DLQ no se consume automáticamente. Alguien (una alerta de 06-05 y una persona) mira
inventario.pedidos.dlq, entiende por qué falló y decide si reprocesar (mover el mensaje de vuelta) o descartar.
- Enganche con outbox e idempotencia
Con lo visto, la cadena completa de garantías para "un cliente hace un pedido" queda así, y merece la pena verla junta aunque cada pieza se diseñó en 02-05:
| Riesgo | Pieza que lo cubre | Dónde vive |
|---|---|---|
Se guarda el pedido pero no se publica pedido.creado (o al revés) |
Outbox transaccional: pedido y evento en la misma transacción; un relay lee outbox y llama a publicarEvento |
Pedidos (guardarConEventos()) |
| El relay publica y RabbitMQ cae antes de guardarlo | persistent: true + cola duradera + canal de confirmaciones |
publicador.js |
Inventario procesa y muere antes del ack → reentrega |
Idempotencia del consumidor: procesarUnaVez(eventoId, 'inventario', fn) con eventos_procesados |
Cada consumidor |
| El relay reenvía la misma fila de outbox dos veces | Mismo eventoId en ambas copias → el consumidor la deduplica |
Contrato del sobre |
| Un mensaje imposible de procesar bloquea la cola | DLQ | Topología |
El usuario reenvía POST /pedidos |
Idempotency-Key (03-01) |
Pedidos |
Nada de esto es exótico: es at-least-once en cada salto más una deduplicación por eventoId en el destino. Es el precio de no tener transacciones distribuidas, y es barato.
- RabbitMQ frente a Kafka y colas gestionadas
| Criterio | RabbitMQ | Apache Kafka | Colas gestionadas (AWS SQS/SNS, Google Pub/Sub) |
|---|---|---|---|
| Modelo | Broker AMQP: exchanges, colas, enrutamiento flexible; el mensaje se borra al hacer ack | Log distribuido: tópicos particionados, los mensajes se retienen (días o para siempre); los consumidores guardan su offset | Cola (SQS) + pub-sub (SNS) o tópico con suscripciones (Pub/Sub); sin operar nada |
| Enrutamiento | Muy rico (topic con comodines, headers) | Por tópico y partición; el filtrado lo hace el consumidor | Filtros de suscripción sencillos |
| Orden | Por cola con un consumidor | Por partición, garantizado; muy fuerte | Solo con colas FIFO / ordering keys |
| Reproducir mensajes antiguos | No (una vez consumido, se fue) | Sí: releer desde un offset; base del event sourcing | No (o retención corta) |
| Rendimiento | Decenas de miles de msg/s por nodo | Millones de msg/s; diseñado para streaming | Elástico; pagas por uso |
| Operación | Un clúster sencillo; consola web excelente | Más complejo (particiones, ZooKeeper/KRaft, rebalanceos) | Nula, pero lock-in con el proveedor |
| Latencia | Muy baja (ms) | Baja, pero orientada a lotes | Variable (decenas de ms) |
| Curva de aprendizaje | Suave; conceptos intuitivos | Empinada; nuevo modelo mental | Suave |
| Encaja cuando | Eventos de negocio y comandos entre servicios, enrutamiento variado, volumen medio | Streaming, analítica, reprocesado histórico, volúmenes enormes | Ya estás en esa nube y no quieres operar un broker |
Por qué TechCorp elige RabbitMQ:
- El volumen (~3.000 pedidos/día, unos pocos eventos por pedido) está a órdenes de magnitud del umbral en que Kafka compensa su complejidad operativa. Con ~25 técnicos, el equipo de Plataforma no puede dedicar a nadie a cuidar particiones.
- El enrutamiento por topic con una cola por consumidor modela de forma directa el mapa de contextos de 02-03: cada servicio se suscribe a lo que le interesa y nadie más se entera.
- La DLQ, el prefetch y las confirmaciones cubren las garantías que la saga necesita sin código extra.
- En 02-05 descartamos event sourcing "por ahora"; si algún día se adopta, la retención y el reprocesado de Kafka serían el motivo para migrar, y el sobre de evento estándar hará esa migración menos dolorosa.
- Marta descartó las colas gestionadas para no atar el sistema a una nube en un momento en que aún se decide dónde correrá Kubernetes.
Queda un tema que esta lección ha rozado en cada fragmento: la forma de la carga de cada evento (qué campos lleva pedido.creado, qué pasa cuando hay que añadir uno o cambiar otro) y ese campo version del sobre. Es el contrato de los eventos, y se trata junto con el de las APIs en 03-06.
Errores Comunes y Consejos
- Publicar directamente en una cola (
sendToQueue) en lugar de en el exchange. Funciona hasta que un segundo servicio necesita el mismo evento; entonces hay que tocar al productor. Con el exchange solo se añade un binding. ackantes de procesar "para que vaya rápido". Convierte at-least-once en at-most-once: un reinicio a mitad y el pedido se queda sin reserva para siempre.- Reencolar sin límite (
nack(msg, false, true)en todocatch). Un mensaje venenoso monopoliza la cola. Una reentrega y a la DLQ. - Olvidar
prefetch. La primera réplica que arranca se lleva todos los mensajes pendientes; las demás se quedan mirando. - Cola no duradera o mensaje no persistente. Todo se ve bien hasta el primer reinicio del broker. Ambos flags, siempre, para eventos de negocio.
- Cargas de evento que apuntan a datos (
{"pedidoId": "ped-88213"}y que el consumidor llame aGET /pedidos/ped-88213). Reintroduce el acoplamiento temporal que queríamos quitar. El evento lleva lo que el consumidor necesita (por esopedido.creadoincluye líneas, precios y correo). - Confiar en el orden entre colas o entre réplicas. Diseña los consumidores para tolerar
pago.confirmadoantes de que su propia BD reflejestock.reservado(la máquina de estados de 02-05 yprocesarUnaVezayudan). - Ignorar la DLQ. Sin alerta sobre
*.dlq, los mensajes se acumulan meses. En 06-05 se define la alerta; desde hoy, mírala en la consola de RabbitMQ. - Una conexión por petición. Las conexiones AMQP son caras. Una conexión por proceso, un canal por consumidor o publicador, reutilizados.
Ejercicios
Ejercicio 1. El equipo de Pagos y comunicaciones quiere que Notificaciones envíe también un correo "hemos recibido tu pedido" en cuanto se crea (además del de confirmación). Indica qué binding hay que añadir, en qué cola, y qué no hay que tocar. Después, razona: ¿tendría sentido que Notificaciones usara la misma cola notificaciones.pedidos o una nueva notificaciones.pedidos-recibidos? Da un argumento a favor de cada opción.
Ejercicio 2. Un consumidor de Pagos procesa stock.reservado, llama a la pasarela externa, que cobra 79,70 € a Ana, y justo antes del ack la instancia muere. RabbitMQ reentrega el mensaje a otra réplica. Explica paso a paso qué ocurre con y sin procesarUnaVez, y qué garantía adicional necesita el consumidor de Pagos respecto a la pasarela (pista: 03-01 habló de ello con otro nombre).
Ejercicio 3. Escribe una función moverDeDlqACola(canal, nombreCola, maximo) con amqplib que lea hasta maximo mensajes de ${nombreCola}.dlq con canal.get() (obtención síncrona, sin suscripción) y los vuelva a publicar en el exchange techcorp.eventos con la routing key original (disponible en msg.fields.routingKey o, tras el dead-lettering, en la cabecera x-death), conservando persistent: true, y haga ack de cada uno en la DLQ solo tras republicarlo. Comenta cada línea.
Soluciones
Solución 1.
Basta con añadir el binding pedido.creado a la cola de Notificaciones (bindQueue('notificaciones.pedidos', 'techcorp.eventos', 'pedido.creado')) y registrar un manejador 'pedido.creado' en su consumidor. No hay que tocar Pedidos (sigue publicando exactamente igual), ni Inventario, ni el exchange: eso es lo que compra el pub-sub por topic. La carga de pedido.creado ya lleva cliente.email y cliente.nombre precisamente para que Notificaciones no tenga que llamar a nadie.
Misma cola: más simple, un solo consumidor, un solo prefetch, una sola DLQ que vigilar; el orden relativo "recibido → confirmado" para un mismo pedido se conserva mejor. Cola nueva: aísla el fallo (si el correo de "recibido" se rompe y llena la DLQ, los de confirmación siguen saliendo) y permite escalar y priorizar por separado (los de confirmación son más importantes). Para TechCorp, con el volumen actual, la misma cola; separarlas cuando haya un motivo operativo.
Solución 2.
Sin procesarUnaVez: la segunda réplica recibe el mismo stock.reservado (mismo eventoId), no sabe que ya se procesó, vuelve a llamar a la pasarela y cobra 79,70 € dos veces; además publica dos pago.confirmado. Con procesarUnaVez: la segunda réplica consulta eventos_procesados para (evt-..., 'pagos')... y aquí está la trampa: si la primera réplica murió antes de confirmar la transacción que inserta en eventos_procesados, el evento no consta como procesado, y la segunda réplica también cobrará. procesarUnaVez protege contra "procesé y morí antes del ack" solo si el efecto y el registro en eventos_procesados están en la misma transacción local, y la llamada a la pasarela externa no puede estar dentro de esa transacción.
La garantía adicional: la pasarela debe aceptar una clave de idempotencia (la mayoría de pasarelas reales lo hacen), y Pagos debe usar como clave algo derivado del pedido (ped-88213) o del eventoId, exactamente el mismo concepto que la cabecera Idempotency-Key de 03-01, pero de Pagos hacia fuera. Así, la segunda llamada a la pasarela devuelve el mismo cargo en lugar de crear otro. Regla general: la idempotencia se necesita en cada frontera donde hay un efecto no reversible.
Solución 3.
async function moverDeDlqACola(canal, nombreCola, maximo = 100) {
const dlq = `${nombreCola}.dlq`;
let movidos = 0;
for (let i = 0; i < maximo; i++) {
// get() obtiene UN mensaje sin suscribirse; devuelve false si la DLQ está vacía
const msg = await canal.get(dlq, { noAck: false });
if (!msg) break;
// Tras el dead-lettering, msg.fields.routingKey es la del DLX (= nombreCola).
// La routing key original (p. ej. 'pedido.creado') queda registrada en la cabecera x-death.
const xDeath = msg.properties.headers?.['x-death']?.[0];
const routingKeyOriginal = xDeath?.['routing-keys']?.[0] ?? msg.properties.type;
if (!routingKeyOriginal) {
// Sin forma de saber a dónde iba: lo dejamos en la DLQ (nack con requeue) y seguimos
canal.nack(msg, false, true);
continue;
}
// Republicamos en el exchange principal con las mismas propiedades (persistent, messageId, type...)
canal.publish('techcorp.eventos', routingKeyOriginal, msg.content, {
...msg.properties,
persistent: true,
headers: { ...msg.properties.headers, 'x-reprocesado-desde-dlq': dlq }
});
// Solo tras republicar hacemos ack en la DLQ: si el proceso muere entre medias,
// el mensaje sigue en la DLQ (podría duplicarse, pero procesarUnaVez lo absorbe)
canal.ack(msg);
movidos++;
}
return movidos;
}Notas: msg.properties.type lo rellenamos en publicarEvento con sobre.tipo, así que sirve de respaldo si x-death no estuviera. Con un canal de confirmaciones sería aún más seguro esperar waitForConfirms() antes del ack. Y sí, este script puede duplicar un evento en el peor caso; como todo en esta lección, la idempotencia del consumidor es lo que lo hace inofensivo.
Conclusión
La mensajería asíncrona es lo que hace posible la saga de 02-05: elimina el acoplamiento temporal entre servicios, amortigua picos y permite que un mismo hecho (pedido.confirmado) llegue a varios interesados sin que el productor sepa de ellos. Hemos fijado el vocabulario (evento frente a comando, cola frente a pub-sub, ack/nack, reentrega, DLQ), aceptado at-least-once + idempotencia como la garantía realista, y construido la topología de TechCorp en RabbitMQ: un exchange topic techcorp.eventos, routing keys iguales al tipo de evento, una cola duradera por consumidor (inventario.pedidos, pagos.stock, notificaciones.pedidos, pedidos.saga, pedidos.clientes) con su DLQ, y el código amqplib para publicar el sobre estándar con persistent: true y para consumir con prefetch, ack tras procesar y nack a la DLQ. Con REST para lo síncrono y RabbitMQ para lo asíncrono, TechCorp ya tiene sus dos canales principales.
Pero REST/JSON no es la única forma de llamada síncrona ni siempre la mejor: cuando dos servicios internos se hablan miles de veces por minuto, la verbosidad de JSON y la falta de contrato tipado pesan, y cuando un front-end necesita componer datos de varios servicios, REST obliga a muchas peticiones o a respuestas enormes. Para el primer caso existe gRPC; para el segundo, GraphQL. En la siguiente lección veremos ambos, con su código en Node.js, y decidiremos dónde encajan (y dónde no) en TechCorp.
Curso de Microservicios
Módulo 1: Introducción a los Microservicios
- Conceptos Básicos de Microservicios
- Ventajas y Desventajas de los Microservicios
- Comparación con la Arquitectura Monolítica
- Cuándo Adoptar Microservicios: Criterios de Decisión
- El Caso Práctico del Curso: la Tienda Online de TechCorp
Módulo 2: Diseño de Microservicios
- Principios de Diseño de Microservicios
- Descomposición de Aplicaciones Monolíticas
- Definición de Bounded Contexts
- Gestión de Datos: una Base de Datos por Servicio
- Consistencia Distribuida: Sagas, CQRS y Event Sourcing
Módulo 3: Comunicación entre Microservicios
- APIs RESTful
- Mensajería Asíncrona
- Protocolos de Comunicación: gRPC, GraphQL
- API Gateway y Backend for Frontend
- Descubrimiento de Servicios y Balanceo de Carga
- Contratos y Versionado de APIs
Módulo 4: Implementación de Microservicios
- Elección de Tecnologías y Herramientas
- Desarrollo de un Microservicio Simple
- Gestión de Configuración
- Integración Práctica: Consumir APIs y Publicar Eventos
- Pruebas en Microservicios: Unitarias, de Integración y de Contrato
Módulo 5: Despliegue y Orquestación
- Contenedores y Docker
- Orquestación con Kubernetes
- CI/CD para Microservicios
- Estrategias de Despliegue: Rolling, Blue-Green y Canary
- Service Mesh: Istio y Linkerd
Módulo 6: Monitoreo y Mantenimiento
- Monitoreo y Logging
- Trazabilidad Distribuida con OpenTelemetry
- Gestión de Errores y Recuperación
- Escalabilidad y Rendimiento
- SLOs, Alertas y Gestión de Incidentes
Módulo 7: Seguridad en Microservicios
- Autenticación y Autorización
- Seguridad en la Comunicación
- Prácticas de Seguridad
- Seguridad en Contenedores y Kubernetes
