O que é RabbitMQ?
O RabbitMQ é um message broker de código aberto que utiliza o AMQP (Advanced Message Queuing Protocol). Os producers publicam mensagens para ele, os consumers as recebem, e o broker é responsável por reter cada mensagem até que alguém a confirme.
A frase que melhor descreve seu design é: um broker inteligente com consumers simples. O RabbitMQ gerencia exchanges, regras de roteamento, acknowledgements, retries e dead letters. Um consumer é um programa pequeno que lê uma mensagem, executa uma tarefa e informa que ela foi concluída. Essa divisão de trabalho é o motivo pelo qual o RabbitMQ se encaixa em sistemas onde o roteamento e as garantias de entrega são mais importantes do que o throughput bruto.
Ele foi escrito em Erlang em 2007, razão pela qual é excepcionalmente bom em lidar com muitas conexões simultâneas e possui uma reputação de estabilidade. Ele executa um protocolo completo em vez de uma API de fila minimalista, e esse protocolo é a fonte tanto de seu poder quanto de sua curva de aprendizado.
O modelo AMQP
A ideia mais importante do RabbitMQ é que um producer nunca publica diretamente em uma queue. Ele publica em um exchange, e o exchange decide quais queues receberão uma cópia.
producer -> exchange --binding--> queue -> consumer
\--binding--> queue -> consumer
O modelo é composto por cinco elementos:
- Um producer abre um channel e chama
basic.publishcom o nome de um exchange, uma routing key e um corpo. - Um exchange recebe cada mensagem publicada e a roteia. Ele não armazena nada, a menos que esteja vinculado a uma queue.
- Um binding é uma regra que conecta um exchange a uma queue. Ele pode conter uma routing key ou um padrão, e um único exchange pode ter bindings para várias queues.
- Uma queue armazena as mensagens em ordem até que um consumer as consuma.
- Um consumer se inscreve em uma queue, recebe as entregas e as confirma (acknowledges).
Essa indireção é o ponto principal. Um producer que publica order.created.eu não sabe se um serviço, cinco serviços ou ninguém está ouvindo. Novos consumers são adicionados vinculando uma nova queue, sem a necessidade de alterar o producer.
Todo o trabalho acontece em um channel, que é uma conexão virtual leve multiplexada sobre uma única conexão TCP. Channels não são thread-safe, portanto, o padrão comum é utilizar um channel por tarefa ou por consumer.
import amqp from "amqplib";
const conn = await amqp.connect(process.env.AMQP_URL!);
const ch = await conn.createChannel();
Tipos de exchange e roteamento
O tipo de uma exchange determina como ela associa uma routing key aos seus bindings. Existem quatro tipos, e cada um possui um uso claro.
Uma direct exchange roteia para filas cuja binding key seja exatamente igual à routing key. Use-a para enviar uma mensagem para uma fila específica ou para um pequeno conjunto que compartilhe a mesma chave.
await ch.assertExchange("logs", "direct", { durable: true });
await ch.bindQueue("logs.errors", "logs", "error");
await ch.publish("logs", "error", body);
Uma fanout exchange ignora completamente a routing key e copia a mensagem para todas as filas vinculadas. É o pub/sub em sua forma mais pura: uma publicação, muitos consumidores independentes.
await ch.assertExchange("events", "fanout", { durable: true });
await ch.bindQueue("search-indexer", "events", "");
await ch.bindQueue("email-notifier", "events", "");
Uma topic exchange compara a routing key com um padrão. As chaves são palavras separadas por pontos. * corresponde a exatamente uma palavra e # corresponde a zero ou mais, o que torna as topic exchanges as mais flexíveis das quatro.
await ch.assertExchange("orders", "topic", { durable: true });
await ch.bindQueue("eu-orders", "orders", "order.*.eu");
await ch.bindQueue("all-orders", "orders", "order.#");
await ch.publish("orders", "order.created.eu", body);
Uma headers exchange ignora a routing key e faz a correspondência com base nos atributos do cabeçalho da mensagem. Raramente é a escolha certa, pois as topic exchanges são mais fáceis de ler e compreender, mas é útil quando o roteamento depende de vários atributos independentes.
Se você deve lembrar de apenas uma regra, que seja esta: escolha o tipo de exchange com base na pergunta que você está fazendo. “Qual fila exata?” é direct. “Todo mundo?” é fanout. “Qual família de eventos?” é topic.
Acknowledgements, nack e prefetch
A garantia de entrega do RabbitMQ é baseada em acknowledgements. Quando um consumidor recebe uma mensagem, o broker a marca como unacked, mas a mantém armazenada. A mensagem só é removida quando o consumidor chama ack. Se a conexão cair antes disso, o broker a entrega novamente.
await ch.consume("orders.created", async (msg) => {
if (!msg) return;
try {
await handleOrder(JSON.parse(msg.content.toString()));
ch.ack(msg);
} catch (err) {
ch.nack(msg, false, false); // requeue: false, dead-letter instead
}
});
nack (ou o antigo reject) aceita uma flag requeue. O requeue coloca a mensagem de volta no início da fila, o que é ideal para falhas transitórias, mas incorreto para “poison messages” que falharão permanentemente. A solução usual nesses casos é enviá-la para uma dead-letter exchange.
O Prefetch, configurado com basic.qos, limita quantas mensagens não confirmadas um consumidor pode reter simultaneamente. Sem isso, o broker envia as mensagens o mais rápido possível e um consumidor lento pode acabar com milhares delas em buffer na memória.
await ch.prefetch(20);
Configure o prefetch para um pequeno múltiplo da sua concorrência real. Se for alto demais, um único consumidor monopoliza a fila; se for baixo demais, o consumidor ficará ocioso aguardando a próxima viagem de ida e volta (round trip).
Durabilidade e garantias de entrega
Uma mensagem sobrevive ao reinício de um broker apenas se três condições forem verdadeiras, e é fácil esquecer uma delas.
- A queue deve ser declarada como
durable: true. - A message deve ser publicada com
persistent: true. - A exchange também deve ser
durable: true.
await ch.assertExchange("orders", "topic", { durable: true });
await ch.assertQueue("orders.created", { durable: true });
ch.publish("orders", "order.created.eu", body, { persistent: true });
Durabilidade trata de sobreviver a um reinício, não de garantir a entrega. Para isso, habilite os publisher confirms. Sem confirms, publish funciona no modelo fire-and-forget: se o broker cair antes de gravar a mensagem, o produtor nunca saberá.
const ch = await conn.createConfirmChannel();
ch.publish("orders", "order.created.eu", body, { persistent: true }, (err) => {
if (err) console.error("broker did not confirm", err);
else console.log("message is safely queued");
});
Mesmo com confirms e queues duráveis, a entrega é at-least-once. Um consumidor pode travar após realizar o trabalho, mas antes de enviar o ack, e o broker fará a redelivery. Exactly-once através de uma rede é efetivamente impossível, por isso o RabbitMQ torna essa escolha explícita e solicita que você torne seus handlers idempotentes.
Dead-letter exchanges e filas de retry
Um dead-letter exchange (DLX) é um exchange comum que recebe mensagens que uma fila rejeita, expira ou descarta. É o mecanismo por trás das filas de retry e do tratamento de poison-messages.
Uma mensagem é enviada para o dead-letter quando ocorre um destes eventos:
- O consumer envia um nack ou a rejeita com
requeue: false. - O TTL por mensagem ou por fila expira.
- A fila excede seu limite de comprimento e descarta a mensagem mais antiga.
Você configura o DLX na fila que detém o trabalho, e não no consumer.
await ch.assertQueue("orders.created", {
durable: true,
deadLetterExchange: "orders.dlx",
deadLetterRoutingKey: "failed",
});
O padrão clássico de delayed retry utiliza uma segunda fila com um TTL e seu próprio DLX apontando de volta para o exchange de trabalho. Uma mensagem que falhou é enviada via dead-letter para a fila de retry, permanece lá durante o TTL, expira e é enviada via dead-letter de volta para ser processada novamente. Isso gera um retry com atraso sem a necessidade de timers no seu código.
await ch.assertExchange("orders.dlx", "direct", { durable: true });
await ch.assertQueue("orders.retry", {
durable: true,
messageTtl: 30_000,
deadLetterExchange: "orders",
deadLetterRoutingKey: "order.created.retry",
});
await ch.bindQueue("orders.retry", "orders.dlx", "retry");
Uma mensagem que continua falhando entrará em loop. Monitore a contagem de retries nos headers da mensagem e, após um limite, direcione-a para uma fila de failed permanente que ninguém consome automaticamente. Trate essa fila como uma superfície operacional: configure alertas para quando ela crescer e crie um caminho de replay.
TTL de mensagens e limites de fila
O Time-to-live controla por quanto tempo uma mensagem pode aguardar. Configure-o por fila, por mensagem, ou ambos.
await ch.assertQueue("verification", {
durable: true,
messageTtl: 600_000, // 10 minutes for every message
maxLength: 10_000, // keep at most 10k messages
overflow: "reject-publish", // backpressure instead of dropping
});
O TTL por mensagem é definido no momento da publicação e é frequentemente utilizado para valores que variam entre as mensagens.
ch.publish("orders", key, body, {
expiration: "30000", // milliseconds, as a string
});
Os limites de tamanho da fila servem como backpressure. maxLength limita a memória, e overflow decide o que acontece quando o limite é atingido: drop-head descarta silenciosamente a mensagem mais antiga, enquanto reject-publish recusa novas publicações para que o produtor sinta a pressão. Para uma fila que não pode perder trabalho, reject-publish é o padrão mais seguro e serve como um sinal de alerta.
Filas de trabalho (work queues) vs pub/sub
O mesmo broker cobre dois modelos de comunicação muito diferentes, e confundi-los causa bugs.
Uma work queue distribui cada mensagem para exatamente um consumidor. Vários workers competem na mesma fila, e o broker realiza as entregas via round-robin. Adicione um prefetch e acknowledgements e você terá uma fila de tarefas escalável e tolerante a falhas.
await ch.assertQueue("jobs.thumbnails", { durable: true });
await ch.prefetch(5);
Pub/sub entrega cada mensagem para todos os consumidores interessados. Cada assinante possui sua própria fila vinculada ao mesmo exchange, portanto, um assinante lento ou offline nunca “rouba” uma mensagem dos demais.
await ch.assertExchange("events", "fanout", { durable: true });
await ch.assertQueue("events.billing", { durable: true });
await ch.bindQueue("events.billing", "events", "");
await ch.assertQueue("events.analytics", { durable: true });
await ch.bindQueue("events.analytics", "events", "");
A regra geral: uma fila compartilhada por muitos consumidores é uma work queue; uma fila por consumidor é pub/sub. Um topic exchange permite que você tenha ambos simultaneamente, com alguns consumidores compartilhando uma fila e outros possuindo a sua própria.
O padrão RPC
O RabbitMQ também suporta requisição/resposta. O cliente publica uma requisição com uma fila replyTo e um correlationId, e o servidor publica a resposta nessa fila com o mesmo id.
import { randomUUID } from "node:crypto";
const { queue } = await ch.assertQueue("", { exclusive: true });
const correlationId = randomUUID();
ch.consume(queue, (msg) => {
if (msg?.properties.correlationId === correlationId) {
console.log("reply", JSON.parse(msg.content.toString()));
}
}, { noAck: true });
ch.publish("rpc.inventory", "check", Buffer.from("{}"), {
replyTo: queue,
correlationId,
});
O assertQueue("") cria uma fila temporária, exclusiva e com auto-deleção apenas para este cliente. RPC via fila é útil quando o chamador realmente precisa de uma resposta, mas isso reintroduz o acoplamento síncrono e o problema de timeout. Para a maioria dos sistemas, um evento seguido de um evento de callback é mais fácil de operar do que RPC, e uma chamada HTTP simples é ainda mais fácil quando a dependência está disponível.
Clustering e quorum queues
Um único nó representa um ponto único de falha, por isso o RabbitMQ em produção roda como um cluster. As filas podem ser replicadas entre nós, e os clientes se reconectam a outro nó quando um deles falha.
A estratégia de replicação moderna é a quorum queue, construída sobre o algoritmo de consenso Raft. Uma quorum queue possui um líder e seguidores, e uma escrita é confirmada apenas quando a maioria a possui. Isso a torna segura contra partições de rede que faziam com que as mirrored queues clássicas perdessem dados, sendo a escolha padrão para filas duráveis hoje em dia.
await ch.assertQueue("orders.created", {
durable: true,
arguments: { "x-queue-type": "quorum" },
});
As quorum queues preferem um número pequeno e ímpar de réplicas, tipicamente três ou cinco. Elas são mais pesadas que as filas clássicas, portanto, use-as para dados que você não pode perder e mantenha as filas transitórias como clássicas. Para topologias multi-região, a federation e o plugin shovel movem mensagens entre brokers independentes, em vez de estender um único cluster através de um link lento.
A UI de gerenciamento e monitoramento
Cada node do RabbitMQ vem com um plugin de gerenciamento que fornece uma UI web e uma API HTTP. Ele exibe exchanges, queues, bindings, conexões e channels, e permite que você publique uma mensagem de teste ou reexecute uma de uma queue.
rabbitmq-plugins enable rabbitmq_management
curl -u guest:guest http://localhost:15672/api/queues/%2F/orders.created
Quatro números são os mais importantes em um dashboard.
- Queue depth — mensagens prontas. Uma profundidade crescente significa que os consumers não estão conseguindo acompanhar a demanda.
- Unacked count — mensagens entregues, mas não confirmadas. Um número que apenas cresce indica que os consumers estão travados.
- Publish and deliver rates — o formato do tráfego e se os consumers mantêm o ritmo.
- Redelivery rate — um aumento na contagem de redeliveries aponta para crashes ou falhas repetitivas.
O RabbitMQ também emite métricas do Prometheus, portanto, a queue depth e a utilização dos consumers devem estar no mesmo dashboard que o restante dos seus serviços. Configure alertas para a profundidade e a idade da mensagem mais antiga, e não para falhas individuais, que são esperadas.
Conexões, canais e recuperação
Uma conexão é cara; um canal é barato. Abra uma conexão por processo e, em seguida, crie um canal por produtor, por consumidor ou por unidade de trabalho. Um canal não é thread-safe, portanto, compartilhá-lo entre handlers concorrentes causa a intercalação de frames e erros confusos.
import amqp from "amqplib";
let conn: amqp.Connection;
let ch: amqp.Channel;
async function connect() {
conn = await amqp.connect(process.env.AMQP_URL!);
conn.on("error", (err) => console.error("connection error", err));
conn.on("close", () => setTimeout(connect, 5_000));
ch = await conn.createChannel();
await ch.assertExchange("orders", "topic", { durable: true });
}
Trate a reconexão de forma deliberada. amqplib não reconecta automaticamente para você, então monitore o close e reconstrua a conexão, o canal e cada consumidor. Como uma conexão interrompida deixa mensagens em trânsito sem confirmação (unacked), o broker as redeliver, que é o comportamento desejado. Este é outro ponto onde a entrega at-least-once se manifesta: uma reconexão pode repetir o trabalho, portanto, os handlers devem ser idempotentes.
Propriedades e prioridades de mensagens
Toda mensagem publicada pode carregar propriedades junto ao seu corpo. Elas viajam com a mensagem e ficam visíveis para os consumidores, o que as torna um local leve para armazenar metadados.
ch.publish("orders", "order.created.eu", body, {
persistent: true,
contentType: "application/json",
messageId: order.id,
correlationId: traceId,
timestamp: Math.floor(Date.now() / 1000),
type: "order.created",
headers: { "x-retry-count": 0, "x-source": "checkout" },
});
contentType e messageId auxiliam os consumidores e as ferramentas; correlationId vincula uma mensagem a um trace; headers são onde você mantém metadados da aplicação, como a contagem de tentativas (retry count). Não coloque grandes volumes de dados nos headers, pois o broker precisa indexá-los e exibi-los.
O RabbitMQ também suporta priority queues, onde mensagens de maior prioridade são entregues primeiro. Declare a fila com uma prioridade máxima e publique com um valor de priority.
await ch.assertQueue("jobs", {
durable: true,
maxPriority: 10,
});
ch.publish("", "jobs", body, { priority: 9 });
As prioridades apenas reordenam as mensagens que já estão aguardando. Se os consumidores mantiverem a fila vazia, a prioridade não terá efeito, e uma mensagem de alta prioridade que chegue após uma de baixa prioridade ainda terá que esperar por ela. Use a prioridade para distinções reais de negócio, e não como um mecanismo geral de agendamento.
Escolhendo entre RabbitMQ e Kafka
Os dois são frequentemente comparados, mas resolvem problemas diferentes.
O RabbitMQ é um smart broker para roteamento e filas de trabalho. As mensagens são tarefas que são removidas assim que confirmadas (acknowledged). Os exchanges fazem o roteamento por padrão, as confirmações são por mensagem e o dead-lettering já vem integrado. Ele se destaca quando uma mensagem deve chegar a um conjunto específico de filas, ou quando uma unidade de trabalho deve ser repetida e estacionada em caso de falha.
O Kafka é um durable log. As mensagens são anexadas a um log ordenado e particionado, sendo retidas por um tempo configurado, independentemente de quem as leia. Diversos consumer groups podem ler o mesmo tópico de forma independente, e qualquer consumidor pode retroceder para um offset anterior. Ele se destaca em throughputs muito altos e no replay de eventos.
Uma heurística útil: se você ficaria chateado por uma mensagem ter sido consumida e deletada, você precisa de um log. Se você se importa que uma tarefa tenha sido roteada corretamente e concluída, você precisa de um broker. Muitos sistemas utilizam ambos — Kafka para o event stream e RabbitMQ para o trabalho.
Melhores práticas
- Declare exchanges, queues e bindings de forma idempotente na inicialização, e use
durable: truepara qualquer dado que não possa ser perdido. - Publique mensagens persistentes em queues duráveis e utilize publisher confirms para tarefas que não podem desaparecer.
- Sempre faça o acknowledge manualmente após o sucesso do efeito colateral, nunca antes.
- Defina um limite de prefetch para que um único consumer não faça o buffer de toda a queue.
- Atribua a cada queue durável uma dead-letter exchange e construa um caminho de replay para ela.
- Implemente retentativas atrasadas com uma retry queue de TTL em vez de usar sleep no handler.
- Torne os consumers idempotentes, pois a entrega é do tipo at-least-once.
- Limite o comprimento da queue e escolha
reject-publishquando descartar mensagens for inaceitável. - Use quorum queues para dados duráveis e classic queues para dados transitórios.
- Feche channels e conexões no
SIGTERM, e pare de consumir antes de realizar o draining. - Monitore ready, unacked, a taxa de redelivery e a idade da mensagem mais antiga.
Erros comuns
- Publicar em uma queue em vez de um exchange e depois se perguntar por que os bindings não funcionam.
- Usar auto-ack (
noAck: true) e perder mensagens sempre que um handler crasha. - Esquecer o
persistent: truee perder mensagens ao reiniciar o broker. - Declarar uma queue como durable, mas publicar mensagens não persistentes.
- Definir o prefetch com um valor muito alto, fazendo com que um único consumer “mate” os demais.
- Fazer o requeue de uma poison message infinitamente com
nack(msg, false, true). - Criar um dead-letter exchange sem nenhum consumer e nunca monitorá-lo.
- Usar um fanout exchange onde era necessário um topic exchange, ou vice-versa.
- Rodar um único broker em produção e dizer que ele é altamente disponível.
- Misturar semânticas de work-queue e pub/sub na mesma queue.
- Bloquear o event loop em um consumer, o que interrompe os heartbeats e dispara uma desconexão falsa.
Próximos passos
Se você busca a mesma ideia de processamento confiável, mas com menos infraestrutura, o guia de Redis Queues aborda o BullMQ sobre o Redis. Para entender a disciplina mais ampla de remover o processamento do caminho da requisição, leia sobre Batch Processing. Quando seus eventos são um log durável que muitos consumidores reproduzem, em vez de tarefas que são roteadas e removidas, o guia de Kafka é o próximo passo natural, e o de Node.js basics cobre o runtime em que cada consumidor amqplib é executado.