RabbitMQ消息确认机制剖析

目录
  • 前言
  • 消息确认
    • 基本流程
    • 消息确认模式
      • ConfirmCallback确认模式
      • ReturnCallback退回模式
      • 消息发送者确认
      • 消息接收者确认
      • basicAck模式
      • basicNack模式
      • basicReject模式
    • 测试
      • 解决办法
  • 消费者确认失败
  • 总结

前言

上一章讲解了RabbitMq的三种Exchange消息发送的模式,但是在默认情况下RabbitMQ并不能保证消息是否发送成功,以及是否能被成功消费,为了保证消息在传递过程中不丢失,需要对消息进行确认机制,来提高消息的可靠性。

消息确认

基本流程

说明:

  • 生产者发送消息到RabbitMQ Server后,RabbitMQ Server需要对生产者进行消息Confirm确认。
  • 消费者消费消息后需要对 RabbitMQ Server进行消息ACK确认。

消息确认模式

RabbitMq提供了两种消息发送者确认模式分别为: ConfirmCallback确认模式和 ReturnCallback退回模式。

ConfirmCallback确认模式

@Component
public class RabbitConfirmConfig implements ConfirmCallback
{
    private Logger logger = LoggerFactory.getLogger(RabbitConfirmConfig.class);
    public void confirm(CorrelationData correlationData, boolean ack,
            String cause)
    {
        logger.info("数据内容:{}",correlationData);
        logger.info("是否确认成功:{}",ack);
        logger.info("错误原因:{}",cause);
        if (!ack)
        {
            logger.info("exchange produce confirm message send error" + cause);
        }
        else
        {
            logger.info("exchange produce confirm message send success");
        }
    }
}

说明:ConfirmCallback模式确认,需要重写confirm接方法,此方法的三个参数分别为:CorrelationData、ack、cause

  • CorrelationData:对象内部只有一个id属性,用来表示当前消息的唯一性。
  • ack:消息投递状态,true表示投递成功
  • cause: 消息投递失败原因

虽然消息被broker接收到只能表示已经到达MQ服务器,但是并不能保证消息一定会被投递到目标 queue里。所以我们需要实现returnCallback来进行相关处理。

ReturnCallback退回模式

@Component
public class RabbitReturnConfig implements ReturnCallback
{
    private Logger logger = LoggerFactory.getLogger(RabbitReturnConfig.class);
    public void returnedMessage(Message message, int replyCode,
            String replyText, String exchange, String routingKey)
    {
       logger.info("消息发送送到队列信息:");
       logger.info("发生消息:{}",message);
       logger.info("回应码:{}",replyCode);
       logger.info("回应信息:{}",replyText);
       logger.info("交换机:{}",exchange);
       logger.info("路由键:{}",routingKey);
    }
}

说明:实现接口ReturnCallback重写returnedMessage()方法,方法有五个参数message(消息体)、replyCode(响应code)、replyText(响应内容)、exchange(交换机)、routingKey(路由键)。

消息发送者确认

@Component
public class MqConfirmProduce
{
    @Autowired
    private RabbitTemplate rabbitTemplate;
    @Autowired
    private RabbitConfirmConfig rabbitConfirmConfig;
    @Autowired
    private RabbitReturnConfig rabbitReturnConfig;
    /**
     *
     * @param exchange 消息交互机名称
     * @param routeKey 消息路由键的名称
     * @param message  消息内容
     */
    public void sendMessage(String exchange ,String routeKey,Object msg)
    {
        //确保消息发送失败后可以重新返回到队列中
        rabbitTemplate.setMandatory(true);
        // 消费者确认收到消息后,手动ack回执回调处理
        rabbitTemplate.setConfirmCallback(rabbitConfirmConfig);
        //消息投递到队列失败回调处理
        rabbitTemplate.setReturnCallback(rabbitReturnConfig);
        //保证消息唯一性
        CorrelationData correlationData =new CorrelationData(UUID.randomUUID().toString());
        //发送消息
        rabbitTemplate.convertAndSend(exchange,routeKey,msg,
                message -> {
                    message.getMessageProperties().setDeliveryMode(MessageDeliveryMode.PERSISTENT);
                    return message;
                },
                correlationData);
    }
}

说明:注意需要开启消息确认的配置:

  rabbitmq:
    host: localhost
    port: 5672
    username: guest
    password: guest
    virtual-host: /
    #开启发送确认
    publisher-confirms: true
    # 开启发送失败退回
    publisher-returns: true
    listener:
      simple:
       # 手动确认
        acknowledge-mode: manual
        retry:
          enabled: true

消息接收者确认

@Component
@RabbitListener(queues = "testQueue")
public class MqConfirmConsumer
{
    private static final Logger logger = LoggerFactory.getLogger(MqConfirmConsumer.class);
    @RabbitHandler
    public void receive(String msg, Channel channel, Message message) throws IOException
    {
        logger.info("receive message content:{}",message);
        try
        {
            logger.info("开始消息确认");
            channel.basicAck(message.getMessageProperties().getDeliveryTag(), false);
            logger.info("消息确认成功");
        }
        catch (Exception e)
        {
             logger.error("消息确认失败,即将再次返回队列中");               channel.basicNack(message.getMessageProperties().getDeliveryTag(), true, true);
        }
    }
}

说明:消息者确认消息有三种模式,分别为basicAck、basicNack、basicReject。

basicAck模式

表示成功确认,使用此回执方法后,消息会被rabbitmq broker删除。

void basicAck(long deliveryTag, boolean multiple)
  • deliveryTag:消息投递序号,
  • multiple:是否批量确认,值为 true则会一次性ack所有小于当前消息deliveryTag的消息。

basicNack模式

表示失败确认,一般在消费消息异常时用到此方法,可以将消息重新投递入队列。

void basicNack(long deliveryTag, boolean multiple, boolean requeue)
  • deliveryTag:表示消息投递序号。
  • requeue: 表示消息是否重新入队列,true表示重新投入队列中。
  • multiple:是否批量确认,true表示会一次性ack所有小于当前消息deliveryTag的消息。

basicReject模式

basicReject:拒绝消息,与basicNack区别在于不能进行批量操作,其他用法很相似。

void basicReject(long deliveryTag, boolean requeue)
  • deliveryTag:消息投递序号。
  • requeue:值为true表示消息重新入队列

测试

测试发送消息,消息发送者的确认信息如下:

c.s.f.r.config.RabbitConfirmConfig - exchange produce confirm message send success
c.s.f.r.config.RabbitConfirmConfig - 数据内容:CorrelationData [id=88ea47a5-726d-44c5-9839-1f2a6bf942ed]
c.s.f.r.config.RabbitConfirmConfig - 是否确认成功:true
c.s.f.r.config.RabbitConfirmConfig - 错误原因:null
c.s.f.r.config.RabbitConfirmConfig - exchange produce confirm message send success

消费者的确认信息如下:

receive message content:(Body:'this is test message' MessageProperties [headers={spring_listener_return_correlation=0fcefb6d-acea-4eb2-8484-e3a82f8c584f, spring_returned_message_correlation=88ea47a5-726d-44c5-9839-1f2a6bf942ed}, contentType=text/plain, contentEncoding=UTF-8, contentLength=0, receivedDeliveryMode=PERSISTENT, priority=0, redelivered=false, receivedExchange=testDirect, receivedRoutingKey=testDirectRouting, deliveryTag=2, consumerTag=amq.ctag-dOwkSPuI1e0HR_1Ufu3Erw, consumerQueue=testQueue])
 c.s.f.r.consumer.MqConfirmConsumer - 开始消息确认
c.s.f.r.consumer.MqConfirmConsumer - 消息确认成功

消费者确认失败

如果消息确认在消费者确认失败,那么消息将会重写投递导导消息队列的首部。模拟消费者确认失败场景:

@Component
@RabbitListener(queues = "testQueue")
public class MqConfirmConsumer
{
    private static final Logger logger = LoggerFactory.getLogger(MqConfirmConsumer.class);
    @RabbitHandler
    public void receive(String msg, Channel channel, Message message) throws IOException
    {
        logger.info("receive message content:{}",message);
        try
        {
            logger.info("开始消息确认");
            int c=1/0;
            channel.basicAck(message.getMessageProperties().getDeliveryTag(), false);
            logger.info("消息确认成功");
        }
        catch (Exception e)
        {
          logger.error("消息确认失败,即将再次返回队列中");                   channel.basicNack(message.getMessageProperties().getDeliveryTag(), true, true);
        }
    }
}

查看执行结果:

c.s.f.r.consumer.MqConfirmConsumer - receive message content:(Body:'this is test message' MessageProperties [headers={spring_listener_return_correlation=0fcefb6d-acea-4eb2-8484-e3a82f8c584f, spring_returned_message_correlation=39d4cdd1-cbeb-4090-91ea-9e5d0bed785c}, contentType=text/plain, contentEncoding=UTF-8, contentLength=0, receivedDeliveryMode=PERSISTENT, priority=0, redelivered=false, receivedExchange=testDirect, receivedRoutingKey=testDirectRouting, deliveryTag=1, consumerTag=amq.ctag-e5GtG455pkm7eWfY3xGleg, consumerQueue=testQueue])
c.s.f.r.consumer.MqConfirmConsumer - 开始消息确认
c.s.f.r.consumer.MqConfirmConsumer - 消息确认失败,即将再次返回队列中

消息已经重新返回队列中。我们查看队列信息具体如下:

说明:我们可以看到消息为Unacked状态,消息又会重新会被消费,然后确认失败,又重新被消费,导致死循环。

解决办法

针对这种情况,我们将如何处理呢?我们手动确认失败后,并将消息持久入到MySQL中通过定时任务做补偿。然后删除消息队列。具体修改如下:

 @RabbitHandler
    public void receive(String msg, Channel channel, Message message) throws IOException
    {
        logger.info("receive message content:{}",message);
        try
        {
            logger.info("开始消息确认");
            int c=1/0;
            channel.basicAck(message.getMessageProperties().getDeliveryTag(), false);
            logger.info("消息确认成功");
        }
        catch (Exception e)
        {
            if (message.getMessageProperties().getRedelivered())
            {
                logger.error("消息确认失败,拒绝处理");
              //执行持久化处理                channel.basicReject(message.getMessageProperties().getDeliveryTag(), false);
          }
            else
            {
                logger.error("消息确认失败,即将再次返回队列中");
                channel.basicNack(message.getMessageProperties().getDeliveryTag(), false, true);
            }
        }
    }

修改后执行结果如下:

总结

本文讲解了RabbitMQ消息确认机制,消息是否需要确认,我们需要根据业务的场景来分析,如有疑问,请随时反馈,更多关于RabbitMQ消息确认的资料请关注我们其它相关文章!

(0)

相关推荐

  • 详解SpringBoot整合RabbitMQ如何实现消息确认

    目录 简介 生产者消息确认 介绍 流程 配置 ConfirmCallback ReturnCallback 注册ConfirmCallback和ReturnCallback 消费者消息确认 介绍 手动确认三种方式 简介 本文介绍SpringBoot整合RabbitMQ如何进行消息的确认. 生产者消息确认 介绍 发送消息确认:用来确认消息从 producer发送到 broker 然后broker 的 exchange 到 queue过程中,消息是否成功投递. 如果消息和队列是可持久化的,那么确认消

  • springboot + rabbitmq 如何实现消息确认机制(踩坑经验)

    本文收录在个人博客:www.chengxy-nds.top,技术资源共享,一起进步 最近部门号召大伙多组织一些技术分享会,说是要活跃公司的技术氛围,但早就看穿一切的我知道,这 T M 就是为了刷KPI.不过,话说回来这的确是件好事,与其开那些没味的扯皮会,多做技术交流还是很有助于个人成长的. 于是乎我主动报名参加了分享,咳咳咳~ ,真的不是为了那点KPI,就是想和大伙一起学习学习! 这次我分享的是 springboot + rabbitmq 如何实现消息确认机制,以及在实际开发中的一点踩坑经验,

  • SpringBoot整合RabbitMQ实现消息确认机制

    前面几篇案例已经将常用的交换器(DirectExchange.TopicExchange.FanoutExchange)的用法介绍完了,现在我们来看一下消息的回调,也就是消息确认. 在rabbitmq-provider项目的application.yml文件上加上一些配置 server: port: 8021 spring: #给项目来个名字 application: name: rabbitmq-provider #配置rabbitMq 服务器 rabbitmq: host: 127.0.0.

  • RabbitMQ消息确认机制剖析

    目录 前言 消息确认 基本流程 消息确认模式 ConfirmCallback确认模式 ReturnCallback退回模式 消息发送者确认 消息接收者确认 basicAck模式 basicNack模式 basicReject模式 测试 解决办法 消费者确认失败 总结 前言 上一章讲解了RabbitMq的三种Exchange消息发送的模式,但是在默认情况下RabbitMQ并不能保证消息是否发送成功,以及是否能被成功消费,为了保证消息在传递过程中不丢失,需要对消息进行确认机制,来提高消息的可靠性.

  • 利用Python学习RabbitMQ消息队列

    RabbitMQ可以当做一个消息代理,它的核心原理非常简单:即接收和发送消息,可以把它想象成一个邮局:我们把信件放入邮箱,邮递员就会把信件投递到你的收件人处,RabbitMQ就是一个邮箱.邮局.投递员功能综合体,整个过程就是:邮箱接收信件,邮局转发信件,投递员投递信件到达收件人处. RabbitMQ和邮局的主要区别就是RabbitMQ接收.存储和发送的是二进制数据----消息. rabbitmq基本管理命令: 一步启动Erlang node和Rabbit应用:sudo rabbitmq-serv

  • Java RabbitMQ消息队列详解常见问题

    目录 消息堆积 保证消息不丢失 死信队列 延迟队列 RabbitMQ消息幂等问题 RabbitMQ消息自动重试机制 合理的选择重试机制 消费者开启手动ack模式 rabbitMQ如何解决消息幂等问题 RabbitMQ解决分布式事务问题 基于RabbitMQ解决分布式事务的思路 消息堆积 消息堆积的产生场景: 生产者产生的消息速度大于消费者消费的速度.解决:增加消费者的数量或速度. 没有消费者进行消费的时候.解决:死信队列.设置消息有效期.相当于对我们的消息设置有效期,在规定的时间内如果没有消费的

  • RabbitMq消息防丢失功能实现方式讲解

    目录 1.概述 1.1.数据丢失的原因 1.2.如何防止数据丢失 2.手动应答 3.消息确认机制 3.1.AMQP事务 3.2.confirm 1.概述 1.1.数据丢失的原因 在消息中有三种可能性造成数据丢失: 消费者消费消息失败 生产者生产消息失败 MQ数据丢失 消费者消费消息失败: RabbitMq存在应答机制,默认为自动应答,MQ向消费者推送一条消息,消费者收到这条消息后会返回一个ack(应答)给MQ,MQ收到应答后会删除这条消息. 自动应答存在一个问题,就是消费者收到消息后立马就会给M

  • springboot中rabbitmq实现消息可靠性机制详解

    1. 生产者模块通过publisher confirm机制实现消息可靠性 1.1 生产者模块导入rabbitmq相关依赖 <!--AMQP依赖,包含RabbitMQ--> <dependency> <groupId>org.springframework.boot</groupId> <artifactId>spring-boot-starter-amqp</artifactId> </dependency> <!-

  • SpringBoot整合RabbitMQ消息队列的完整步骤

    SpringBoot整合RabbitMQ 主要实现RabbitMQ以下三种消息队列: 简单消息队列(演示direct模式) 基于RabbitMQ特性的延时消息队列 基于RabbitMQ相关插件的延时消息队列 公共资源 1. 引入pom依赖 <dependency> <groupId>org.springframework.boot</groupId> <artifactId>spring-boot-starter-amqp</artifactId>

  • 聊聊RabbitMQ发布确认高级问题

    目录 1.发布确认高级 1.1.发布确认SpringBoot版本 1.1.1.确认机制方案 1.1.2.代码架构图 1.1.3.配置文件 1.1.4.配置类 1.1.5.回调接口 1.1.6.生产者 1.1.7.消费者 1.1.8.测试结果 1.2.回退消息 1.2.1.Mandatory参数 1.2.2.配置文件 1.2.3.生产者代码 1.2.4.回调接口代码 1.2.5.测试结果 1.3.备份交换机 1.3.1.代码架构图 1.3.2.配置类代码 1.3.3.消费者代码 1.3.4.测试结

随机推荐