Voltar para Artigos
Back-end★ Destaque12 min de leitura

Escalando com Mensageria: RabbitMQ com Exchanges, Dead Letter Queues e Retry

Direct, Fanout e Topic Exchanges. Dead Letter Queue para mensagens que falham. Retry com backoff exponencial. Prefetch e confirmação manual de mensagens. Quando usar RabbitMQ vs Bull/BullMQ.

12 de agosto de 2026

Emitir 1.000 notas fiscais em loop síncrono dentro de uma requisição HTTP resulta em timeout. Enviar e-mails de forma síncrona bloqueia a resposta da API. Processar um vídeo de 500MB durante o request mata o servidor. Esses problemas têm a mesma raiz: acoplamento temporal — o cliente precisa esperar o trabalho pesado terminar para receber a resposta.

Mensageria desacopla o produtor do consumidor: a API publica uma mensagem ('processe esta NF') e retorna imediatamente (202 Accepted). Workers em background pegam a mensagem e processam assincronamente. O RabbitMQ garante que nenhuma mensagem seja perdida, mesmo que os workers caiam — e com Dead Letter Queues, você captura e reprocessa as que falharem.

Conceitos Fundamentais do RabbitMQ

  • Producer — publica mensagens. Normalmente a sua API.
  • Exchange — recebe mensagens do producer e roteia para as queues baseado em regras.
  • Queue — armazena as mensagens até um consumer processá-las.
  • Consumer — processa mensagens da queue. Pode haver múltiplos consumers na mesma queue (distribuição de carga).
  • Binding — regra que conecta uma exchange a uma queue.
  • Acknowledgement (ack) — confirmação de que a mensagem foi processada com sucesso. Sem ack, o RabbitMQ recoloca a mensagem na fila.

Tipos de Exchange

  • Direct — roteia para a queue cujo binding key é exatamente igual à routing key da mensagem. Use para roteamento direto por tipo.
  • Fanout — envia para TODAS as queues vinculadas, ignorando a routing key. Use para broadcast (ex: notificar todos os serviços de um evento).
  • Topic — roteia por padrão com wildcards (* = uma palavra, # = zero ou mais). Use quando um consumer precisa de múltiplos tipos de mensagem.
  • Headers — roteia por atributos do header da mensagem. Raro, use apenas quando routing key não é suficiente.
bash
npm install amqplib
npm install -D @types/amqplib
RabbitMQClient.ts
import { connect, type Connection, type Channel, type ConsumeMessage } from 'amqplib';
import { logger } from '../logger';

export class RabbitMQClient {
  private connection: Connection | null = null;
  private channel: Channel | null = null;

  async connect(): Promise<void> {
    this.connection = await connect(process.env.RABBITMQ_URL as string);
    this.channel = await this.connection.createChannel();

    // Prefetch: limita quantas mensagens o consumer pega de uma vez
    // Sem isso, o consumer pega TODAS as mensagens da queue de uma vez,
    // sobrecarregando o processo com processamento em paralelo sem controle
    await this.channel.prefetch(5); // Processa até 5 mensagens por vez

    this.connection.on('error', (err) => {
      logger.error('RabbitMQ connection error', err);
      this.reconnect();
    });

    logger.info('Conectado ao RabbitMQ');
  }

  private async reconnect(): Promise<void> {
    await new Promise((resolve) => setTimeout(resolve, 5000));
    await this.connect();
  }

  getChannel(): Channel {
    if (!this.channel) throw new Error('Canal RabbitMQ não inicializado');
    return this.channel;
  }

  async close(): Promise<void> {
    await this.channel?.close();
    await this.connection?.close();
  }
}

Dead Letter Queue: Tratando Falhas

O prefetch(5) é um detalhe de configuração crítico: sem ele, o RabbitMQ entrega todas as mensagens disponíveis ao consumer de uma vez. Imagine 10.000 mensagens na fila e um worker que só processa 10/segundo — a memória do processo Node.js explodiria. O prefetch limita que o consumer nunca tenha mais de N mensagens sem ACK ao mesmo tempo, funcionando como um throttle automático.

Sem uma DLQ, uma mensagem que falha repetidamente fica em loop infinito na fila, consumindo recursos. Configure uma Dead Letter Exchange para capturar mensagens que falharam após N tentativas:

setupQueues.ts
import type { Channel } from 'amqplib';

export async function setupQueues(channel: Channel): Promise<void> {
  // 1. Declara a Dead Letter Exchange
  await channel.assertExchange('dlx.notifications', 'direct', { durable: true });

  // 2. Declara a Dead Letter Queue (onde as mensagens falhas vão parar)
  await channel.assertQueue('notifications.dead-letter', {
    durable: true,
  });

  // 3. Vincula DLQ à DLX
  await channel.bindQueue('notifications.dead-letter', 'dlx.notifications', 'notifications');

  // 4. Declara a fila principal com referência à DLX
  await channel.assertQueue('notifications', {
    durable: true, // Sobrevive ao restart do RabbitMQ
    arguments: {
      'x-dead-letter-exchange': 'dlx.notifications',   // Para onde vai se falhar
      'x-dead-letter-routing-key': 'notifications',     // Routing key da DLQ
      'x-message-ttl': 30 * 60 * 1000,                 // TTL: 30 minutos
      'x-max-retries': 3,                               // Tentativas máximas
    },
  });

  // 5. Exchange principal
  await channel.assertExchange('notifications', 'topic', { durable: true });
  await channel.bindQueue('notifications', 'notifications', 'notification.*');
}

// Quando uma mensagem é nack() sem requeue, ela vai para a DLQ automaticamente
// A DLQ pode ter um consumer separado para analisar falhas e disparar alertas

Publisher com Confirmação

NotificationPublisher.ts
import type { Channel } from 'amqplib';

interface NotificationMessage {
  userId: string;
  type: 'email' | 'push' | 'sms';
  title: string;
  body: string;
  metadata?: Record<string, unknown>;
}

export class NotificationPublisher {
  constructor(private channel: Channel) {}

  async publish(message: NotificationMessage): Promise<void> {
    // Routing key com tipo: 'notification.email', 'notification.push', etc.
    const routingKey = `notification.${message.type}`;

    // persistent: true — mensagem sobrevive ao restart do RabbitMQ
    const published = this.channel.publish(
      'notifications',
      routingKey,
      Buffer.from(JSON.stringify(message)),
      {
        persistent: true,
        contentType: 'application/json',
        // Timestamp para auditoria e TTL do lado do consumidor
        timestamp: Date.now(),
        messageId: crypto.randomUUID(),
      }
    );

    // publish() retorna false se o buffer interno do canal está cheio
    // Neste caso, aguarde o evento 'drain' antes de publicar mais
    if (!published) {
      await new Promise<void>((resolve) => this.channel.once('drain', resolve));
    }
  }
}

Consumer com Retry e Backoff

O noAck: false é obrigatório em produção. Com noAck: true, o RabbitMQ remove a mensagem da fila assim que a entrega — se o worker cair no meio do processamento, a mensagem se perde para sempre. Com confirmação manual (channel.ack(msg)), a mensagem só sai da fila quando o processamento completou com sucesso. O channel.nack(msg, false, false) manda para a DLQ sem requeue direto, evitando loops infinitos.

NotificationWorker.ts
import type { Channel, ConsumeMessage } from 'amqplib';
import { logger } from '../infra/logger';

export class NotificationWorker {
  constructor(private channel: Channel) {}

  async start(): Promise<void> {
    await this.channel.consume('notifications', this.handleMessage.bind(this), {
      noAck: false, // Confirmação manual — NUNCA use noAck: true em produção
    });

    logger.info('NotificationWorker iniciado');
  }

  private async handleMessage(msg: ConsumeMessage | null): Promise<void> {
    if (!msg) return;

    const attempts = (msg.properties.headers?.['x-attempts'] as number) ?? 0;
    const MAX_RETRIES = 3;

    try {
      const payload = JSON.parse(msg.content.toString());
      await this.processNotification(payload);

      // Confirma: mensagem processada com sucesso e pode ser removida da fila
      this.channel.ack(msg);
    } catch (err) {
      logger.error('Falha ao processar notificação', err, {
        attempts,
        messageId: msg.properties.messageId,
      });

      if (attempts < MAX_RETRIES) {
        // Retry com backoff exponencial: 1s, 2s, 4s
        const delayMs = Math.pow(2, attempts) * 1000;

        setTimeout(() => {
          // Recoloca na fila com contador de tentativas incrementado
          this.channel.nack(msg, false, false); // false, false = não requeue direto
          // Republica com header atualizado
          this.channel.publish(
            'notifications',
            msg.fields.routingKey,
            msg.content,
            {
              ...msg.properties,
              headers: { ...msg.properties.headers, 'x-attempts': attempts + 1 },
            }
          );
        }, delayMs);
      } else {
        // Esgotou as tentativas: vai para a Dead Letter Queue
        logger.error('Mensagem enviada para DLQ após 3 tentativas', undefined, {
          messageId: msg.properties.messageId,
        });
        this.channel.nack(msg, false, false);
      }
    }
  }

  private async processNotification(payload: unknown): Promise<void> {
    // Lógica de envio de notificação
  }
}

RabbitMQ vs Bull/BullMQ: Para processamento de jobs em background em um monólito, o BullMQ (baseado em Redis) é mais simples — sem servidor adicional, UI embutida, e API de alto nível. O RabbitMQ brilha em arquiteturas de microsserviços onde diferentes serviços (em linguagens diferentes) precisam se comunicar via protocolo AMQP, com roteamento complexo por Exchange/Binding e garantias de entrega mais sofisticadas.

Conclusão

RabbitMQ resolve o acoplamento temporal entre produtor e consumidor com durabilidade de mensagens, confirmação manual de processamento e Dead Letter Queues para capturar falhas. O prefetch controla o paralelismo do consumer. O retry com backoff exponencial evita sobrecarregar serviços em falha. A combinação de exchanges tipadas com DLQ é o padrão de mensageria para sistemas que não podem perder mensagens.