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
Producer
→
📡 campaign-events
Exchange: fanout
Exchange: fanout
→
📧 email-queue
durable · priority · DLQ
durable · priority · DLQ
📊 analytics-queue
durable
durable
💰 billing-queue
durable
durable
→
✉️ email-worker
prefetch: 50
prefetch: 50
📈 analytics-svc
prefetch: 100
prefetch: 100
💳 billing-svc
prefetch: 20
prefetch: 20
Eventos de estado (topic routing)
🔔 WebhookReceiver
Producer
Producer
→
🗺️ email-events
Exchange: topic
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
x-retry-count ≥ 3
→
💀 dlx
Exchange: direct
Exchange: direct
→
⚰️ dead-letters
inspect + replay
inspect + replay
→
🔔 Slack Alert
🔁 Replay CLI
Estado actual de colas
Checklist de producción
-
Acknowledge obligatorio en todos los consumersMensajes sin ACK bloquean la cola — usar
ch.ack(msg)ych.nack(msg) -
Prefetch configurado por tipo de consumeremail-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 idempotentesDuplicar un envío no duplica el cargo ni el email — deduplicación por
messageId -
Reconexión automática con backoffEvento
connection.on('close')+setTimeout(reconnect, 5000) -
Alerta en profundidad de colaPrometheus gauge
rabbitmq_queue_messages— alerta si email-queue > 10 000 -
Plugin delayed messages pendienteNecesario 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