0x09【RabbitMQ】Return 消息机制
1、简介
-
Return Listener用于处理一些不可路由的消息; -
我们的消息生产者,通过指定一个Exchange和RoutingKey,把消息送达到某一个队列中去,然后我们的消费者监听队列,进行消费处理。
-
但是在某些情况下,如果我们在发送消息的时候,当前的exchange不存在或者指定的路由key路由不到,这个时候如果我们需要监听这种不可达的消息,就要使用
Return Listener; -
在基础QPI中有一个关键的配置项:
channel.basicPublish(x,x,Mandatory,x,x)Mandatory:- 如果为true,则监听器会接收到路由不可达的消息,然后进行后续处理;
- 如果为false,那么broker端自动删除该消息( 不处理你了,直接删掉。 )。
2、Return消息机制流程

3、代码编写
生产者
public class ReturnProd {
public static void main(String[] args) throws Exception{
ConnectionFactory connectionFactory = new ConnectionFactory();
connectionFactory.setHost("49.235.24.110");
connectionFactory.setPort(5672);
connectionFactory.setVirtualHost("/");
connectionFactory.setAutomaticRecoveryEnabled(true);
connectionFactory.setNetworkRecoveryInterval(3000);
Connection connection = connectionFactory.newConnection();
Channel channel = connection.createChannel();
String exchangeName = "test_return_exchange";
String routingKey = "return.save";
String routingKeyError = "return.err.save";
String msg = "hello RabbitMQ return test";
// 添加Return监听
channel.addReturnListener(new ReturnListener() {
@Override
public void handleReturn(int replyCode, String replyText, String exchange, String routingKey,
AMQP.BasicProperties properties, byte[] body) throws IOException {
System.err.println("-----------handle return--------------------");
System.err.println("replacyCode:"+replyCode);
System.err.println("replyText:"+replyText);
System.err.println("exchange:"+exchange);
System.err.println("routingKey"+routingKey);
System.err.println("properties:"+properties);
System.err.println("body:"+new String(body));
}
});
// 注意关键在于mandatory参数设置为true
// channel.basicPublish(exchangeName,routingKeyError,true,null,msg.getBytes());
channel.basicPublish(exchangeName,routingKey,true,null,msg.getBytes());
}
}
消费者
public class RetrunCon {
public static void main(String[] args) throws Exception{
ConnectionFactory connectionFactory = new ConnectionFactory();
connectionFactory.setHost("49.235.24.110");
connectionFactory.setPort(5672);
connectionFactory.setVirtualHost("/");
connectionFactory.setAutomaticRecoveryEnabled(true);
connectionFactory.setNetworkRecoveryInterval(3000);
Connection connection = connectionFactory.newConnection();
Channel channel = connection.createChannel();
String exchangeName = "test_return_exchange";
String exchangeType = "topic";
String routingKey = "return.*";
String queueName = "test_return_queue";
channel.exchangeDeclare(exchangeName,exchangeType,true,false,null);
channel.queueDeclare(queueName,true,false,false,null);
channel.queueBind(queueName,exchangeName,routingKey);
QueueingConsumer consumer = new QueueingConsumer(channel);
channel.basicConsume(queueName,true,consumer);
while(true){
QueueingConsumer.Delivery delivery = consumer.nextDelivery();
String msg = new String(delivery.getBody());
System.out.println("收到消息:"+msg);
}
}
}
==测试时注意:发送错误消息,并且Mandatory设置为true,才会有反应==
本文是原创文章,采用 CC BY-NC-ND 4.0 协议,完整转载请注明来自 程序员小航
评论
匿名评论
隐私政策
你无需删除空行,直接评论以获取最佳展示效果