前言

上一章讲解了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消息确认的资料请关注Devmax其它相关文章!

RabbitMQ消息确认机制剖析的更多相关文章

  1. Swift回调及notifition消息机制

    overrideinit(){}//定义一个方法执行协议的方法funcdoSomething(){iflet_=self.delegate{delegate!

  2. 深入理解 Swift 派发机制

    然而,只要缓存建立了起来,这个查找过程就会通过缓存来把性能提高到和函数表派发一样快.但这只是消息机制的原理,这里有一篇文章很深入的讲解了具体的技术细节.Swift的派发机制那么,到底Swift是怎么派发的呢?

  3. python操作RabbitMq的三种工作模式

    这篇文章主要为大家介绍了python操作RabbitMq的三种工作模式,有需要的朋友可以借鉴参考下,希望能够有所帮助,祝大家多多进步早日升职加薪

  4. Android消息机制Handler深入理解

    这篇文章介绍了深入理解Android消息机制Handler,对大家的学习或者工作具有一定的参考学习价值,需要的朋友们下面随着小编来一起学习学习吧

  5. springboot-rabbitmq-reply 消息直接回复模式详情

    这篇文章主要介绍了springboot-rabbitmq-reply消息直接回复模式详情,文章通过围绕主题展开详细的内容介绍,具有一定的参考价值,感兴趣的小伙伴可以参考一下

  6. SpringBoot+RabbitMQ 实现死信队列的示例

    本文主要介绍了SpringBoot+RabbitMQ 实现死信队列的示例,文中通过示例代码介绍的非常详细,对大家的学习或者工作具有一定的参考学习价值,需要的朋友们下面随着小编来一起学习学习吧

  7. PHP实现RabbitMQ消息列队的示例代码

    众所周知,php本身的运行效率存在一定的缺陷,所以如果有一个很复杂很耗时的业务时,必须开发一个常驻内存的程序。本文将利用PHP实现RabbitMQ消息列队,感兴趣的可以了解一下

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

    消息队列是最古老的中间件之一,从系统之间有通信需求开始,就自然产生了消息队列。本文告诉什么是消息队列,为什么需要消息队列,常见的消息队列有哪些,RabbitMQ的部署和使用

  9. SpringBoot整合Canal与RabbitMQ监听数据变更记录

    这篇文章主要介绍了SpringBoot整合Canal与RabbitMQ监听数据变更记录,文章围绕主题展开详细的内容介绍,具有一定的参考价值,需要的小伙伴可以参考一下

  10. RabbitMQ消息确认机制剖析

    这篇文章主要为大家介绍了RabbitMQ消息确认机制剖析,有需要的朋友可以借鉴参考下,希望能够有所帮助,祝大家多多进步,早日升职加薪

随机推荐

  1. 基于EJB技术的商务预订系统的开发

    用EJB结构开发的应用程序是可伸缩的、事务型的、多用户安全的。总的来说,EJB是一个组件事务监控的标准服务器端的组件模型。基于EJB技术的系统结构模型EJB结构是一个服务端组件结构,是一个层次性结构,其结构模型如图1所示。图2:商务预订系统的构架EntityBean是为了现实世界的对象建造的模型,这些对象通常是数据库的一些持久记录。

  2. Java利用POI实现导入导出Excel表格

    这篇文章主要为大家详细介绍了Java利用POI实现导入导出Excel表格,文中示例代码介绍的非常详细,具有一定的参考价值,感兴趣的小伙伴们可以参考一下

  3. Mybatis分页插件PageHelper手写实现示例

    这篇文章主要为大家介绍了Mybatis分页插件PageHelper手写实现示例,有需要的朋友可以借鉴参考下,希望能够有所帮助,祝大家多多进步,早日升职加薪

  4. (jsp/html)网页上嵌入播放器(常用播放器代码整理)

    网页上嵌入播放器,只要在HTML上添加以上代码就OK了,下面整理了一些常用的播放器代码,总有一款适合你,感兴趣的朋友可以参考下哈,希望对你有所帮助

  5. Java 阻塞队列BlockingQueue详解

    本文详细介绍了BlockingQueue家庭中的所有成员,包括他们各自的功能以及常见使用场景,通过实例代码介绍了Java 阻塞队列BlockingQueue的相关知识,需要的朋友可以参考下

  6. Java异常Exception详细讲解

    异常就是不正常,比如当我们身体出现了异常我们会根据身体情况选择喝开水、吃药、看病、等 异常处理方法。 java异常处理机制是我们java语言使用异常处理机制为程序提供了错误处理的能力,程序出现的错误,程序可以安全的退出,以保证程序正常的运行等

  7. Java Bean 作用域及它的几种类型介绍

    这篇文章主要介绍了Java Bean作用域及它的几种类型介绍,Spring框架作为一个管理Bean的IoC容器,那么Bean自然是Spring中的重要资源了,那Bean的作用域又是什么,接下来我们一起进入文章详细学习吧

  8. 面试突击之跨域问题的解决方案详解

    跨域问题本质是浏览器的一种保护机制,它的初衷是为了保证用户的安全,防止恶意网站窃取数据。那怎么解决这个问题呢?接下来我们一起来看

  9. Mybatis-Plus接口BaseMapper与Services使用详解

    这篇文章主要为大家介绍了Mybatis-Plus接口BaseMapper与Services使用详解,有需要的朋友可以借鉴参考下,希望能够有所帮助,祝大家多多进步,早日升职加薪

  10. mybatis-plus雪花算法增强idworker的实现

    今天聊聊在mybatis-plus中引入分布式ID生成框架idworker,进一步增强实现生成分布式唯一ID,具有一定的参考价值,感兴趣的小伙伴们可以参考一下

返回
顶部