🔗 Mensajería y microservicios en Node.js
Colas con BullMQ, RabbitMQ, NestJS microservices, patrones saga y outbox, Hono y edge computing, y un flujo completo de job encolado con retries.
Mensajería y microservicios en Node.js
Cuando una aplicación crece, dividirla en servicios independientes que se comunican por mensajes (en lugar de llamadas HTTP síncronas) resuelve problemas de acoplamiento, escalado y resiliencia. En este artículo dominas las colas de trabajo con BullMQ, el broker RabbitMQ, los microservicios de NestJS, los patrones de consistencia (saga, outbox) y el edge computing con Hono.
REST vs message-driven
La decisión fundamental es cómo se comunican los servicios:
| Aspecto | REST (síncrono) | Message-driven (asíncrono) |
|---|---|---|
| Acoplamiento | El productor conoce al consumidor | Bajo: nadie se conoce |
| Disponibilidad | Si el consumidor cae, falla | La cola absorbe la caída |
| Latencia | Inmediata | Diferida (cola) |
| Garantías | Difíciles de retomar | Retries y dead-letter |
| Caso de uso | Queries/CRUD en vivo | Tareas pesadas, eventos, integraciones |
💡 Regla práctica: deja síncrono lo que el usuario espera en la respuesta (leer datos) y mueve a cola lo que puede esperar o puede fallar (emails, PDFs, indexado, pagos externos).
BullMQ: colas de trabajo con Redis
BullMQ es la librería de colas más usada en el ecosistema Node. Se apoya en Redis y distingue dos roles: el productor encola jobs y el worker los procesa.
El productor (encola un job)
import { Queue } from "bullmq";
const colaEmails = new Queue("emails", {
connection: { host: "localhost", port: 6379 },
});
await colaEmails.add("enviar-email", {
to: "ana@dev.io",
subject: "Bienvenida",
body: "Hola Ana",
}, {
attempts: 3, // reintentos automáticos
backoff: { type: "exponential", delay: 5000 },
removeOnComplete: true,
});
colaEmails.add(nombre, datos, opciones): el nombre permite rutear a distintos procesos.attempts: cuántas veces se reintenta si el worker lanza un error.backoff: espera entre reintentos (exponencial condelaybase).
El worker (procesa el job)
import { Worker } from "bullmq";
import { sendEmail } from "./email";
const worker = new Worker(
"emails",
async (job) => {
const { to, subject, body } = job.data;
console.log(`Procesando job ${job.id}`);
await sendEmail(to, subject, body);
},
{ connection: { host: "localhost", port: 6379 } },
);
worker.on("completed", (job) => console.log(`ok: ${job.id}`));
worker.on("failed", (job, err) => console.error(`fallo: ${job.id}`, err.message));
worker.on("error", (err) => console.error("error de worker", err));
⚠️ Si el worker lanza una excepción, BullMQ reintenta según
attempts; cuando se agotan, el job pasa a dead-letter (estadofailed) para revisión manual. No dejes jobs muertos sin unworker.on("failed")que los loguee.
Cron jobs y repetición
BullMQ también programa tareas periódicas:
import { Queue } from "bullmq";
const colaReportes = new Queue("reportes", { connection: { host: "localhost", port: 6379 } });
await colaReportes.upsertJobScheduler(
"reporte-diario",
{ pattern: "0 8 * * *" }, // cron: 8:00 cada día
{ name: "generar-reporte", data: { tipo: "ventas" } },
);
RabbitMQ: broker con exchanges y queues
RabbitMQ es un message broker maduro (protocolo AMQP). A diferencia de BullMQ (cola simple), RabbitMQ introduce los exchanges: un componente que recibe mensajes y los enruta a una o más colas según reglas.
import amqplib from "amqplib";
const conn = await amqplib.connect("amqp://localhost");
const ch = await conn.createChannel();
// Declarar exchange de tipo topic y una cola
await ch.assertExchange("eventos", "topic", { durable: true });
await ch.assertQueue("pedidos.cobrados", { durable: true });
await ch.bindQueue("pedidos.cobrados", "eventos", "pedido.cobrado.*");
// Publicar un mensaje enrutado
ch.publish("eventos", "pedido.cobrado.123", Buffer.from(JSON.stringify({ pedidoId: 123 })));
El consumidor escucha la cola:
import amqplib from "amqplib";
const conn = await amqplib.connect("amqp://localhost");
const ch = await conn.createChannel();
await ch.assertQueue("pedidos.cobrados", { durable: true });
ch.consume("pedidos.cobrados", (msg) => {
if (!msg) return;
const evento = JSON.parse(msg.content.toString());
console.log("Pedido cobrado:", evento);
ch.ack(msg); // confirmar: solo borrar tras procesar bien
}, { noAck: false });
| Tipo de exchange | Enrutamiento |
|---|---|
direct |
Clave exacta (pedido.cobrado) |
topic |
Patrón con * (una palabra) y # (varias) |
fanout |
A todas las colas enlazadas |
headers |
Por cabeceras del mensaje |
💡
ch.ack(msg)confirma el mensaje: si el consumidor muere sin confirmar, RabbitMQ re-entrega. ConnoAck: trueel mensaje se pierde si el proceso cae — para trabajo importante usanoAck: falsey confirma explícitamente.
NestJS microservices
NestJS no solo hace HTTP: sus microservicios usan el mismo modelo de decoradores con distintos transportes (TCP, Redis, NATS, Kafka, gRPC). El servicio escucha con @MessagePattern y el cliente envía con ClientProxy.
// pedidos.service.ts — dentro de un microservicio
import { Controller } from "@nestjs/common";
import { MessagePattern, Payload } from "@nestjs/microservices";
@Controller()
export class PedidosController {
@MessagePattern({ cmd: "cobrar" }) // patrón que atiende
async cobrar(@Payload() data: { pedidoId: number }) {
return { cobrado: true, pedidoId: data.pedidoId };
}
}
// cliente — se conecta por el transporte TCP
import { Injectable } from "@nestjs/common";
import { ClientProxy, ClientProxyFactory, Transport } from "@nestjs/microservices";
@Injectable()
export class PedidosClient {
private client: ClientProxy = ClientProxyFactory.create({
transport: Transport.TCP,
options: { host: "localhost", port: 3001 },
});
async cobrar(pedidoId: number) {
return this.client.send({ cmd: "cobrar" }, { pedidoId }).toPromise();
}
}
💡 La ventaja de NestJS microservices: el código del servicio es idéntico cambies el transporte. Switches entre TCP, Redis o Kafka cambiando solo la configuración del
ClientProxyy del módulo de escucha.
Patrones: saga, outbox y event-driven
Saga (orquestación simple)
Una operación distribuida que toca varios servicios (cobrar + reservar stock + notificar) no puede ser una transacción atómica. La saga la divide en pasos con compensación: si un paso falla, se ejecutan los pasos inversos.
async function sagaCompra(pedidoId: number) {
const pasos = [
{ hacer: () => cobrar(pedidoId), compensar: () => reembolsar(pedidoId) },
{ hacer: () => reservarStock(pedidoId), compensar: () => liberarStock(pedidoId) },
];
const ejecutados: typeof pasos = [];
for (const paso of pasos) {
try {
await paso.hacer();
ejecutados.push(paso);
} catch (e) {
for (const hecho of ejecutados.reverse()) await hecho.compensar();
throw e;
}
}
}
Patrón outbox
Cuando guardas un evento de dominio y a la vez lo publicas a una cola, si publicas antes de commitear y el commit falla, tienes un mensaje fantasma. El outbox garantiza “publicar al menos una vez”: escribes el evento en una tabla outbox en la misma transacción que los datos de negocio, y un relay lo publica a la cola después.
BEGIN;
INSERT INTO pedidos (id, total) VALUES (1, 100);
INSERT INTO outbox (id, agregado, payload) VALUES (gen_random_uuid(), 'pedido', '{"id":1}');
COMMIT;
Un worker lee de outbox, publica a la cola y marca la fila como enviada. La consistencia entre DB y cola viene del hecho de compartir la transacción de la base de datos.
⚠️ El patrón outbox no elimina las entregas duplicadas: da “al menos una vez”. Tu consumidor debe ser idempotente (procesar el mismo evento dos veces produce el mismo resultado).
Event-driven
En el diseño event-driven los servicios publican eventos del pasado (“pedido.cobrado”) y los demás reaccionan sin que nadie orqueste. Esto desacopla productores de consumidores y permite añadir nuevos consumidores sin tocar el emisor.
Hono y el edge computing
No todos los servicios viven en un servidor largo. El edge computing ejecuta lógica en el borde de la red (cerca del usuario) en runtimes como Cloudflare Workers, Deno o Bun. Hono es un framework web ultra-ligero que corre en todos ellos con el mismo código.
// worker de Cloudflare: responde en el edge
import { Hono } from "hono";
const app = new Hono();
app.get("/", (c) => c.text("Hola desde el edge"));
app.get("/usuarios/:id", (c) => c.json({ id: c.req.param("id"), nombre: "Ana" }));
export default app;
// Bun: runtime con APIs nativas rápidas
Bun.serve({
port: 3000,
fetch(req) {
return new Response(`Hola ${req.url}`);
},
});
| Runtime | Dónde corre | Característica clave |
|---|---|---|
| Node.js | Servidor propio / VPS | Madurez, ecosistema npm completo |
| Cloudflare Workers | Edge (CDN) | Arranque en frío mínimo, escala global |
| Deno | Edge / server | Seguridad, TypeScript nativo, APIs web |
| Bun | Server | Rendimiento, Bun.serve, bundler integrado |
💡 Edge no sustituye a Node: es para respuestas rápidas, cache y autorización cerca del usuario. La lógica pesada (colas, DB relacional, jobs) sigue viviendo en Node. Hono te permite compartir código entre ambos mundos.
Ejemplo completo: flujo con BullMQ
Montamos un flujo real: una API encola un job de “procesar pedido” y un worker lo procesa con reintentos. Separamos en dos procesos para ver la arquitectura.
Producer (producer.ts):
import { Queue } from "bullmq";
const cola = new Queue("pedidos", {
connection: { host: "localhost", port: 6379 },
});
export async function encolarPedido(pedidoId: number) {
return cola.add(
"procesar",
{ pedidoId },
{
attempts: 3,
backoff: { type: "exponential", delay: 2000 },
removeOnComplete: { age: 3600 },
},
);
}
// simula un request
await encolarPedido(42);
console.log("Pedido 42 encolado");
Worker (worker.ts):
import { Worker } from "bullmq";
async function procesarPedido(pedidoId: number) {
console.log(`Procesando pedido ${pedidoId}...`);
if (pedidoId % 2 === 0) throw new Error("saldo insuficiente"); // fuerza reintento
console.log(`Pedido ${pedidoId} completado`);
}
const worker = new Worker(
"pedidos",
async (job) => {
await procesarPedido(job.data.pedidoId);
},
{ connection: { host: "localhost", port: 6379 }, concurrency: 5 },
);
worker.on("completed", (job) => console.log(`OK ${job.id} tras ${job.attemptsMade + 1} intentos`));
worker.on("failed", (job, err) => console.error(`FALLÓ ${job.id}: ${err.message}`));
Cómo funciona el reintento:
- El worker lanza un error → el job se retrasa según
backoff. - Se reintenta hasta
attemptsveces (aquí 3). - Si agota los intentos, queda en estado
failed(dead-letter). job.attemptsMadete dice en qué intento está.
⚠️ Separa productor y worker en procesos o deploys distintos. Así puedes escalar el número de workers sin tocar la API, y la API nunca bloquea por el trabajo pesado.
Buenas prácticas en mensajería
- Idempotencia: el consumidor debe tolerar duplicados.
- Retries con backoff exponencial y dead-letter para errores permanentes.
- Payloads pequeños y versionados: añade un campo
versiona tus eventos. - Confirma mensajes solo tras procesar correctamente (RabbitMQ
ack, BullMQ completar). - Monitoriza colas: longitud, jobs fallidos y latencia de procesado.
- Graceful shutdown: detén el worker antes de matar el proceso para no perder jobs a medias.
Para profundizar
- BullMQ — Documentación oficial: jobs, workers, schedulers, retries y dead-letter en profundidad.
- RabbitMQ — Tutorials: el tutorial “Work Queues” y “Topics” explican exchanges y acuses.
- NestJS — Microservices: transportes TCP/Redis/NATS/Kafka/gRPC y patterns.
- Hono — Documentación: framework para edge en Cloudflare Workers, Deno y Bun.
- Cloudflare Workers: la plataforma de edge computing y sus runtimes.