CULTIVA IA
FlowSend — Arquitectura RabbitMQ
Diseño de mensajería para microservicios · Plataforma Email Marketing B2B
RabbitMQ 3.13 Node.js 20 · TypeScript 5 2 M emails/día 5 Microservicios
Métricas del sistema
2M
Emails procesados/día
▲ +18% vs. mes anterior
99.97%
Uptime mensajería
▲ SLA objetivo: 99.9%
~4ms
Latencia media de cola
▲ Reducción de 380ms (monolito)
12
Dead-letters hoy
▼ 0.0006% tasa de fallo
Diagrama de flujo de mensajes
Flujo campañas masivas (fanout)
📤 CampaignService
Producer
📡 campaign-events
Exchange: fanout
📧 email-queue
durable · priority · DLQ
📊 analytics-queue
durable
💰 billing-queue
durable
✉️ email-worker
prefetch: 50
📈 analytics-svc
prefetch: 100
💳 billing-svc
prefetch: 20
Eventos de estado (topic routing)
🔔 WebhookReceiver
Producer
🗺️ email-events
Exchange: topic
email.delivered
delivery-log
email.opened
engagement-track
email.bounced
bounce-handler
#
audit-log
📊 AnalyticsSvc
📊 AnalyticsSvc
🧹 ListCleaner
📝 AuditLogger
Dead-letter queue (reintentos agotados)
📧 email-queue
x-retry-count ≥ 3
💀 dlx
Exchange: direct
⚰️ dead-letters
inspect + replay
🔔 Slack Alert
🔁 Replay CLI
Estado actual de colas
Cola Exchange / Tipo Profundidad Estado
email-queue fanout
3,621
Activa
analytics-queue fanout
1,200
Activa
billing-queue fanout
587
Activa
delivery-log topic
312
Activa
bounce-handler topic
47
Procesando
dead-letters DLQ
12
Revisar
Checklist de producción
  • Acknowledge obligatorio en todos los consumers
    Mensajes sin ACK bloquean la cola — usar ch.ack(msg) y ch.nack(msg)
  • Prefetch configurado por tipo de consumer
    email-worker: 50 · analytics: 100 · billing: 20 · sin prefetch = un consumer acapara todo
  • Colas durable + mensajes persistent
    durable: true + persistent: true — sobreviven reinicios del broker
  • Dead-letter exchange en todas las colas críticas
    x-dead-letter-exchange: 'dlx' — ningún mensaje falla en silencio
  • Consumers idempotentes
    Duplicar un envío no duplica el cargo ni el email — deduplicación por messageId
  • Reconexión automática con backoff
    Evento connection.on('close') + setTimeout(reconnect, 5000)
  • Alerta en profundidad de cola
    Prometheus gauge rabbitmq_queue_messages — alerta si email-queue > 10 000
  • ⚠️
    Plugin delayed messages pendiente
    Necesario para scheduled campaigns — habilitar rabbitmq_delayed_message_exchange
Código TypeScript — producción
src/messaging/campaign-producer.ts TypeScript
import amqp, { Channel } from 'amqplib';
import { randomUUID } from 'crypto';

// Conexión compartida con reconexión automática
let channel: Channel;

async function getChannel(): Promise<Channel> {
  if (channel) return channel;

  const conn = await amqp.connect(
    'amqp://admin:secret@rabbitmq:5672'
  );
  channel = await conn.createChannel();
  await channel.prefetch(50);

  conn.on('close', () => {
    channel = null!;
    setTimeout(getChannel, 5000); // backoff 5s
  });

  return channel;
}

/** Publica campaña masiva al exchange fanout */
export async function publishCampaign(
  campaignId: string,
  recipients: string[],
  template: string,
  priority = 0
) {
  const ch = await getChannel();

  await ch.assertExchange(
    'campaign-events', 'fanout',
    { durable: true }
  );

  for (const email of recipients) {
    ch.publish(
      'campaign-events', '',
      Buffer.from(JSON.stringify({
        campaignId, email, template,
        scheduledAt: new Date().toISOString(),
      })),
      {
        persistent: true,
        priority,
        contentType: 'application/json',
        messageId: randomUUID(),
        timestamp: Date.now(),
        headers: { 'x-retry-count': 0 },
      }
    );
  }

  console.log(`[CampaignProducer] ${recipients.length} mensajes publicados`);
}
src/workers/email-worker.ts TypeScript
import { getChannel } from '../messaging/connection';
import { sendEmail } from '../services/mailer';

const QUEUE = 'email-queue';
const MAX_RETRIES = 3;

export async function startEmailWorker() {
  const ch = await getChannel();

  // DLQ: mensajes fallidos van a 'dlx'
  await ch.assertExchange('dlx', 'direct', { durable: true });
  await ch.assertQueue('dead-letters', { durable: true });
  await ch.bindQueue('dead-letters', 'dlx', '');

  await ch.assertQueue(QUEUE, {
    durable: true,
    arguments: {
      'x-max-priority': 10,
      'x-dead-letter-exchange': 'dlx',
      'x-message-ttl': 86_400_000, // 24h
    },
  });

  ch.consume(QUEUE, async (msg) => {
    if (!msg) return;

    const task = JSON.parse(msg.content.toString());
    const retries = msg.properties.headers?.['x-retry-count'] ?? 0;

    try {
      await sendEmail(task.email, task.template, task);
      ch.ack(msg); // ✓ éxito
    } catch (err) {
      if (retries >= MAX_RETRIES) {
        // Enviar a DLQ tras 3 intentos
        ch.nack(msg, false, false);
        console.error(`[DLQ] ${task.email} —`, err);
      } else {
        // Reintentar con contador incrementado
        ch.nack(msg, false, false);
        ch.sendToQueue(QUEUE, msg.content, {
          ...msg.properties,
          headers: {
            ...msg.properties.headers,
            'x-retry-count': retries + 1,
          },
        });
      }
    }
  });

  console.log(`[EmailWorker] Escuchando en ${QUEUE}`);
}
src/monitoring/dlq-monitor.ts  —  Monitor de dead-letters con alerta Slack TypeScript
import { getChannel } from '../messaging/connection';
import { sendSlackAlert } from '../services/slack';

export async function startDLQMonitor() {
  const ch = await getChannel();

  ch.consume('dead-letters', async (msg) => {
    if (!msg) return;
    const payload = JSON.parse(msg.content.toString());
    const deathReason = msg.properties.headers?.['x-death']?.[0]?.reason ?? 'unknown';
    const deathCount = msg.properties.headers?.['x-death']?.[0]?.count ?? 0;

    console.error(`[DLQ] reason=${deathReason} count=${deathCount}`, payload);

    await sendSlackAlert({
      channel: '#ops-alertas',
      text: `⚰️ *Dead letter* en email-queue\n` +
            `• Destinatario: \`${payload.email}\`\n` +
            `• Campaña: \`${payload.campaignId}\`\n` +
            `• Razón: ${deathReason} (${deathCount} intentos)\n` +
            `• Ver: https://rabbitmq.flowsend.io/#/queues/%2F/dead-letters`,
    });

    ch.ack(msg); // ack para no bloquear el monitor
  });
}
Patrones implementados en FlowSend
📡
Fanout — Broadcast
Cada mensaje de campaña llega simultáneamente a email-service, analytics-service y billing-service. Ningún servicio bloquea a otro.
campaign-events exchange
🗺️
Topic — Enrutamiento selectivo
Webhooks de SendGrid llegan al exchange y se enrutan: email.bounced solo al limpiador de listas, # al audit log.
email-events exchange
Priority Queue — Urgencias primero
Emails transaccionales (reset de contraseña, confirmación de pago) se publican con priority: 9, saltando la cola de campañas.
x-max-priority: 10
💀
Dead-letter Queue — Fallos controlados
Tras 3 reintentos con backoff, el mensaje pasa a dead-letters. El monitor envía alerta a Slack con campaignId y destinatario.
dlx exchange · x-death headers
🔄
Work Queue — Balanceo de carga
Múltiples instancias de email-worker consumen de la misma cola. RabbitMQ distribuye round-robin respetando el prefetch de 50.
prefetch: 50 · N workers
🕐
Delayed Messages — Campañas programadas
Plugin rabbitmq_delayed_message_exchange permite publicar con x-delay: 3600000 para campañas con hora de envío definida.
x-delayed-message · PENDIENTE