Implementação de Mensageria com RabbitMQ no Spring Boot

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);
    }
}

Tags: RabbitMQ spring-boot AMQP Mensageria java

Publicado em 9-17 08:01