Para iniciar o desenvolvimento com o RabbitMQ utilizando o cliente Java nativo, é necessário configurar um projeto Maven e adicionar a dependência do driver AMQP. Recomenda-se o uso do JDK 8 ou superior para versões 5.x do cliente.
<dependency>
<groupId>com.rabbitmq</groupId>
<artifactId>amqp-client</artifactId>
<version>5.16.0</version>
</dependency>
O RabbitMQ trabalha com mecanismos de confirmação para garantir que as mensagens não sejam perdidas durante o trajeto entre o produtor, o servidor (broker) e o consumidor. 1. Confirmação Automática do Consumidor (Auto-Ack)
No modo de confirmação automática, o RabbitMQ considera que a mensagem foi entregue com sucesso assim que ela é enviada ao consumidor, independentemente de o processamento ter sido concluído ou não. ### Exemplo com Exchange Direct
Neste modelo, a mensagem é entregue apenas às filas cuja chave de roteamento (routing key) coincide exatamente com a chave definida pelo produtor. Produtor Direct:```
package com.tech.rabbit.direct;
import com.rabbitmq.client.BuiltinExchangeType; import com.rabbitmq.client.Channel; import com.rabbitmq.client.Connection; import com.rabbitmq.client.ConnectionFactory;
public class EmissorLogDirect { private static final String NOME_EXCHANGE = "logs_diretos";
public static void main(String[] args) throws Exception {
ConnectionFactory fabrica = new ConnectionFactory();
fabrica.setHost("localhost");
try (Connection conexao = fabrica.newConnection();
Channel canal = conexao.createChannel()) {
canal.exchangeDeclare(NOME_EXCHANGE, BuiltinExchangeType.DIRECT);
String[] niveis = {"info", "warning", "error"};
for (String nivel : niveis) {
String conteudo = "Log de sistema: " + nivel;
canal.basicPublish(NOME_EXCHANGE, nivel, null, conteudo.getBytes("UTF-8"));
System.out.println(" [x] Enviado: '" + nivel + "':'" + conteudo + "'");
}
}
}
}
**Consumidor Seletivo (Apenas Erros):**```
package com.tech.rabbit.direct;
import com.rabbitmq.client.*;
public class ReceptorApenasErro {
private static final String NOME_EXCHANGE = "logs_diretos";
public static void main(String[] args) throws Exception {
ConnectionFactory fabrica = new ConnectionFactory();
fabrica.setHost("localhost");
Connection conexao = fabrica.newConnection();
Channel canal = conexao.createChannel();
canal.exchangeDeclare(NOME_EXCHANGE, BuiltinExchangeType.DIRECT);
String nomeFila = canal.queueDeclare().getQueue();
// Faz o bind apenas para a chave 'error'
canal.queueBind(nomeFila, NOME_EXCHANGE, "error");
DeliverCallback callback = (consumerTag, entrega) -> {
String msg = new String(entrega.getBody(), "UTF-8");
System.out.println(" [!] Erro Recebido: " + msg);
};
// autoAck = true
canal.basicConsume(nomeFila, true, callback, consumerTag -> {});
}
}
- Confirmação Manual do Consumidor (Explicit Ack)
Para sistemas onde a integridade do processamento é crítica, utiliza-se a confirmação manual. O RabbitMQ só removerá a mensagem da fila após receber um sinal explícito (basicAck) do consumidor. Consumidor com Ack Manual e Rejeição:```
package com.tech.rabbit.confirm;
import com.rabbitmq.client.*; import java.io.IOException;
public class ReceptorComConfirmacao { private final static String FILA_TAREFAS = "fila_confirmacao_manual";
public static void main(String[] args) throws Exception {
ConnectionFactory factory = new ConnectionFactory();
factory.setHost("localhost");
Connection connection = factory.newConnection();
Channel channel = connection.createChannel();
channel.queueDeclare(FILA_TAREFAS, true, false, false, null);
System.out.println(" [*] Aguardando mensagens...");
// Define que o consumidor processará apenas 1 mensagem por vez
channel.basicQos(1);
Consumer consumer = new DefaultConsumer(channel) {
@Override
public void handleDelivery(String consumerTag, Envelope envelope,
AMQP.BasicProperties properties, byte[] body) throws IOException {
String mensagem = new String(body, "UTF-8");
try {
System.out.println(" [x] Processando: " + mensagem);
if (mensagem.contains("fail")) {
// Rejeita a mensagem. Se requeue = true, ela volta para a fila.
System.out.println(" [!] Falha detectada. Rejeitando...");
getChannel().basicReject(envelope.getDeliveryTag(), true);
} else {
// Confirma o processamento com sucesso
getChannel().basicAck(envelope.getDeliveryTag(), false);
System.out.println(" [v] Ack enviado.");
}
} catch (Exception e) {
getChannel().basicNack(envelope.getDeliveryTag(), false, true);
}
}
};
// autoAck = false é obrigatório para ack manual
channel.basicConsume(FILA_TAREFAS, false, consumer);
}
}
3. Confirmação do Publicador (Publisher Confirms)
-------------------------------------------------
O produtor também precisa ter certeza de que a mensagem chegou ao broker. Existem duas abordagens principais: síncrona e assíncrona. ### Confirmação Síncrona
Neste modo, o produtor envia a mensagem e aguarda a confirmação do servidor antes de prosseguir. ```
channel.confirmSelect(); // Habilita confirmações no canal
channel.basicPublish(exchange, routingKey, null, corpo);
if (channel.waitForConfirms(5000)) {
System.out.println("Mensagem confirmada pelo broker.");
} else {
System.err.println("Erro: Mensagem não confirmada.");
}
Confirmação Assíncrona
Mais eficiente para alto rendimento, utiliza callbacks para tratar confirmações (Acks) e falhas (Nacks) do broker de forma não bloqueante. ```
package com.tech.rabbit.publisher;
import com.rabbitmq.client.*;
public class EmissorAssincrono { private final static String NOME_EXCHANGE = "confirm_async_exchange";
public static void main(String[] args) throws Exception {
ConnectionFactory factory = new ConnectionFactory();
factory.setHost("localhost");
try (Connection connection = factory.newConnection();
Channel channel = connection.createChannel()) {
channel.exchangeDeclare(NOME_EXCHANGE, BuiltinExchangeType.DIRECT);
channel.confirmSelect();
// Listener para gerenciar as confirmações do broker
channel.addConfirmListener(new ConfirmListener() {
@Override
public void handleAck(long deliveryTag, boolean multiple) {
System.out.println("Broker recebeu mensagem: " + deliveryTag);
}
@Override
public void handleNack(long deliveryTag, boolean multiple) {
System.err.println("Broker falhou ao receber mensagem: " + deliveryTag);
}
});
// Listener para mensagens que não puderam ser roteadas (Mandatory)
channel.addReturnListener((replyCode, replyText, exchange, routingKey, properties, body) -> {
System.out.println("Mensagem retornada: " + new String(body));
});
for (int i = 0; i < 10; i++) {
String payload = "Dados sequenciais " + i;
// mandatory = true faz a mensagem retornar se não houver fila vinculada
channel.basicPublish(NOME_EXCHANGE, "chave_teste", true, null, payload.getBytes());
}
// Pequena pausa para observar os logs dos callbacks
Thread.sleep(2000);
}
}
}