🔎 Buscar

🔗 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.

Wiki / Apuntes📖 Contenido

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 con delay base).

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 (estado failed) para revisión manual. No dejes jobs muertos sin un worker.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. Con noAck: true el mensaje se pierde si el proceso cae — para trabajo importante usa noAck: false y 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 ClientProxy y 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:

  1. El worker lanza un error → el job se retrasa según backoff.
  2. Se reintenta hasta attempts veces (aquí 3).
  3. Si agota los intentos, queda en estado failed (dead-letter).
  4. job.attemptsMade te 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

  1. Idempotencia: el consumidor debe tolerar duplicados.
  2. Retries con backoff exponencial y dead-letter para errores permanentes.
  3. Payloads pequeños y versionados: añade un campo version a tus eventos.
  4. Confirma mensajes solo tras procesar correctamente (RabbitMQ ack, BullMQ completar).
  5. Monitoriza colas: longitud, jobs fallidos y latencia de procesado.
  6. Graceful shutdown: detén el worker antes de matar el proceso para no perder jobs a medias.

Para profundizar

Estudio · Recursos de todo el mundo (inglés, chino, japonés, español, francés, ruso…) curados y traducidos al español.