Sua primeira fila: publicando e consumindo
2 min de leitura
Broker rodando, teoria na cabeça. Agora a gente conecta código nele.
Os exemplos aqui são em Node com a biblioteca amqplib, mas a ideia é
igual em qualquer linguagem - muda o nome dos métodos, não o conceito.
npm install amqplib
Publicando
São sempre os mesmos quatro passos: conecta, abre um canal, garante que a estrutura existe, publica.
import amqp from "amqplib"
const conn = await amqp.connect("amqp://localhost")
const channel = await conn.createChannel()
await channel.assertExchange("pedidos", "direct", { durable: true })
await channel.assertQueue("fila.pagamento", { durable: true })
await channel.bindQueue("fila.pagamento", "pedidos", "order.created")
channel.publish(
"pedidos",
"order.created",
Buffer.from(JSON.stringify({ pedidoId: 42 }))
)
await channel.close()
await conn.close()
Repara nos três assert: eles criam se não existir, e não fazem nada
se já existir. Rodar duas vezes não quebra. É assim que se garante a
estrutura na subida da aplicação, sem depender de alguém ter clicado no
painel.
E note que a mensagem vai como Buffer - o RabbitMQ trafega bytes, não objeto. Serializar pra JSON e converter pra Buffer é responsabilidade sua.
Consumindo
O consumidor é um processo separado que fica escutando:
import amqp from "amqplib"
const conn = await amqp.connect("amqp://localhost")
const channel = await conn.createChannel()
await channel.assertQueue("fila.pagamento", { durable: true })
channel.consume("fila.pagamento", (msg) => {
if (!msg) return
const conteudo = JSON.parse(msg.content.toString())
console.log("processando pedido", conteudo.pedidoId)
channel.ack(msg)
})
O consume não é um loop que você controla: você registra uma função e
o RabbitMQ chama ela toda vez que chega mensagem. O processo fica de pé
esperando - se ele terminar, para de consumir.
O channel.ack(msg) no fim é o que avisa o broker que deu tudo certo e
a mensagem pode ser removida da fila. Sem ack, a mensagem volta -
esse é o assunto inteiro do próximo nó.
Vendo acontecer
Com o painel aberto em localhost:15672, roda o publisher algumas vezes
com o consumidor desligado. A aba Queues mostra o contador da
fila.pagamento subindo - as mensagens estão lá, guardadas, esperando.
Agora sobe o consumidor. O contador zera na hora.
Isso é o desacoplamento no tempo que a gente falou lá no primeiro nó, acontecendo na sua frente.
Três coisas pra fixar:
asserté idempotente - cria se não existe, ignora se existe. Chame sempre na subida.- Mensagem trafega como bytes. Serializar e desserializar é com você.
consumeregistra um callback, não é um loop. O processo precisa ficar vivo.
Dica: publisher e consumer são processos separados de propósito. Roda em dois terminais - é assim que vai ser em produção, com o worker num container diferente da API.
No próximo nó a gente fala do que acontece quando o consumidor falha no
meio do processamento - e por que ack é mais importante do que parece.
// Quiz
O que channel.assertQueue(...) faz se a fila já existir?