RabbitMQ实战-消费端ACK、NACK及重回队列机制

简介: RabbitMQ实战-消费端ACK、NACK及重回队列机制

0 前言

当连接失败时,消息可能还在客户端和服务器之间传输 - 它们可能处于两侧的解码或编码的中间过程,在 TCP 堆栈缓冲区中,或在电线上飞行。在这种情况下,传输中的信息将无法正常投递 - 它们需要被重新投递。Acknowledgements机制让服务器和客户端知道何时需要重新投递。

根据定义,使用消息代理(如RabbitMQ)的系统是分布式的。由于发送的协议方法(消息)不能保证到达协作方或由其成功处理,因此发布者和消费者都需要一个投递和处理确认的机制。

  • 从Consumer到 RabbitMQ 的投递处理确认,在消息协议中即acknowledgements
  • broker对publishers的确认是一个协议扩展,即publisher confirms

这两个功能都启发于 TCP。它们对于从publisher到broker和从broker到consumer的可靠投递都至关重要。即对数据安全至关重要,应用程序对数据安全的责任与broker一样多。

当 RabbitMQ 向 Con 传递消息时,它要知道何时考虑该消息才能成功发送。啥逻辑最佳取决于系统。因此,它主要是应用决定的。在 AMQP 0-9-1 中,当 Con:

  • 使用basicConsume方法进行注册
  • 或使用basicGet 方法按需获取消息

就会进行。

1 消费者确认模式和数据安全考量

当节点向Con传递消息,它必须决定该消息是否应由Con考虑处理(或至少接收)。由于多种内容(客户端连接、消费者应用等)可能会失败,因此此决定是数据安全问题。消息传递协议通常提供一个确认机制,允许Con确认交付到他们连接到的节点。是否使用该机制由Con订阅时决定。

根据使用的确认模式,RabbitMQ可考虑在消息发出后:

  • 立即成功传递(写入 TCP socket)
  • 或收到明确('manual')客户确认时。手动发送的确认可能是ack、nack,并使用以下协议方法之一:
  • basic.ack:积极确认
  • basic.nack:消极确认
  • basicReject:消极确认,但还有一个limitation

void basicReject(long deliveryTag, boolean requeue) throws IOException;

开启消费确认

spring.rabbitmq.listener.simple.acknowledge-mode=manual

Con ACK就是确认是否消费成功:

  • NONE(自动确认/不确认)- 消费者收到消息后即自动确认,无论消息是否正确处理,都不会进一步检查。可能导致某些情况下消息丢失(如消费者处理失败时,RabbitMQ仍认为消息已成功处理)
  • AUTO(自动处理确认)- RabbitMQ默认的模式。如果消费者处理消息时没有抛出异常,RabbitMQ会自动确认消息;如果处理时出现异常,消息将被重新投递,等待再次消费
  • MANUAL(手动确认)- 若抛异常,消息不会丢失,一直处Unacked状态,消息不会再次发送给其他消费者。可选择显式关闭连接,消息会恢复到Ready状态并重新投递。消费者需要显式调用ack方法确认消息成功处理。如果消费者没有确认(如抛出异常或未处理消息),消息会保持在未确认状态(Unacked),不会再次投递。关闭消费者连接时,未确认的消息会重新回到队列中。

手动确认模式(MANUAL)适用于需要更精细控制的场景,能够确保消息不会因为处理失败而丢失。

2 投递标识:Delivery Tags

如何确定投递(确认表明他们各自的投递)。

当一个 Con(订阅)被注册,MQ将使用basic.deliver方法发送(推送)消息。该方法带有delivery tag,该tag可唯一标识channel上的投递。因此,Delivery tags作用域在每个 channel 内。

Delivery Tags是单调增长的正整数,由客户库提供。客户端库方法,承认交付以交付标签作为参数。由于每个通道的递送标签范围很广,因此必须在接收的同一通道上确认交付。在不同的通道上确认将导致'未知交货标签'协议异常并关闭通道。

3 ACK投递

用于交付确认的 API 方法通常暴露为客户库中通道上的操作。Java 客户端用户将使用channel:

// 假设已有channel实例

boolean autoAck = false;

channel.basicConsume(queueName, autoAck, "a-consumer-tag",

    new DefaultConsumer(channel) {

        @Override

        public void handleDelivery(String consumerTag,

                                   Envelope envelope,

                                   AMQP.BasicProperties properties,

                                   byte[] body)

            throws IOException

        {

            long deliveryTag = envelope.getDeliveryTag();

            // positively acknowledge a single delivery, the message will

            // be discarded

            channel.basicAck(deliveryTag, false);

        }

    });

4 Acknowledging Multiple Deliveries at Once

Manual确认模式可批量进行,以减少网络流量。basicReject史上都无该字段,这就是为啥basicNack被MQ引入作为协议的扩展。

将acknowledgement方法的multiple字段置true来实现:

  • multiple=true:MQ 将确认所有未完成的delivery tag,并包括确认中指定的tag。与确认相关其他内容一样,这个作用域是channel内。比如,若channel Ch有未确认的delivery tag 5、6、7、8,当一个delivery tag=8multiple=true的acknowledgement frame到达该channel,则从 5 到 8 的所有投递都将被确认
  • multiple=false:仍不确认投递 5、6 和 7

要确认与MQ Java客户端的多次投递,将Channel#basicAck的multiple参数置true。

boolean autoAck = false;

channel.basicConsume(queueName, autoAck, "a-consumer-tag",

    new DefaultConsumer(channel) {

        @Override

        public void handleDelivery(String consumerTag,

                                   Envelope envelope,

                                   AMQP.BasicProperties properties,

                                   byte[] body)

            throws IOException

        {

            long deliveryTag = envelope.getDeliveryTag();

            // positively acknowledge all deliveries up to

            // this delivery tag

            channel.basicAck(deliveryTag, true);

        }

    });

5 NACK和Requeuing of Deliveries

有时,消费者无法及时处理投递,但其他实例可能能够处理。这时可能更想让它重新入队,让其他Con接收和处理它。basicRejectbasicNack就是用于实现这种想法的两个协议方法。这些方法通常用于消极地确认投递。

此类投递可被Broker丢弃或重新入队。此行为由requeue字段控制:

  • 当字段设置为true,Broker将用指定的delivery tag重新入队投递(或多个投递)。

这两个方法通常暴露作为客户端库中channel上的操作。Java 客户端用户可以调用:

  • Channel#basicReject
  • Channel#basicNack

boolean autoAck = false;

channel.basicConsume(queueName, autoAck, "a-consumer-tag",

    new DefaultConsumer(channel) {

        @Override

        public void handleDelivery(String consumerTag,

                                   Envelope envelope,

                                   AMQP.BasicProperties properties,

                                   byte[] body)

            throws IOException

        {

            long deliveryTag = envelope.getDeliveryTag();

            // negatively acknowledge, the message will

            // be discarded

            channel.basicReject(deliveryTag, false);

        }

    });

消费端的重回队列

重回队列针对没有处理成功的消息,将消息重新投递给Broker。重回队列会把消费失败的消息重新添加到队列尾端,供Con重新消费。一般在实际应用中,都会关闭重回队列,即设置为false。

6 RabbitMQ ACK 机制的意义

ACK机制可保证Con拉取到了消息,若处理失败了,则队列中还有这个消息,仍然可以给Con处理。

ack机制是 Con 告诉 Broker 当前消息是否成功消费,至于 Broker 如何处理 NACK,取决于 Con 是否设置了 requeue:若 requeue=false, 则NACK 后 Broker 还是会删除消息的。

但一般处理消息失败都是因为代码逻辑出bug,即使队列中后来仍然保留该消息,然后再给Con消费,依旧报错。当然,若一台机器宕机,消息还有,还可以给另外机器消费,这种情景下 ACK 很有用。

若不使用 ACK 机制,直接把出错消息存库,便于日后查bug或重新执行。 参考 Quartz 定时任务调度,Quartz可以让失败的任务重新执行一次,或者不管,或者怎么怎么样,但是 RabbitMQ 好像缺了这一点。

7 ACK和NACK

autoACK=false 时,就可用手工ACK。手工方式包括:

  • 手工ACK,会发送给Broker一个应答,代表消息处理成功,Broker就可回送响应给Pro
  • 手工NACK,表示消息处理失败,若设置了重回队列,Broker端就会将没有成功处理的消息重新发送

使用

Con消费时,若由于业务异常,可手工 NACK 记录日志,然后进行补偿

void basicNack(long deliveryTag,

         boolean multiple,

         boolean requeue)

若由于服务器宕机等严重问题,就需要手工 ACK 保障Con消费成功

void basicAck(long deliveryTag, boolean multiple)

8 实战

Con,关闭自动签收功能

/**

* ACK & 重回队列 - Con

*

* @author JavaEdge

*/

public class Consumer {

   public static void main(String[] args) throws Exception {

      ConnectionFactory connectionFactory = new ConnectionFactory();

      connectionFactory.setHost("localhost");

      connectionFactory.setPort(5672);

      connectionFactory.setVirtualHost("/");

      Connection connection = connectionFactory.newConnection();

      Channel channel = connection.createChannel();

 

      String exchangeName = "test_ack_exchange";

      String queueName = "test_ack_queue";

      String routingKey = "ack.#";

      channel.exchangeDeclare(exchangeName, "topic", true, false, null);

      channel.queueDeclare(queueName, true, false, false, null);

      channel.queueBind(queueName, exchangeName, routingKey);

      // 手工签收须关闭:autoAck = false

      channel.basicConsume(queueName, false, new MyConsumer(channel));

   }

}

对第一条消息(序号0)进行NACK,并设置重回队列:

/**

* ACK & 重回队列 - 自定义Con

*

* @author JavaEdge

*/

public class MyConsumer extends DefaultConsumer {

   private final Channel channel;

 

   public MyConsumer(Channel channel) {

       super(channel);

       this.channel = channel;

   }

 

   @Override

   public void handleDelivery(String consumerTag, Envelope envelope,

                              AMQP.BasicProperties properties, byte[] body) throws IOException {

       System.err.println("-----------Consume Message----------");

       System.err.println("body: " + new String(body));

       try {

           Thread.sleep(2000);

       } catch (InterruptedException e) {

           e.printStackTrace();

       }

       if ((Integer) properties.getHeaders().get("num") == 0) {

           channel.basicNack(envelope.getDeliveryTag(), false, true);

       } else {

           channel.basicAck(envelope.getDeliveryTag(), false);

       }

   }

}

Pro 对消息设置序号,以便区分:

/**

* ACK & 重回队列 - Pro

*

* @author JavaEdge

*/

public class Producer {

   public static void main(String[] args) throws Exception {

       ConnectionFactory connectionFactory = new ConnectionFactory();

       connectionFactory.setHost("localhost");

       connectionFactory.setPort(5672);

       connectionFactory.setVirtualHost("/");

 

       Connection connection = connectionFactory.newConnection();

       Channel channel = connection.createChannel();

 

       String exchange = "test_ack_exchange";

       String routingKey = "ack.save";

 

       for (int i = 0; i < 3; i++) {

           Map<String, Object> headers = new HashMap<>(16);

           headers.put("num", i);

           AMQP.BasicProperties properties = new AMQP.BasicProperties.Builder()

                   .deliveryMode(2)

                   .contentEncoding("UTF-8")

                   .headers(headers)

                   .build();

           String msg = "JavaEdge RabbitMQ ACK Message " + i;

           channel.basicPublish(exchange, routingKey, true, properties, msg.getBytes());

       }

   }

}

启动Con、启动Pro。这里第一条消息由于调用NACK,并设置重回队列,导致该条消息一直重复发送,消费端就会一直循环消费:


相关实践学习
深入解析Docker容器化技术
Docker是一个开源的应用容器引擎,让开发者可以打包他们的应用以及依赖包到一个可移植的容器中,然后发布到任何流行的Linux机器上,也可以实现虚拟化,容器是完全使用沙箱机制,相互之间不会有任何接口。Docker是世界领先的软件容器平台。开发人员利用Docker可以消除协作编码时“在我的机器上可正常工作”的问题。运维人员利用Docker可以在隔离容器中并行运行和管理应用,获得更好的计算密度。企业利用Docker可以构建敏捷的软件交付管道,以更快的速度、更高的安全性和可靠的信誉为Linux和Windows Server应用发布新功能。 在本套课程中,我们将全面的讲解Docker技术栈,从环境安装到容器、镜像操作以及生产环境如何部署开发的微服务应用。本课程由黑马程序员提供。 &nbsp; &nbsp; 相关的阿里云产品:容器服务 ACK 容器服务 Kubernetes 版(简称 ACK)提供高性能可伸缩的容器应用管理能力,支持企业级容器化应用的全生命周期管理。整合阿里云虚拟化、存储、网络和安全能力,打造云端最佳容器化应用运行环境。 了解产品详情: https://www.aliyun.com/product/kubernetes
目录
相关文章
|
消息中间件 大数据 关系型数据库
RocketMQ实战—3.基于RocketMQ升级订单系统架构
本文主要介绍了基于MQ实现订单系统核心流程的异步化改造、基于MQ实现订单系统和第三方系统的解耦、基于MQ实现将订单数据同步给大数据团队、秒杀系统的技术难点以及秒杀商详页的架构设计和基于MQ实现秒杀系统的异步化架构。
947 64
RocketMQ实战—3.基于RocketMQ升级订单系统架构
|
消息中间件 Java 数据库
RocketMQ实战—9.营销系统代码初版
本文主要介绍了实现营销系统四大促销场景的代码初版:全量用户推送促销活动、全量用户发放优惠券、特定用户推送领取优惠券消息、热门商品定时推送。
RocketMQ实战—9.营销系统代码初版
|
消息中间件 搜索推荐 调度
RocketMQ实战—8.营销系统业务和方案介绍
本文详细介绍了电商营销系统的业务流程、技术架构及挑战解决方案。涵盖核心交易与支付后履约流程,优惠券和促销活动的发券、领券、用券、销券机制,以及会员与推送的数据库设计。技术架构基于Nacos服务注册中心、Dubbo RPC框架、RocketMQ消息中间件和XXLJob分布式调度工具,实现系统间高效通信与任务管理。针对千万级用户量下的推送和发券场景,提出异步化、分片处理与惰性发券等优化方案,解决高并发压力。同时,通过RocketMQ实现系统解耦,提升扩展性,并利用XXLJob完成爆款商品推荐的分布式调度推送。整体设计确保系统在大规模用户场景下的性能与稳定性。
RocketMQ实战—8.营销系统业务和方案介绍
|
消息中间件 存储 NoSQL
RocketMQ实战—6.生产优化及运维方案
本文围绕RocketMQ集群的使用与优化,详细探讨了六个关键问题。首先,介绍了如何通过ACL配置实现RocketMQ集群的权限控制,防止不同团队间误用Topic。其次,讲解了消息轨迹功能的开启与追踪流程,帮助定位和排查问题。接着,分析了百万消息积压的处理方法,包括直接丢弃、扩容消费者或通过新Topic间接扩容等策略。此外,提出了针对RocketMQ集群崩溃的金融级高可用方案,确保消息不丢失。同时,讨论了为RocketMQ增加限流功能的重要性及实现方式,以提升系统稳定性。最后,分享了从Kafka迁移到RocketMQ的双写双读方案,确保数据一致性与平稳过渡。
|
10月前
|
消息中间件 Ubuntu Java
SpringBoot整合MQTT实战:基于EMQX实现双向设备通信
本教程指导在Ubuntu上部署EMQX 5.9.0并集成Spring Boot实现MQTT双向通信,涵盖服务器搭建、客户端配置及生产实践,助您快速构建企业级物联网消息系统。
3029 1
|
消息中间件 Java 中间件
RocketMQ实战—2.RocketMQ集群生产部署
本文主要介绍了大纲什么是消息中间件、消息中间件的技术选型、RocketMQ的架构原理和使用方式、消息中间件路由中心的架构原理、Broker的主从架构原理、高可用的消息中间件生产部署架构、部署一个小规模的RocketMQ集群进行压测、如何对RocketMQ集群进行可视化的监控和管理、进行OS内核参数和JVM参数的调整、如何对小规模RocketMQ集群进行压测、消息中间件集群生产部署规划梳理。
RocketMQ实战—2.RocketMQ集群生产部署
|
消息中间件 NoSQL 大数据
RocketMQ实战—5.消息重复+乱序+延迟的处理
本文围绕RocketMQ的使用与优化展开,分析了优惠券重复发放的原因及解决方案。首先,通过案例说明了优惠券系统因消息重复、数据库宕机或消费失败等原因导致重复发券的问题,并提出引入幂等性机制(如业务判断法、Redis状态判断法)来保证数据唯一性。其次,探讨了死信队列在处理消费失败时的作用,以及如何通过重试和死信队列解决消息处理异常。接着,分析了订单库同步中消息乱序的原因,提出了基于顺序消息机制的代码实现方案,确保消息按序处理。此外,介绍了利用Tag和属性过滤数据提升效率的方法,以及延迟消息机制优化定时退款扫描的功能。最后,总结了RocketMQ生产实践中的经验.
RocketMQ实战—5.消息重复+乱序+延迟的处理
|
存储 Kubernetes 监控
K8s集群实战:使用kubeadm和kuboard部署Kubernetes集群
总之,使用kubeadm和kuboard部署K8s集群就像回归童年一样,简单又有趣。不要忘记,技术是为人服务的,用K8s集群操控云端资源,我们不过是想在复杂的世界找寻简单。尽管部署过程可能遇到困难,但朝着简化复杂的目标,我们就能找到意义和乐趣。希望你也能利用这些工具,找到你的乐趣,满足你的需求。
1285 33
|
消息中间件 Java 测试技术
RocketMQ实战—7.生产集群部署和生产参数
本文详细介绍了RocketMQ生产集群的部署与调优过程,包括集群规划、环境搭建、参数配置和优化策略。
RocketMQ实战—7.生产集群部署和生产参数

推荐镜像

更多