¿Qué es RabbitMQ?
RabbitMQ es un message broker de código abierto que utiliza AMQP (Advanced Message Queuing Protocol). Los productores publican mensajes en él, los consumidores los reciben y el broker se encarga de retener cada mensaje hasta que alguien lo haya confirmado.
La frase que mejor describe su diseño es un broker inteligente con consumidores simples. RabbitMQ gestiona exchanges, reglas de enrutamiento, acknowledgements, reintentos y dead letters. Un consumidor es un programa pequeño que lee un mensaje, realiza una tarea y notifica que ha terminado. Esta división del trabajo es la razón por la cual RabbitMQ encaja en sistemas donde el enrutamiento y las garantías de entrega son más importantes que el rendimiento bruto (throughput).
Fue escrito en Erlang en 2007, razón por la cual es excepcionalmente bueno manejando múltiples conexiones concurrentes y tiene una reputación de estabilidad. Ejecuta un protocolo completo en lugar de una API de colas mínima, y ese protocolo es la fuente tanto de su potencia como de su curva de aprendizaje.
El modelo AMQP
La idea más importante de RabbitMQ es que un producer nunca publica directamente en una queue. Publica en un exchange, y es el exchange quien decide qué queues reciben una copia.
producer -> exchange --binding--> queue -> consumer
\--binding--> queue -> consumer
El modelo se compone de cinco elementos:
- Un producer abre un channel y llama a
basic.publishcon el nombre de un exchange, una routing key y un cuerpo. - Un exchange recibe cada mensaje publicado y lo enruta. No almacena nada a menos que esté vinculado a una queue.
- Un binding es una regla que conecta un exchange con una queue. Puede contener una routing key o un patrón, y un mismo exchange puede vincularse a varias queues.
- Una queue almacena los mensajes en orden hasta que un consumer los procesa.
- Un consumer se suscribe a una queue, recibe las entregas y las confirma (acknowledges).
Esta indirección es el objetivo principal. Un producer que publica order.created.eu no sabe si hay un servicio, cinco servicios o nadie escuchando. Se pueden añadir nuevos consumers vinculando una nueva queue, sin necesidad de modificar el producer.
Todo el trabajo ocurre en un channel, que es una conexión virtual ligera multiplexada sobre una única conexión TCP. Los channels no son thread-safe, por lo que el patrón común es utilizar un channel por tarea o por consumer.
import amqp from "amqplib";
const conn = await amqp.connect(process.env.AMQP_URL!);
const ch = await conn.createChannel();
Tipos de exchange y enrutamiento
El tipo de un exchange determina cómo empareja una routing key con sus bindings. Existen cuatro tipos, y cada uno tiene un uso claro.
Un direct exchange enruta a las colas cuya binding key sea exactamente igual a la routing key. Úsalo para enviar un mensaje a una cola específica, o a un grupo pequeño que comparta la misma key.
await ch.assertExchange("logs", "direct", { durable: true });
await ch.bindQueue("logs.errors", "logs", "error");
await ch.publish("logs", "error", body);
Un fanout exchange ignora completamente la routing key y copia el mensaje a todas las colas vinculadas. Es pub/sub en su forma más pura: una publicación, muchos consumidores independientes.
await ch.assertExchange("events", "fanout", { durable: true });
await ch.bindQueue("search-indexer", "events", "");
await ch.bindQueue("email-notifier", "events", "");
Un topic exchange empareja la routing key basándose en un patrón. Las keys son palabras separadas por puntos. * coincide exactamente con una palabra y # coincide con cero o más, lo que convierte a los topic exchanges en los más flexibles de los cuatro.
await ch.assertExchange("orders", "topic", { durable: true });
await ch.bindQueue("eu-orders", "orders", "order.*.eu");
await ch.bindQueue("all-orders", "orders", "order.#");
await ch.publish("orders", "order.created.eu", body);
Un headers exchange ignora la routing key y, en su lugar, realiza el emparejamiento basándose en los atributos del encabezado del mensaje. Rara vez es la opción correcta, ya que los topic exchanges son más fáciles de leer y analizar, pero es útil cuando el enrutamiento depende de varios atributos independientes.
Si recuerdas una sola regla, que sea esta: elige el tipo de exchange según la pregunta que te estés haciendo. “¿Qué cola exacta?” es direct. “¿Todos?” es fanout. “¿Qué familia de eventos?” es topic.
Acknowledgements, nack y prefetch
La garantía de entrega de RabbitMQ se basa en los acknowledgements. Cuando un consumidor recibe un mensaje, el broker lo marca como unacked pero lo conserva. El mensaje solo se elimina cuando el consumidor llama a ack. Si la conexión se interrumpe antes de eso, el broker vuelve a entregar el mensaje.
await ch.consume("orders.created", async (msg) => {
if (!msg) return;
try {
await handleOrder(JSON.parse(msg.content.toString()));
ch.ack(msg);
} catch (err) {
ch.nack(msg, false, false); // requeue: false, dead-letter instead
}
});
nack (o el antiguo reject) acepta un flag de requeue. El requeueing coloca el mensaje nuevamente al principio de la cola, lo cual es correcto para un fallo transitorio, pero incorrecto para un “poison message” que fallará permanentemente. La solución habitual es enviarlo a un dead-letter exchange.
El prefetch, que se configura con basic.qos, limita cuántos mensajes sin confirmar puede retener un consumidor a la vez. Sin esto, el broker envía los mensajes tan rápido como puede y un consumidor lento podría acumular miles de ellos en memoria.
await ch.prefetch(20);
Configura el prefetch como un múltiplo pequeño de tu concurrencia real. Si es demasiado alto, un solo consumidor acaparará la cola; si es demasiado bajo, el consumidor se quedará inactivo esperando el siguiente viaje de ida y vuelta.
Durabilidad y garantías de entrega
Un mensaje sobrevive al reinicio de un broker solo si se cumplen tres condiciones, y es fácil olvidar alguna.
- La queue debe estar declarada como
durable: true. - El mensaje debe publicarse con
persistent: true. - El exchange también debería ser
durable: true.
await ch.assertExchange("orders", "topic", { durable: true });
await ch.assertQueue("orders.created", { durable: true });
ch.publish("orders", "order.created.eu", body, { persistent: true });
La durabilidad se refiere a sobrevivir a un reinicio, no a garantizar la entrega. Para eso, activa los publisher confirms. Sin confirms, publish es un proceso de “disparar y olvidar” (fire-and-forget): si el broker muere antes de escribir el mensaje, el productor nunca lo sabrá.
const ch = await conn.createConfirmChannel();
ch.publish("orders", "order.created.eu", body, { persistent: true }, (err) => {
if (err) console.error("broker did not confirm", err);
else console.log("message is safely queued");
});
Incluso con confirms y queues durables, la entrega es at-least-once (al menos una vez). Un consumidor puede fallar después de realizar el trabajo pero antes de enviar el ack, y el broker volverá a entregar el mensaje. La entrega exactamente una vez (exactly-once) a través de una red es efectivamente imposible, por lo que RabbitMQ hace que este compromiso sea explícito y te pide que diseñes tus handlers para que sean idempotentes.
Dead-letter exchanges y colas de reintento
Un dead-letter exchange (DLX) es un exchange ordinario que recibe mensajes que una cola rechaza, que expiran o que se descartan. Es el mecanismo que sustenta las colas de reintento y el manejo de poison-messages.
Un mensaje se envía al dead-letter exchange cuando ocurre uno de estos eventos:
- El consumidor hace un nack o lo rechaza con
requeue: false. - Expira su TTL por mensaje o por cola.
- La cola excede su límite de longitud y descarta el mensaje más antiguo.
El DLX se configura en la cola que contiene el trabajo, no en el consumidor.
await ch.assertQueue("orders.created", {
durable: true,
deadLetterExchange: "orders.dlx",
deadLetterRoutingKey: "failed",
});
El patrón clásico de reintento retrasado (delayed retry) utiliza una segunda cola con un TTL y su propio DLX que apunta de vuelta al exchange de trabajo. Un mensaje fallido se envía al dead-letter exchange hacia la cola de reintento, permanece allí durante el TTL, expira y se envía nuevamente mediante dead-letter para ser procesado otra vez. Esto produce un reintento con retraso sin necesidad de implementar temporizadores en tu código.
await ch.assertExchange("orders.dlx", "direct", { durable: true });
await ch.assertQueue("orders.retry", {
durable: true,
messageTtl: 30_000,
deadLetterExchange: "orders",
deadLetterRoutingKey: "order.created.retry",
});
await ch.bindQueue("orders.retry", "orders.dlx", "retry");
Un mensaje que falla repetidamente entrará en un bucle. Rastrea un contador de reintentos en los headers del mensaje y, tras alcanzar un límite, redirígelo a una cola de failed permanente que nadie consuma automáticamente. Trata esa cola como una superficie operativa: genera alertas cuando crezca y construye una ruta de replay.
TTL de mensajes y límites de cola
El Time-to-live controla cuánto tiempo puede esperar un mensaje. Puedes configurarlo por cola, por mensaje, o ambos.
await ch.assertQueue("verification", {
durable: true,
messageTtl: 600_000, // 10 minutes for every message
maxLength: 10_000, // keep at most 10k messages
overflow: "reject-publish", // backpressure instead of dropping
});
El TTL por mensaje se establece al publicar y se utiliza a menudo para valores que varían según el mensaje.
ch.publish("orders", key, body, {
expiration: "30000", // milliseconds, as a string
});
Los límites de longitud de la cola sirven como backpressure. maxLength limita la memoria, y overflow decide qué sucede cuando se alcanza el límite: drop-head descarta silenciosamente el mensaje más antiguo, mientras que reject-publish rechaza nuevas publicaciones para que el productor sienta la presión. Para una cola que no debe perder trabajo, reject-publish es la opción predeterminada más segura y funciona como una señal de alerta.
Colas de trabajo vs pub/sub
El mismo broker cubre dos formas de comunicación muy diferentes, y confundirlas provoca errores.
Una cola de trabajo (work queue) distribuye cada mensaje a exactamente un consumidor. Muchos workers compiten por la misma cola y el broker distribuye las entregas mediante round-robin. Si añades un prefetch y acknowledgements, tienes una cola de tareas escalable y tolerante a fallos.
await ch.assertQueue("jobs.thumbnails", { durable: true });
await ch.prefetch(5);
Pub/sub entrega cada mensaje a todos los consumidores interesados. Cada suscriptor tiene su propia cola vinculada al mismo exchange, por lo que un suscriptor lento o desconectado nunca le “roba” un mensaje a los demás.
await ch.assertExchange("events", "fanout", { durable: true });
await ch.assertQueue("events.billing", { durable: true });
await ch.bindQueue("events.billing", "events", "");
await ch.assertQueue("events.analytics", { durable: true });
await ch.bindQueue("events.analytics", "events", "");
La regla de oro: una cola compartida por muchos consumidores es una cola de trabajo; una cola por consumidor es pub/sub. Un topic exchange te permite tener ambos a la vez, con algunos consumidores compartiendo una cola y otros teniendo la suya propia.
El patrón RPC
RabbitMQ también puede manejar solicitudes y respuestas (request/response). El cliente publica una solicitud con una cola replyTo y un correlationId, y el servidor publica la respuesta en esa cola con el mismo id.
import { randomUUID } from "node:crypto";
const { queue } = await ch.assertQueue("", { exclusive: true });
const correlationId = randomUUID();
ch.consume(queue, (msg) => {
if (msg?.properties.correlationId === correlationId) {
console.log("reply", JSON.parse(msg.content.toString()));
}
}, { noAck: true });
ch.publish("rpc.inventory", "check", Buffer.from("{}"), {
replyTo: queue,
correlationId,
});
El assertQueue("") crea una cola temporal, exclusiva y de auto-eliminación solo para este cliente. El RPC sobre una cola es útil cuando el llamador realmente necesita una respuesta, pero reintroduce el acoplamiento síncrono y el problema de los timeouts. Para la mayoría de los sistemas, un evento más un evento de callback es más sencillo de operar que el RPC, y una llamada HTTP simple es aún más fácil cuando la dependencia está disponible.
Clustering y quorum queues
Un único nodo representa un punto único de fallo, por lo que RabbitMQ en producción se ejecuta como un cluster. Las colas pueden replicarse entre nodos, y los clientes se reconectan a otro nodo cuando uno falla.
La estrategia de replicación moderna es la quorum queue, basada en el algoritmo de consenso Raft. Una quorum queue tiene un líder y seguidores, y una escritura se confirma solo una vez que la mayoría la ha recibido. Esto la hace segura ante particiones de red que causaban pérdida de datos en las classic mirrored queues, y es la opción predeterminada para colas duraderas hoy en día.
await ch.assertQueue("orders.created", {
durable: true,
arguments: { "x-queue-type": "quorum" },
});
Las quorum queues prefieren un número pequeño e impar de réplicas, típicamente tres o cinco. Son más pesadas que las classic queues, por lo que debes usarlas para datos que no puedes perder y mantener las colas transitorias como classic. Para topologías multirregión, la federation y el plugin shovel mueven mensajes entre brokers independientes en lugar de extender un solo cluster a través de un enlace lento.
La UI de gestión y el monitoreo
Cada nodo de RabbitMQ incluye un plugin de gestión que ofrece una UI web y una API HTTP. Esta muestra los exchanges, queues, bindings, conexiones y canales, y permite publicar un mensaje de prueba o reproducir uno desde una queue.
rabbitmq-plugins enable rabbitmq_management
curl -u guest:guest http://localhost:15672/api/queues/%2F/orders.created
Hay cuatro números que son los más importantes en un dashboard.
- Queue depth — mensajes listos. Una profundidad creciente significa que los consumidores no pueden seguir el ritmo.
- Unacked count — mensajes entregados pero no confirmados. Un número que solo crece indica que los consumidores están bloqueados.
- Tasas de publicación y entrega — la forma del tráfico y si los consumidores mantienen el ritmo.
- Tasa de redelivery — un recuento de redelivery creciente apunta a crashes o fallos repetidos.
RabbitMQ también emite métricas de Prometheus, por lo que la queue depth y la utilización de los consumidores deben estar en el mismo dashboard que el resto de tus servicios. Configura alertas basadas en la profundidad y en la antigüedad del mensaje más viejo, no en fallos individuales, ya que estos son esperados.
Conexiones, canales y recuperación
Una conexión es costosa; un canal es económico. Abre una sola conexión por proceso y luego crea un canal por productor, por consumidor o por unidad de trabajo. Un canal no es thread-safe, por lo que compartir uno entre manejadores concurrentes provoca frames entrelazados y errores confusos.
import amqp from "amqplib";
let conn: amqp.Connection;
let ch: amqp.Channel;
async function connect() {
conn = await amqp.connect(process.env.AMQP_URL!);
conn.on("error", (err) => console.error("connection error", err));
conn.on("close", () => setTimeout(connect, 5_000));
ch = await conn.createChannel();
await ch.assertExchange("orders", "topic", { durable: true });
}
Gestiona la reconexión de forma deliberada. amqplib no reconecta automáticamente, así que escucha el evento close y reconstruye la conexión, el canal y cada consumidor. Debido a que una conexión caída deja mensajes en tránsito sin confirmar (unacked), el broker los vuelve a entregar, que es el comportamiento deseado. Este es otro escenario donde aparece la entrega at-least-once: una reconexión puede repetir el trabajo, por lo que los manejadores deben ser idempotentes.
Propiedades y prioridades de los mensajes
Cada mensaje publicado puede llevar propiedades junto con su cuerpo. Estas viajan con el mensaje y son visibles para los consumidores, lo que las convierte en un lugar ligero para almacenar metadatos.
ch.publish("orders", "order.created.eu", body, {
persistent: true,
contentType: "application/json",
messageId: order.id,
correlationId: traceId,
timestamp: Math.floor(Date.now() / 1000),
type: "order.created",
headers: { "x-retry-count": 0, "x-source": "checkout" },
});
contentType y messageId ayudan a los consumidores y a las herramientas; correlationId vincula un mensaje a un trace; headers es donde guardas los metadatos de la aplicación, como el contador de reintentos. No coloques datos voluminosos en los headers, ya que el broker debe indexarlos y mostrarlos.
RabbitMQ también soporta priority queues, donde los mensajes con mayor prioridad se entregan primero. Declara la cola con una prioridad máxima y publica con un valor de priority.
await ch.assertQueue("jobs", {
durable: true,
maxPriority: 10,
});
ch.publish("", "jobs", body, { priority: 9 });
Las prioridades solo reordenan los mensajes que ya están esperando. Si los consumidores mantienen la cola vacía, la prioridad no surte efecto, y un mensaje de alta prioridad que llegue después de uno de baja prioridad seguirá esperando detrás de este. Utiliza la prioridad para una distinción de negocio real, no como un mecanismo de programación general.
Eligiendo entre RabbitMQ y Kafka
A menudo se comparan ambos, pero resuelven problemas diferentes.
RabbitMQ es un smart broker para el enrutamiento y colas de trabajo. Los mensajes son tareas que se eliminan una vez confirmadas. Los exchanges enrutan por patrones, los acknowledgements son por mensaje y el dead-lettering viene integrado. Destaca cuando un mensaje debe llegar a un conjunto específico de colas, o cuando una unidad de trabajo debe reintentarse y quedar en espera en caso de fallo.
Kafka es un durable log. Los mensajes se añaden a un log ordenado y particionado, y se conservan durante un tiempo configurado independientemente de quién los lea. Muchos consumer groups pueden leer el mismo topic de forma independiente, y cualquier consumidor puede retroceder a un offset anterior. Destaca en escenarios de throughput muy alto y en la reproducción de eventos (event replay).
Una heurística útil: si te molestaría que un mensaje fuera consumido y eliminado, necesitas un log. Si te importa que una tarea haya sido enrutada correctamente y completada, necesitas un broker. Muchos sistemas utilizan ambos: Kafka para el flujo de eventos y RabbitMQ para el trabajo.
Mejores prácticas
- Declara los exchanges, queues y bindings de forma idempotente al iniciar, y utiliza
durable: truepara cualquier dato que no puedas perder. - Publica mensajes persistentes en queues duraderas y utiliza publisher confirms para el trabajo que no debe desaparecer.
- Realiza siempre el acknowledge manualmente después de que el efecto secundario haya tenido éxito, nunca antes.
- Establece un límite de prefetch para que un solo consumidor no pueda almacenar en búfer toda la queue.
- Asigna a cada queue duradera un dead-letter exchange y construye una ruta de reintento (replay path) para ella.
- Implementa reintentos retardados con una TTL retry queue en lugar de usar sleep en el handler.
- Haz que los consumidores sean idempotentes, ya que la entrega es at-least-once.
- Limita la longitud de la queue y elige
reject-publishcuando no sea aceptable descartar mensajes. - Utiliza quorum queues para datos duraderos y classic queues para datos transitorios.
- Cierra los channels y connections en
SIGTERM, y detén el consumo antes del drenado. - Monitorea los mensajes ready, unacked, la tasa de redelivery y la antigüedad del mensaje más antiguo (oldest-message age).
Errores comunes
- Publicar en una cola en lugar de en un exchange, y luego preguntarse por qué los bindings no funcionan.
- Usar auto-ack (
noAck: true) y perder mensajes cada vez que un handler falla. - Olvidar
persistent: true, perdiendo así los mensajes al reiniciar el broker. - Declarar una cola como durable pero publicar mensajes que no son persistentes.
- Configurar el prefetch demasiado alto, permitiendo que un consumidor deje sin recursos al resto.
- Reencolar un “poison message” infinitamente con
nack(msg, false, true). - Crear un dead-letter exchange sin ningún consumidor y nunca revisarlo.
- Usar un fanout exchange cuando se necesitaba un topic exchange, o viceversa.
- Ejecutar un único broker en producción y decir que es altamente disponible.
- Mezclar semánticas de work-queue y pub/sub en la misma cola.
- Bloquear el event loop en un consumidor, lo que detiene los heartbeats y provoca una desconexión falsa.
Próximos pasos
Si buscas la misma idea de trabajo fiable pero con menos infraestructura, la guía de Redis Queues cubre BullMQ sobre Redis. Para profundizar en la disciplina de mover el trabajo fuera de la ruta de la solicitud, lee Batch Processing. Cuando tus eventos sean un log duradero que muchos consumidores reproducen, en lugar de tareas que se enrutan y eliminan, la guía de Kafka es el siguiente paso natural, y Node.js basics cubre el runtime en el que se ejecuta cada consumidor amqplib.