1、简介

  • Return Listener用于处理一些不可路由的消息;

  • 我们的消息生产者,通过指定一个Exchange和RoutingKey,把消息送达到某一个队列中去,然后我们的消费者监听队列,进行消费处理。

  • 但是在某些情况下,如果我们在发送消息的时候,当前的exchange不存在或者指定的路由key路由不到,这个时候如果我们需要监听这种不可达的消息,就要使用Return Listener;

  • 在基础QPI中有一个关键的配置项: channel.basicPublish(x,x,Mandatory,x,x)

    • Mandatory
      • 如果为true,则监听器会接收到路由不可达的消息,然后进行后续处理;
      • 如果为false,那么broker端自动删除该消息( 不处理你了,直接删掉。 )。

2、Return消息机制流程

8Hb5QO.png

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,才会有反应==