RabbitMQ学习笔记:生产者消息确认

2024-08-27 15:48

本文主要是介绍RabbitMQ学习笔记:生产者消息确认,希望对大家解决编程问题提供一定的参考价值,需要的开发者们随着小编来一起学习吧!

环境

window10
虚拟机、secureCRT
Intellij IDEA

消息确认

消息端的消息确认 我们知道可以使用basic.ack()来确认消息被成功消费了;

但是发送端,即生产者如何知道自己的消息成功的发送到了RabbitMQ服务器上呢?

RabbitMQ提供了两种方法:
① 事务确认机制
② 发送方确认机制

事务确认机制

流程如下:
在这里插入图片描述
RabbitMQ客户端中和事务机制有关的方法有三个:
channel.txSelect : 将当前信道设置为事务模式;
channel.txCommit:用于提交事务;
channel.txRollback:用于事务回滚;

channel.txSelect();
channel.basicPublish("", "", MessageProperties.PERSISTENT_TEXT_PLAIN,"yutao".getBytes());
channel.txCommit();

结合上图我们可以得到如下信息:

  • 客户端发送Tx.Select, 将信道设置为事务模式;
  • Broker回复Tx.Select-Ok,确认已将信道设置为事务模式
  • 在发送消息之后,客户端发送Tx.Commit提交事务;
  • Broker回复Tx.Commit-Ok,确认事务提交

接下来看看,带回滚的代码:

try {
channel.txSelect();
channel.basicPublish("", "", MessageProperties.PERSISTENT_TEXT_PLAIN,"yutao".getBytes());
// 分母不能为0 ,肯定会抛异常
int result = 1/0;
channel.txCommit();
} catch (Exception e) {e.printStackTrace();channel.txRollback();
}

回滚流程:
在这里插入图片描述

可以看出,事务机制可以解决生产者消息确认的问题,但是使用事务机制会吸干RabbitMQ的性能;

发送方确认机制

针对事务机制带来的性能问题,rabbitmq又提供确认confirm模式来解决上面的问题;

当将信道设置为confirm模式后,所有在该信道上面发布的消息都会被指派一个ID(从1开始),一旦消息被投递到所有匹配的队列之后,RabbitMQ就会发送一个确认Basic.Ack给生产者(包含唯一ID);如果消息和队列设置了持久化,那么就会等到消息落盘后才会发送确认;

具体代码:

/*** 单个确认* 只给出关键代码*/
public void confirm() {try {// 将信道设置为publisher confirm模式channel.confirmSelect();channel.basicPublish("exchange", "rountingKey", null, "publisher confirm test".getBytes());if (!channel.waitForConfirms()) {System.out.println("send message failed");}} catch (IOException | InterruptedException e) {e.printStackTrace();}
}

注意地方:

channel.waitForConfirms() 这个方法比较关键;
官方文档上注释是:等到自上次调用后发布的所有消息都被代理确认或取消。注意,当在非确认通道上调用时,waitForConfirms抛出一个illeglastateException。
也就是说这个方法 在发送消息后,调用这个方法会阻塞线程,等待rabbitmq的回复确认。

但是如果我们想批量发送消息呢?很简单,只需要把上面的代码用循环包裹住就可以了。

但是上面的代码,很容易看出,是串行的,即发送一条,确认一条;这样时非常影响性能的,和上面讲到的事务机制类似—耗性能。

针对这种情况,我们可以先发送一批数据,再去调用channel.waitForConfirms()来等到确认;

批量确认

public void batchConfirm() {int BATCH_COUNT = 100;List<String> messages = new ArrayList<>();String message = "batch confirm test";try {channel.confirmSelect();int msgCount = 0;while (true) {channel.basicPublish("", "", null, message.getBytes());messages.add(message);if (++msgCount >= BATCH_COUNT) {msgCount = 0;try {// 先发送100条,然后在等待确认// 失败就得进行重发,所以可能会有消息重复问题if (channel.waitForConfirms()) {// 清空缓存里的消息messages.clear();continue;}// 将缓存里的消息进行重发} catch (InterruptedException e) {e.printStackTrace();// 将缓存里的消息进行重发}}}} catch (IOException e) {e.printStackTrace();}
}

思路:先发送一批数据,然后再去确认;但是这样有缺点,如果这一批里面失败了,就得重新发送,这样消息就可能会有重复;如果消息经常丢失时,性能不升反降;

异步确认

针对批量确认代码的问题,rabbitmq提供了异步确认的方法;

private static SortedSet<Object> confirmSet = Collections.synchronizedSortedSet(new TreeSet<>());while(true) {long nextPublishSeqNo = channel.getNextPublishSeqNo();channel.basicPublish("", "", MessageProperties.PERSISTENT_TEXT_PLAIN,"yutao".getBytes());confirmSet.add(nextPublishSeqNo);
}/**
* 异步确认*/
public void asyncConfirm() {try {channel.confirmSelect();channel.addConfirmListener(new ConfirmListener() {@Overridepublic void handleAck(long deliveryTag, boolean multiple) throws IOException {System.out.println("Nack, SeqNo: " + deliveryTag + ", multiple: " + multiple);if (multiple) {confirmSet.headSet(deliveryTag + 1).clear();} else {confirmSet.remove(deliveryTag);}}@Overridepublic void handleNack(long deliveryTag, boolean multiple) throws IOException {if (multiple) {confirmSet.headSet(deliveryTag + 1).clear();} else {confirmSet.remove(deliveryTag);}}});} catch (IOException e) {e.printStackTrace();}
}

异步confirm模式的编程实现最复杂,Channel对象提供的ConfirmListener()回调方法只包含deliveryTag(当前Chanel发出的消息序号),我们需要自己为每一个Channel维护一个unconfirm的消息序号集合,每publish一条数据,集合中元素加1,每回调一次handleAck方法,unconfirm集合删掉相应的一条(multiple=false)或多条(multiple=true)记录。从程序运行效率上看,这个unconfirm集合最好采用有序集合SortedSet存储结构。实际上,SDK中的waitForConfirms()方法也是通过SortedSet维护消息序号的。

从代码里可以看出,其关键的地方就是channel.addConfirmListener(ConfirmListener listener)这个方法;
这是Channel对象提供的回调方法,需要实现ConfirmListener这个接口,并通过回调方法中的deliveryTag来确认消息。
其实这也是一个监听器,监听器这种东西,说白了里面就是一个死循环,等待rabbitmq服务器端的确认消息。当收到了确认消息后,业务上开发人员还需要做点事情,所以就暴露了一个回调接口给开发人员用,通过这个回调接口,开发人员可以对知道哪些消息确认了,哪些没有确认。

参考地址:

RabbitMQ之消息确认机制(事务+Confirm)

《RabbitMQ实战指南》

这篇关于RabbitMQ学习笔记:生产者消息确认的文章就介绍到这儿,希望我们推荐的文章对编程师们有所帮助!



http://www.chinasem.cn/article/1112101

相关文章

如何通过Python实现一个消息队列

《如何通过Python实现一个消息队列》这篇文章主要为大家详细介绍了如何通过Python实现一个简单的消息队列,文中的示例代码讲解详细,感兴趣的小伙伴可以跟随小编一起学习一下... 目录如何通过 python 实现消息队列如何把 http 请求放在队列中执行1. 使用 queue.Queue 和 reque

Java深度学习库DJL实现Python的NumPy方式

《Java深度学习库DJL实现Python的NumPy方式》本文介绍了DJL库的背景和基本功能,包括NDArray的创建、数学运算、数据获取和设置等,同时,还展示了如何使用NDArray进行数据预处理... 目录1 NDArray 的背景介绍1.1 架构2 JavaDJL使用2.1 安装DJL2.2 基本操

解读Redis秒杀优化方案(阻塞队列+基于Stream流的消息队列)

《解读Redis秒杀优化方案(阻塞队列+基于Stream流的消息队列)》该文章介绍了使用Redis的阻塞队列和Stream流的消息队列来优化秒杀系统的方案,通过将秒杀流程拆分为两条流水线,使用Redi... 目录Redis秒杀优化方案(阻塞队列+Stream流的消息队列)什么是消息队列?消费者组的工作方式每

使用C/C++调用libcurl调试消息的方式

《使用C/C++调用libcurl调试消息的方式》在使用C/C++调用libcurl进行HTTP请求时,有时我们需要查看请求的/应答消息的内容(包括请求头和请求体)以方便调试,libcurl提供了多种... 目录1. libcurl 调试工具简介2. 输出请求消息使用 CURLOPT_VERBOSE使用 C

微服务架构之使用RabbitMQ进行异步处理方式

《微服务架构之使用RabbitMQ进行异步处理方式》本文介绍了RabbitMQ的基本概念、异步调用处理逻辑、RabbitMQ的基本使用方法以及在SpringBoot项目中使用RabbitMQ解决高并发... 目录一.什么是RabbitMQ?二.异步调用处理逻辑:三.RabbitMQ的基本使用1.安装2.架构

Java中Springboot集成Kafka实现消息发送和接收功能

《Java中Springboot集成Kafka实现消息发送和接收功能》Kafka是一个高吞吐量的分布式发布-订阅消息系统,主要用于处理大规模数据流,它由生产者、消费者、主题、分区和代理等组件构成,Ka... 目录一、Kafka 简介二、Kafka 功能三、POM依赖四、配置文件五、生产者六、消费者一、Kaf

Springboot使用RabbitMQ实现关闭超时订单(示例详解)

《Springboot使用RabbitMQ实现关闭超时订单(示例详解)》介绍了如何在SpringBoot项目中使用RabbitMQ实现订单的延时处理和超时关闭,通过配置RabbitMQ的交换机、队列和... 目录1.maven中引入rabbitmq的依赖:2.application.yml中进行rabbit

SpringBoot 自定义消息转换器使用详解

《SpringBoot自定义消息转换器使用详解》本文详细介绍了SpringBoot消息转换器的知识,并通过案例操作演示了如何进行自定义消息转换器的定制开发和使用,感兴趣的朋友一起看看吧... 目录一、前言二、SpringBoot 内容协商介绍2.1 什么是内容协商2.2 内容协商机制深入理解2.2.1 内容

SpringBoot整合Canal+RabbitMQ监听数据变更详解

《SpringBoot整合Canal+RabbitMQ监听数据变更详解》在现代分布式系统中,实时获取数据库的变更信息是一个常见的需求,本文将介绍SpringBoot如何通过整合Canal和Rabbit... 目录需求步骤环境搭建整合SpringBoot与Canal实现客户端Canal整合RabbitMQSp

HarmonyOS学习(七)——UI(五)常用布局总结

自适应布局 1.1、线性布局(LinearLayout) 通过线性容器Row和Column实现线性布局。Column容器内的子组件按照垂直方向排列,Row组件中的子组件按照水平方向排列。 属性说明space通过space参数设置主轴上子组件的间距,达到各子组件在排列上的等间距效果alignItems设置子组件在交叉轴上的对齐方式,且在各类尺寸屏幕上表现一致,其中交叉轴为垂直时,取值为Vert