• Springboot 集成 RabbitMq 实现消息确认机制


    消息确认主要分为两种

    1. 发送确认,发送确认包含两种情况,一种是消息是否到达交换机,一种是消息是否到达队列
    2. 接收确认

    一、发送方消息确认

    1、ConfirmCallback 接口

    配置文件需要开启配置

    publisher-confirm-type: correlated

    消息发送到交换机回调,当消息发送到交换机时,会触发此接口中的 confirm 回调函数
    confirm 中有三个参数

    1. CorrelationData:包含消息发送时的 id,可以根据 id 快速定位
    2. ack:消息发送成功与否,成功为 true,失败为 false
    3. cause:消息发送失败的原因

    2、ReturnCallback 接口

    配置文件需要开启配置

    publisher-returns: true

    消息发送到队列回调,接口中包含 returnedMessage 函数,如果消息发送成功,则不会调用此回调函数,如果消息发送失败,则会调用此回调函数
    returnedMessage 中有5个参数

    1. message:发送的消息内容
    2. replyCode:回应码
    3. replyText:回应信息
    4. exchange:交换机
    5. routingKey:路由键

    3、配置发送方消息确认

    3.1、application 配置

    server:
      port: 7008
    spring:
      application:
        name: rabbitmq-demo
      #配置rabbitMq 服务器
      rabbitmq:
        host: 192.168.202.128
        port: 5672
        username: guest
        password: guest
        virtual-host: /
        # 开启生产者发布确认,确认消息已发送到交换机 Exchange
        publisher-confirm-type: correlated
        # 开启发布者返回,确认消息已发送到队列 Queue
        publisher-returns: true
        listener:
          simple:
            # 开启消息手动确认(即需要调用channel.basicAck才会从队列中删除消息),默认是 NONE
            acknowledge-mode: manual
            #表示消费者端每次从队列拉取多少个消息进行消费,直到手动确认消费完毕后,才会继续拉取下一条
            prefetch: 1
            #消费被拒绝时 true:重回队列 false为否
            default-requeue-rejected: false
            retry:
              #是否支持重试,默认false
              enabled: true
              #重试最大次数,默认3次
              max-attempts: 3
              #重试最大间隔时间
              max-interval: 1000ms
    
    
    • 1
    • 2
    • 3
    • 4
    • 5
    • 6
    • 7
    • 8
    • 9
    • 10
    • 11
    • 12
    • 13
    • 14
    • 15
    • 16
    • 17
    • 18
    • 19
    • 20
    • 21
    • 22
    • 23
    • 24
    • 25
    • 26
    • 27
    • 28
    • 29
    • 30
    • 31
    • 32

    3.2、配置发送回调函数

    package com.wxw.notes.rabbitmq.demo.config;
    
    import lombok.extern.slf4j.Slf4j;
    import org.springframework.amqp.core.Message;
    import org.springframework.amqp.rabbit.connection.ConnectionFactory;
    import org.springframework.amqp.rabbit.connection.CorrelationData;
    import org.springframework.amqp.rabbit.core.RabbitTemplate;
    import org.springframework.context.annotation.Bean;
    import org.springframework.context.annotation.Configuration;
    
    /**
     * @author wuxiongwei
     * @date 2022/8/3 10:25
     * @Description
     */
    @Slf4j
    @Configuration
    public class RabbitTemplateConfig {
    
        @Bean
        public RabbitTemplate createRabbitTemplate(ConnectionFactory connectionFactory) {
            RabbitTemplate rabbitTemplate = new RabbitTemplate();
            rabbitTemplate.setConnectionFactory(connectionFactory);
            // 设置开启Mandatory,才能触发回调函数,无论消息推送结果怎么样都强制调用回调函数
            rabbitTemplate.setMandatory(true);
    
            // 确认消息送到交换机(Exchange)回调
            // 如果消息没有到 exchange,则 confirm 回调,ack=false; 如果消息到达exchange,则confirm回调,ack=true
            rabbitTemplate.setConfirmCallback(new RabbitTemplate.ConfirmCallback() {
                @Override
                public void confirm(CorrelationData correlationData, boolean ack, String cause) {
                    if(ack){
                        log.info("消息发送交换机成功:correlationData:{},ack:{},cause:{}",correlationData,ack,cause);
                    }else{
                        log.info("消息发送交换机失败:correlationData:{},ack:{},cause:{}",correlationData,ack,cause);
                    }
                }
            });
    
            // 确认消息送到队列(Queue)回调
            // 如果exchange到queue成功,则不回调return;如果exchange到queue失败,则回调return(需设置mandatory=true,否则不回回调,消息就丢了)
            rabbitTemplate.setReturnCallback(new RabbitTemplate.ReturnCallback() {
                @Override
                public void returnedMessage(Message message, int replyCode, String replyText, String exchange, String routingKey) {
                    log.info("确认消息送到队列结果:");
                    log.info("发送消息:{}", message);
                    log.info("回应码:{}", replyCode);
                    log.info("回应信息:{}", replyText);
                    log.info("交换机:{}", exchange);
                    log.info("路由键:{}", routingKey);
                }
            });
            return rabbitTemplate;
        }
    
    }
    
    
    • 1
    • 2
    • 3
    • 4
    • 5
    • 6
    • 7
    • 8
    • 9
    • 10
    • 11
    • 12
    • 13
    • 14
    • 15
    • 16
    • 17
    • 18
    • 19
    • 20
    • 21
    • 22
    • 23
    • 24
    • 25
    • 26
    • 27
    • 28
    • 29
    • 30
    • 31
    • 32
    • 33
    • 34
    • 35
    • 36
    • 37
    • 38
    • 39
    • 40
    • 41
    • 42
    • 43
    • 44
    • 45
    • 46
    • 47
    • 48
    • 49
    • 50
    • 51
    • 52
    • 53
    • 54
    • 55
    • 56
    • 57

    4、配置交换机和队列

        /**
         * 路由模式队列1
         */
        @Bean
        public Queue directQueue1() {
            return new Queue(DIRECT_QUEUE1);
        }
    
        /**
         * 路由模式交换机
         */
        @Bean
        public DirectExchange directExchange() {
            return new DirectExchange(DIRECT_EXCHANGE);
        }
    
        /**
         * 路由模式队列1和交换机绑定
         */
        @Bean
        public Binding bindingDirectExchange1() {
            return BindingBuilder.bind(directQueue1()).to(directExchange()).with(DIRECT_ROUTING_KEY1);
        }
    
    • 1
    • 2
    • 3
    • 4
    • 5
    • 6
    • 7
    • 8
    • 9
    • 10
    • 11
    • 12
    • 13
    • 14
    • 15
    • 16
    • 17
    • 18
    • 19
    • 20
    • 21
    • 22
    • 23

    5、消息生产者

    @RestController
    @Slf4j
    @RequestMapping("/confirm")
    @RequiredArgsConstructor
    public class ConfirmProviderController {
    
        private final RabbitTemplate rabbitTemplate;
    
    
        @RequestMapping("/routingSend")
        public void routingSend() {
            String routingKeyMessage1 = "routing Message 1";
            log.info("路由模式生产者发送消息 :{}", routingKeyMessage1);
            
            rabbitTemplate.convertAndSend(RabbitMqConfig.DIRECT_EXCHANGE, RabbitMqConfig.DIRECT_ROUTING_KEY1, routingKeyMessage1, new CorrelationData(UUID.randomUUID().toString()));
        }
        
    }
    
    • 1
    • 2
    • 3
    • 4
    • 5
    • 6
    • 7
    • 8
    • 9
    • 10
    • 11
    • 12
    • 13
    • 14
    • 15
    • 16
    • 17
    • 18

    6、发送确认测试

    6.1、发送正常

    这个时候,我们正常发送,肯定是没有问题,所以,只会触发 ConfirmCallback 的交换机发送成功回调
    在这里插入图片描述

    6.2、发送交换机失败

    测试一下,更改一下交换机名称,让它找不到交换机
    在这里插入图片描述
    可以看到,触发了 ConfirmCallback 的交换机发送失败回调
    在这里插入图片描述

    6.3、发送队列失败

    测试一下,更改一下 RoutingKey ,让它找不到队列
    在这里插入图片描述
    再次发送,就可以看到无法找到我们的队列了,触发了 ReturnCallback 回调
    在这里插入图片描述

    二、消费者消息确认

    1、消费者正常消费消息

    package com.wxw.notes.rabbitmq.demo.consumer;
    
    import com.rabbitmq.client.Channel;
    import com.wxw.notes.rabbitmq.demo.config.RabbitMqConfig;
    import lombok.extern.slf4j.Slf4j;
    import org.springframework.amqp.core.Message;
    import org.springframework.amqp.rabbit.annotation.RabbitListener;
    import org.springframework.stereotype.Component;
    
    import java.io.IOException;
    
    /**
     * @author wuxiongwei
     * @date 2022/7/28 10:32
     * @Description
     */
    @Slf4j
    @Component
    public class ComfirmConsumer {
        
        /**
         * 路由模式
         *
         * @param message 消息内容
         */
        @RabbitListener(queues = RabbitMqConfig.DIRECT_QUEUE1)
        public void fanoutConsumer1(Message message, Channel channel) throws IOException {
            log.info("发布订阅模式消费者1接收到消息:{}", message);
            long deliveryTag = message.getMessageProperties().getDeliveryTag();
            // 手动进行消息确认
            channel.basicAck(deliveryTag,false);
        }
    }
    
    
    • 1
    • 2
    • 3
    • 4
    • 5
    • 6
    • 7
    • 8
    • 9
    • 10
    • 11
    • 12
    • 13
    • 14
    • 15
    • 16
    • 17
    • 18
    • 19
    • 20
    • 21
    • 22
    • 23
    • 24
    • 25
    • 26
    • 27
    • 28
    • 29
    • 30
    • 31
    • 32
    • 33
    • 34

    2、消费者确认模式和相关参数

    消费者消息确认,我们需要在配置文件中开启消息的手动确认,因为默认是自动确认的

    acknowledge-mode: manual

    确认模式有三种

    1. none:默认,不进行确认
    2. auto:自动确认
    3. manual:手动确认

    消息确认函数中会用到的几个参数

    1. deliveryTag :唯一标识 ID,当一个消费者向 RabbitMQ 注册后,会建立起一个 Channel ,RabbitMQ 会用 basic.deliver 方法向消费者推送消息,这个方法携带了一个 delivery tag, 它代表了 RabbitMQ 向该 Channel 投递的这条消息的唯一标识 ID,是一个单调递增的正整数,delivery tag 的范围仅限于 Channel。
    2. multiple:是否批处理,当该参数为 true 时,则可以一次性确认 delivery_tag 小于等于传入值的所有消息
    3. requeue:被拒绝的是否重新入队列

    3、消息确认 basicAck

    basicAck 包含 deliveryTag 和 multiple 参数

    消息确认的过程

    生产者发送一条消息,当消息开始被消费到被确认时,个人理解分为三种状态
    第一种就是等待消费状态,消费者还未接收到消息,消息会处于 Ready 状态
    第二种就是消费中的状态,消费者开始消费消息,但还未被确认,消息处于 Unacked 状态
    第三种就是消费完成,消费者完成消费确认

    比如下面三张图,第一二张图,是消费者接收到消息,中间打了一个断点,然后消息就是 Unacked 状态,如果这个时候程序异常,或者手动将程序终止,这条消息会再次回到等待消费的状态,不会丢失,也就是说,只要没有进行消费确认,消息就会一直在。

    在这里插入图片描述在这里插入图片描述

    此时,将断点走过,消息被消费确认,就会被移出队列

    在这里插入图片描述

    4、basicNack

    basicNack 方法用于否定当前消息。 可以批量拒绝消息

    channel.basicNack(deliveryTag,false,false) 三个参数
    deliveryTag:唯一参数
    multiple:和相关参数解释一样
    requeue:和相关参数解释一样

    第三个参数最好设定为 false,否则消息重回队列,当前消费者还会一直消费,程序遭不住
    在这里插入图片描述

    5、basicReject

    basicReject 方法用于拒绝当前的消息,相比 basicNack 少了一个批量拒绝的参数,所以,一般用来拒绝当前消息,如果需要批量拒绝,可使用 basicNack

    channel.basicReject(deliveryTag,false)两个参数
    deliveryTag:唯一参数
    requeue:和相关参数解释一样

    第二个参数一般也是设定 false,设定为 true 后果和 basicNack 一样,遭不住,有兴趣可以试试

    其实消息如果拒绝了还想消费,可以通过其他方式来重新消费,比如先加到另外的队列

    三、总结

    1. 消息在未发出时,我们可以通过发送方的回调来确认消息是否推送到交换机、推送到队列,这样就能比较准确的定位问题
    2. 消息只要在消费过程中,没有做消费确认的状况,即使程序出现异常,消息也会重新回到队列,这样消息就不会丢失

    源码

    注意,看代码的话,请仔细一点

    https://github.com/wxwhowever/springboot-notes/tree/master/rabbitmq-demo

  • 相关阅读:
    -bash: ./xxx.sh: /bin/sh^M: bad interpreter: No such file or directory
    (10)(10.8) 固件下载
    解决在uniapp中自定义组件onLoad回调不执行
    JavaSE入门---程序逻辑控制
    【千律】OpenCV基础:Hough圆检测
    企业应用超融合架构的设计实例及超融合应用场景分析
    Python 使用executemany批量向mysql插入数据
    Go常用命令
    Netty入门知识点
    小学期,第三场-下午:WEB_sessionlfi
  • 原文地址:https://blog.csdn.net/wxw1997a/article/details/126170811