Visão Geral do RabbitMQ
O RabbitMQ é um sistema de filas de mensagens open source desenvolvido em Erlang, que implementa o protocolo AMQP (Advanced Message Queuing Protocol). Ele atua como um intermediário (broker) que permite a comunicação assíncrona entre aplicações, eliminando a necessidade de conexões diretas e constantes entre os serviços.
As principais características incluem suporte a diversos modelos de roteamento, garantias de entrega de mensagens, persistência e segurança. A arquitetura baseia-se no conceito de produtores que enviam mensagens para exchanges, que por sua vez as distribuem para filas onde os consumidores as processam.
Componentes Fundamentais
- Producer (Produtor): Aplicação que envia mensagens para o broker.
- Consumer (Consumidor): Aplicação que recebe e processa mensagens das filas.
- Queue (Fila): Estrutura que armazena as mensagens até que sejam consumidas.
- Exchange (Troca): Recebe mensagens dos produtores e as repassa para as filas baseando-se em regras de roteamento.
- Binding (Vínculo): Regra que conecta uma Exchange a uma Queue, definindo como as mensagens são filtradas.
- Routing Key (Chave de Roteamento): Atributo enviado pelo produtor usado pela exchange para decidir o destino da mensagem.
- Channel (Canal): Conexão virtual dentro de uma conexão TCP usada para transmitir dados.
Modelos de Exchange
O comportamento do roteamento depende do tipo de Exchange configurada. Abaixo estão os principais padrões utilizados no Spring Boot.
1. Fanout (Broadcast)
Neste modelo, a exchange ignora a chave de roteamento e envia uma cópia da mensagem para todas as filas vinculadas a ela. É ideal para cenários de publish/subscribe onde múltiplos serviços precisam reagir ao mesmo evento.
Configuração da Infraestrutura:
@Configuration
public class FanoutBroadcastConfig {
@Bean
public Queue alertQueueA() {
return new Queue("alert_queue_email", true);
}
@Bean
public Queue alertQueueB() {
return new Queue("alert_queue_sms", true);
}
@Bean
public FanoutExchange broadcastExchange() {
return new FanoutExchange("alerts.fanout");
}
@Bean
public Binding bindEmailQueue() {
return BindingBuilder.bind(alertQueueA()).to(broadcastExchange());
}
@Bean
public Binding bindSmsQueue() {
return BindingBuilder.bind(alertQueueB()).to(broadcastExchange());
}
}
Publicação de Mensagens:
@RestController
@RequestMapping("/notifications")
public class NotificationController {
@Autowired
private RabbitTemplate rabbitTemplate;
@PostMapping("/broadcast")
public ResponseEntity<String> sendAlert() {
Map<String, Object> payload = new HashMap<>();
payload.put("level", "CRITICAL");
payload.put("timestamp", System.currentTimeMillis());
// O routing key é ignorado no modo fanout
rabbitTemplate.convertAndSend("alerts.fanout", "", payload);
return ResponseEntity.ok("Alerta disparado");
}
}
Consumo:
@RabbitListener(queues = "alert_queue_email")
public void handleEmailAlert(Map<String, Object> data) {
// Lógica de envio de e-mail
System.out.println("Processando alerta de e-mail: " + data.get("level"));
}
2. Direct (Roteamento Exato)
A exchange do tipo Direct encaminha a mensagem apenas para as filas cuja Binding Key corresponde exatamente à Routing Key da mensagem. É o padrão padrão quando não se especifica o tipo de exchange.
Configuração:
@Configuration
public class RoutingKeyConfig {
@Bean
public Queue orderQueue() {
return new Queue("orders.processing", true);
}
@Bean
public DirectExchange directExchange() {
return new DirectExchange("orders.direct");
}
@Bean
public Binding orderBinding() {
return BindingBuilder.bind(orderQueue()).to(directExchange()).with("order.created");
}
}
Publicação e Consumo:
// Envio
rabbitTemplate.convertAndSend("orders.direct", "order.created", orderData);
// Recepção
@RabbitListener(queues = "orders.processing")
public void processOrder(Map<String, Object> order) {
// Processamento do pedido
}
3. Topic (Padrão Wildcard)
Permite roteamento baseado em padrões utilizando wildcards. A Routing Key é uma string separada por pontos. O símbolo * corresponde a uma palavra e # corresponde a zero ou mais palavras.
Confgiuração:
@Configuration
public class WildcardPatternConfig {
@Bean
public Queue logQueue() {
return new Queue("logs.all");
}
@Bean
public TopicExchange topicExchange() {
return new TopicExchange("logs.topic");
}
@Bean
public Binding logBinding() {
// Recebe mensagens com routing key começando com "logs."
return BindingBuilder.bind(logQueue()).to(topicExchange()).with("logs.#");
}
}
Exemplo de envio: rabbitTemplate.convertAndSend("logs.topic", "logs.system.error", data); será recebido pela fila logs.all.
4. Headers (Atributos)
O roteamento é baseado nos cabeçalhos (headers) da mensagem, ignorando a Routing Key. Pode-se configurar para匹配 todos os headers (x-match: all) ou qualquer um (x-match: any).
@Configuration
public class MetadataExchangeConfig {
@Bean
public Queue priorityQueue() {
return new Queue("priority.tasks");
}
@Bean
public HeadersExchange headersExchange() {
return new HeadersExchange("tasks.headers");
}
@Bean
public Binding priorityBinding() {
Map<String, Object> headers = new HashMap<>();
headers.put("priority", "high");
return BindingBuilder.bind(priorityQueue()).to(headersExchange()).whereAll(headers).match();
}
}
Padrão RPC (Remote Procedure Call)
O RabbitMQ pode ser utilizado para implementar chamadas síncronas estilo RPC. O cliente envia uma mensagem com um ID de correlação e uma fila de resposta temporária. O servidor processa e responde na fila indicada.
Configuração do Template RPC:
@Configuration
public class RpcTemplateConfig {
@Autowired
private ConnectionFactory connectionFactory;
@Bean
public RabbitTemplate rpcRabbitTemplate() {
RabbitTemplate template = new RabbitTemplate(connectionFactory);
template.setUseTemporaryReplyQueues(false);
template.setReplyAddress("amq.rabbitmq.reply-to");
template.setUserCorrelationId(true);
template.setReplyTimeout(5000);
return template;
}
}
Execução da Chamada:
@GetMapping("/compute")
public String executeRpc() {
CorrelationData correlation = new CorrelationData(UUID.randomUUID().toString());
MessageProperties props = new MessageProperties();
props.setCorrelationId(correlation.getId());
Message request = new Message("10".getBytes(), props);
Message response = rpcRabbitTemplate().sendAndReceive("rpc.requests", "rpc.key", request, correlation);
return response != null ? new String(response.getBody()) : "Timeout";
}
Otimização e Confiabilidade
Para garantir a entrega das mensagens em ambientes de produção, é essencial configurar confirmações de publicação e consumo manual.
Confirm Callback (Publicador)
Permite que o produtor saiba se o broker recebeu a mensagem.
@Configuration
public class ReliabilityConfig {
@Bean
public RabbitTemplate reliableTemplate(ConnectionFactory factory) {
RabbitTemplate template = new RabbitTemplate(factory);
template.setMandatory(true);
template.setConfirmCallback((correlation, ack, cause) -> {
if (ack) {
System.out.println("Mensaje confirmada pelo broker");
} else {
System.out.println("Falha na entrega: " + cause);
}
});
template.setReturnsCallback(returned -> {
System.out.println("Mensagem retornada: " + returned.getMessage());
});
return template;
}
}
Acknowledgement Manual e QoS (Consumidor)
Configruar o modo de ack manual no application.yml evita a perda de mensagens se o consumidor falhar durante o processamento.
spring:
rabbitmq:
listener:
simple:
acknowledge-mode: manual
prefetch: 1
No código do consumidor, a confirmação deve ser explícita:
@RabbitListener(queues = "orders.processing")
public void handleOrder(Message message, Channel channel) throws IOException {
try {
// Processamento da negócio
channel.basicAck(message.getMessageProperties().getDeliveryTag(), false);
} catch (Exception e) {
// Rejeita e coloca de volta na fila
channel.basicNack(message.getMessageProperties().getDeliveryTag(), false, true);
}
}