RabbitMQ消息应答与发布

2024-01-24 03:28
文章标签 发布 消息 rabbitmq 应答

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

消息应答

RabbitMQ一旦向消费者发送了一个消息,便立即将该消息,标记为删除.

消费者完成一个任务可能需要一段时间,如果其中一个消费者处理一个很长的任务并仅仅执行了一半就突然挂掉了,在这种情况下,我们将丢失正在处理的消息,后续给消费者发送的消息也就无法接收到了.

为了确保消息不丢失,我们引入了消息应答机制.

消息应答就是:消费者在接收到生产者的消息并且处理该消息之后,告诉RabbitMQ已经处理完成了,rabbitMQ可以进行删除了.

一个生产者,两个消费者.

/**
消费者1
*/
public class Consumer1 {public static final String QueueName = "TaskQueueName";public static void main(String[] args) throws IOException, TimeoutException {ConnectionFactory factory = new ConnectionFactory();factory.setHost("127.0.0.1");factory.setUsername("DGZ");factory.setPassword("Dgz@#151");Connection connection = factory.newConnection();Channel channel = connection.createChannel();//声明接收消息System.out.println("消费者1等待接收消息时间较短");DeliverCallback deliverCallback = (consumerTag,message) -> {Util.sleep(1);  //等待1sSystem.out.println("正在处理该消息" + new String(message.getBody()));/*** 手动应答* 1.消息的标记Tag* 2.是否批量应答 false表示不批量应答信道中的消息*/channel.basicAck(message.getEnvelope().getDeliveryTag(),false);};//取消消息时的回调CancelCallback cancelCallback = consumerTag ->{System.out.println("消息消费被中断");};/*** 消费者消费消息* 1.消费哪个队列* 2.消费成功之后是否要自动应答true:代表自动应答      false:代表手动应答* 3.消费者未成功消费的回调* 4.消费者取消消费的回调*/boolean autoAck = false;channel.basicConsume(QueueName,autoAck,deliverCallback,cancelCallback);}
}

/**
消费者2
*/
public class Consumer2 {public static final String QueueName = "TaskQueueName";public static void main(String[] args) throws IOException, TimeoutException {ConnectionFactory factory = new ConnectionFactory();factory.setHost("127.0.0.1");factory.setUsername("DGZ");factory.setPassword("Dgz@#151");Connection connection = factory.newConnection();Channel channel = connection.createChannel();System.out.println("消费者2等待接收消息时间较长");//声明接收消息DeliverCallback deliverCallback = (consumerTag,message) -> {Util.sleep(10); //等待10sSystem.out.println("正在处理该消息" + new String(message.getBody()));/*** 手动应答* 1.消息的标记Tag* 2.是否批量应答 false表示不批量应答信道中的消息*///当该消息还没有被处理的时候,如果此时这个应用挂掉,// 由于这个手动应答的机制,就不会删除该消息,而是将给消息交给其他应用去处理channel.basicAck(message.getEnvelope().getDeliveryTag(),false);};//取消消息时的回调CancelCallback cancelCallback = consumerTag ->{System.out.println("消息消费被中断");};/*** 消费者消费消息* 1.消费哪个队列* 2.消费成功之后是否要自动应答true:代表自动应答     false:代表手动应答* 3.消费者未成功消费的回调* 4.消费者取消消费的回调*/boolean autoAck = false;channel.basicConsume(QueueName,autoAck,deliverCallback,cancelCallback);}
}

消费者1 1s处理一个消息,消费者2 10s处理一个消息

/**
生产者
*/
public class Producer {//队列名称public static final String QueueName = "TaskQueueName";public static void main(String[] args) throws IOException, TimeoutException, InterruptedException {ConnectionFactory factory = new ConnectionFactory();   //创建连接工厂factory.setHost("127.0.0.1");  //主机factory.setUsername("DGZ"); //用户名factory.setPassword("Dgz@#151"); // 密码Connection connection = factory.newConnection();  //通过连接工厂创建一个连接Channel channel = connection.createChannel();  //获取信道/*** 生产一个对列* 1.对列名称* 2.对列里面的消息是否持久化,默认情况下,消息存储在内存中* 3.该队列是否只供一个消费者进行消费,是否进行消息共享,true可以多个消费者消费 false:只能一个消费者消费* 4.是否自动删除,最后一个消费者端开链接以后,该队列是否自动删除,true表示自动删除* 5.其他参数*/channel.queueDeclare(QueueName,false,false,false,null);Scanner scanner = new Scanner(System.in);System.out.println("请输入消息:");while(scanner.hasNext()) {String message = scanner.next();channel.basicPublish("",QueueName,null,message.getBytes());System.out.println("生产者发出消息: " + message);}/*** 发送一个消息* 1.发送到哪个交换机* 2.路由的key值是哪个本次是队列的名称* 3.其他参数信息* 4.发送消息的消息体*/System.out.println("消息发送完毕");}
}

自动应答和手动应答:

自动应答就是MQ只要把消息发出去,不会管消息是否收到,都会立刻把这个消息进行删除.

手动应答就是MQ把消息发送出去之后,消费者在接收到生产者的消息并且处理该消息之后,告诉RabbitMQ已经处理完成了,rabbitMQ可以进行删除了.

RabbitMQ持久化

当RabbitMQ服务突然挂掉之后,消息生产者发送过来的消息如何保证不丢失.

默认情况下当RabbitMQ服务挂掉之后,它会忽略队列和消息,这时刚刚发送的消息和队列都会丢失.

确保消息不丢失我们需要干两件事,就是将队列和消息标记为持久化.

队列持久化

之前创建的队列都是非持久化的,如果RabbitMQ重启的话,队列就是丢失,如果需要实现持久化队列,那么就需要在声明队列的时候把durable参数设置为true.表示代表开启持久化.

public class Producer {public static final String QueueName = "CJ_QUEUE";public static void main(String[] args) throws IOException, TimeoutException {/*** RabbitMQ工具类,用来创建连接等信息*/Channel channel = RabbitMQUtil.RabbitMQ_getChannel();boolean durable = true;//第二个参数为队列持久化的参数,设置为true,表示队列开启持久化,false表示不开启持久化channel.queueDeclare(QueueName,durable,false,false,null);}
}

此时web管理端就会看见:

 表示队列持久化

 消息持久化

消息持久化是指消息生产者发布消息的时候,开启消息持久化,

public class Producer {public static final String QueueName = "CJ_QUEUE";public static void main(String[] args) throws IOException, TimeoutException {/*** RabbitMQ工具类,用来创建连接等信息*/Channel channel = RabbitMQUtil.RabbitMQ_getChannel();boolean durable = true;//第二个参数为队列持久化的参数,设置为true,表示队列开启持久化,false表示不开启持久化channel.queueDeclare(QueueName,durable,false,false,null);Scanner scanner = new Scanner(System.in);System.out.println("请输入要发送的消息");while (scanner.hasNext()) {String message = scanner.next();/*** 将发布消息的basicPublish方法的第三个参数设置为MessageProperties.PERSISTENT_TEXT_PLAIN* 表示开启消息持久化*/channel.basicPublish("",QueueName, MessageProperties.PERSISTENT_TEXT_PLAIN,message.getBytes("UTF-8"));System.out.println("生产者发出消息" + message);}}
}

不公平分发

在最开始的时候RabbitMQ采用的是轮询分发,但是在某种场景下这种策略并不是很好,比如说两个消费者在处理消息,此时消费者1处理消息比较慢,而消费者2处理消息比较快,如果这个时候还是采用轮询分发的方式,那么处理慢的消费者就会一直在处理消息,而处理快的消费者就会有很大时间处于空闲状态.

为了解决这个问题,引入了不公平分发

不公平分发: 如果一个工作队列(消费者)还没有处理完一个消息或者没有应答签收一个消息,则RabbitMQ不会分配新的消息给该队列.

如果所有的消费者都没有完成手上的消息,生产者还在不停地生产消息,队列还在不停地添加新任务,这是不会给消费者分发消息,就有可能导致队列被撑爆.

这时就只能添加新的工作队列或者改变存储策略了.

设置不公平分发:

public class Consumer1 {public static final String QueueName = "CJ_QUEUE";public static void main(String[] args) throws IOException, TimeoutException {Channel channel = RabbitMQUtil.RabbitMQ_getChannel();System.out.println("消费消息时间较长-----------");DeliverCallback deliverCallback = ((consumerTag,message)->{RabbitMQUtil.sleep(10);System.out.println("消费消息:" + new String(message.getBody()));channel.basicAck(message.getEnvelope().getDeliveryTag(),false);});CancelCallback cancelCallback = consumerTag ->{System.out.println("消息消费被中断");};//设置不公平分发int prefetchCount = 1;channel.basicQos(prefetchCount);//手动应答channel.basicConsume(QueueName,false,deliverCallback,cancelCallback);}
}

public class Consumer2 {public static final String QueueName = "CJ_QUEUE";public static void main(String[] args) throws IOException, TimeoutException {Channel channel = RabbitMQUtil.RabbitMQ_getChannel();System.out.println("消费消息时间较短-----------");DeliverCallback deliverCallback = ((consumerTag,message)->{RabbitMQUtil.sleep(1);System.out.println("消费消息:" + new String(message.getBody()));channel.basicAck(message.getEnvelope().getDeliveryTag(),false);});CancelCallback cancelCallback = consumerTag ->{System.out.println("消息消费被中断");};//设置不公平分发int prefetchCount = 1;channel.basicQos(prefetchCount);//采取手动应答channel.basicConsume(QueueName,false,deliverCallback,cancelCallback);}
}

预期值分发

带权的消息分发

默认的消息发送是异步的,所以在任何时候,channel中不止一个来自消费者收到确认的消息,因此这里就存在一个未确认的消息缓冲区.

因此我们希望限制这里的缓冲区的大小,避免缓冲区中无休止的未确认消息.

这时我们就可以通过basicqos()方法来设置(预取计数来完成).

basicqos方法里面设置通道上允许的未确认消息的最大数量,一旦数据达到配置的数量,RabbitMQ将停止在通道上传递更多的消息.除非有未确认的消息被确认.

例如:假设此时通道上未确认的消息有 4,6,7,9,5,10,并且通道上设置预取计数值为6,这时RabbitMQ将不会在该通道上传递消息.除非这里的消息被确认一个,RabbitMQ将感知到这一变化,并且在发送一条信息.

消息应答和Qos预取值对用户的吞吐量有着重大影响.

不公平分发和预取值分发都用到 basic.qos 方法,如果取值为 1,代表不公平分发,取值不为1,代表预取值分发 

public class Consumer2 {public static final String QueueName = "CJ_QUEUE";public static void main(String[] args) throws IOException, TimeoutException {Channel channel = RabbitMQUtil.RabbitMQ_getChannel();System.out.println("消费消息时间较短-----------");DeliverCallback deliverCallback = ((consumerTag,message)->{RabbitMQUtil.sleep(1);System.out.println("消费消息:" + new String(message.getBody()));channel.basicAck(message.getEnvelope().getDeliveryTag(),false);});CancelCallback cancelCallback = consumerTag ->{System.out.println("消息消费被中断");};//设置不公平分发/*int prefetchCount = 1;channel.basicQos(prefetchCount);*///设置预期值分发 值为4int prefetchCount = 4;channel.basicQos(prefetchCount);//采取手动应答channel.basicConsume(QueueName,false,deliverCallback,cancelCallback);}
}

发布确认

生产者发送消息到RabbitMQ后,需要RabbitMQ返回一个ack,表示RabbitMQ已经收到生产者发送的消息.这样生产者就知道自己发送的信息成功了.

发布确认逻辑

如果消息和队列是可持久化的,那么确认消息会在消息写入磁盘之后发出.

broker 回传给生产者的确认消息中 delivery-tag 域包含了确认消息的序列号,此外 broker 也可以设置 basic.ack 的 multiple 域,表示到这个序列号之前的所有消息都已经得到了处理。

confirm 模式最大的好处在于是异步的,一旦发布一条消息,生产者应用程序就可以在等信道返回确认的同时继续发送下一条消息,当消息最终得到确认之后,生产者应用便可以通过回调方法来处理该确认消息,如果RabbitMQ 因为自身内部错误导致消息丢失,就会发送一条 nack 消息, 生产者应用程序同样可以在回调方法中处理该 nack 消息。

发布确认默认是没有开启的,如果要开启需要调用方法confirmSelect(),

channel.confirmSelect();

单个发布确认

这是一种简单的确认发布,他是一种同步确认发布的方式,也就是发布一个消息之后,只有它被确认,后续的消息才能被发布出去.

waitforconfirms()这个方法只有在消息被确认的时候才返回.

如果在指定的时间范围内没有返回,则抛出异常.

这种确认的方式就是发布速度特别慢,因为如果没有确认发布的消息,那么其他消息就只能阻塞等待.

这种方式最多每秒不超过百条的发送量.


/*** Created with IntelliJ IDEA.** @Author: DongGuoZhen* @Date: 2024/01/05/10:28* @Description: 单个确认发布*/
//确认发布指的是成功发送到了队列,并不是消费者消费了消息。
public class Producer1 {public static final int MESSAGE_MAX = 1000;public static void main(String[] args) throws IOException, TimeoutException, InterruptedException {pushMessage();}private static void pushMessage() throws IOException, TimeoutException, InterruptedException {Channel channel = RabbitMQUtil.RabbitMQ_getChannel();String QueueName = UUID.randomUUID().toString();channel.queueDeclare(QueueName,false,true,false,null);//开启发布确认channel.confirmSelect();//开始时间long start = System.currentTimeMillis();//依次发送1000个消息for (int i = 0; i < 1000; i++) {String message = i+"";channel.basicPublish("",QueueName,null,message.getBytes());//发送单个消息立马进行发布确认boolean flag = channel.waitForConfirms();if(flag) {  //如果成功 trueSystem.out.println("消息: "+ message + " 成功发送到队列:" + QueueName);}}long end = System.currentTimeMillis();System.out.println("发布: " + MESSAGE_MAX + "条消息共耗时:"+ (end-start)+ "ms");}
}

 

批量发布确认

单个确认发布的速度非常慢,与单个确认等待相比,如果发送一批然后在一起确认,这样就大大提高了消息的发送速度和吞吐量.

当然这种方式的缺点就是,当有一个消息出现问题时,我们无法知道是那个消息没有发布出去.

/*** Created with IntelliJ IDEA.** @Author: DongGuoZhen* @Date: 2024/01/05/10:28* @Description: 批量确认发布*/
//确认发布指的是成功发送到了队列,并不是消费者消费了消息。
public class Producer2 {public static final int MESSAGE_MAX = 1000;public static void main(String[] args) throws IOException, TimeoutException, InterruptedException {pushMessage();}private static void pushMessage() throws IOException, TimeoutException, InterruptedException {Channel channel = RabbitMQUtil.RabbitMQ_getChannel();String QueueName = UUID.randomUUID().toString();channel.queueDeclare(QueueName,false,true,false,null);//开启发布确认channel.confirmSelect();int bachSize = 100;  //批量发布确认的条数//开始时间long start = System.currentTimeMillis();//依次发送1000个消息for (int i = 0; i < 1000; i++) {String message = i+"";channel.basicPublish("",QueueName,null,message.getBytes());if((i+1) % bachSize == 0) {  //当发送的消息达到100条,进行确认发布一次channel.waitForConfirms();   //发布确认System.out.println("第"+ i+"条消息被确认");}}long end = System.currentTimeMillis();System.out.println("发布: " + MESSAGE_MAX + "条消息共耗时:"+ (end-start)+ "ms");}
}

 

异步发布确认

异步确认虽然编程上逻辑比较复杂,但是性价比最高,可靠性和效率都很好,利用了回调函数来达到消息的可靠性传递.


/*** Created with IntelliJ IDEA.** @Author: DongGuoZhen* @Date: 2024/01/05/11:02* @Description: 异步确认发布*/
public class Producer3 {public static final int MESSAGE_MAX = 1000;public static void main(String[] args) throws IOException, TimeoutException {pushMessage();}private static void pushMessage() throws IOException, TimeoutException {Channel channel = RabbitMQUtil.RabbitMQ_getChannel();String QueueName = UUID.randomUUID().toString();channel.queueDeclare(QueueName,false,true,false,null);channel.confirmSelect();long start = System.currentTimeMillis();/*** deliveryTag 消息的标记* multiple 是否为批量确认*/ConfirmCallback ackCallback = (deliveryTag,multiple) ->{System.out.println("确认的消息:" + deliveryTag);};ConfirmCallback nackCallback = (deliveryTag,multiple) ->{System.out.println("未确认的消息:" + deliveryTag);};//消息监听器  监听那些消息成功了,那些消息失败了channel.addConfirmListener(ackCallback,nackCallback);//批量发送消息for (int i = 0; i < 1000; i++) {String message = i+"消息";channel.basicPublish("", QueueName,null,message.getBytes());}long end = System.currentTimeMillis();System.out.println("异步确认发布: " + MESSAGE_MAX + "条消息共耗时:"+ (end-start)+ "ms");}
}
  • 单独发布消息

同步等待确认,简单,但吞吐量非常有限。

  • 批量发布消息

批量同步等待确认,简单,合理的吞吐量,一旦出现问题但很难推断出是哪条消息出现了问题。

  • 异步处理

最佳性能和资源使用,在出现错误的情况下可以很好地控制,但是实现起来稍微难些

注意:应答和发布的区别:

应答功能属于消费者,消费者消费完消息后告诉RabbitMQ已经消费成功

发布功能属于生产者,生产者生产消息到RabbitMQ,RabbitMQ需要告诉生产者已经收到消息

这篇关于RabbitMQ消息应答与发布的文章就介绍到这儿,希望我们推荐的文章对编程师们有所帮助!



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

相关文章

RabbitMQ工作模式中的RPC通信模式详解

《RabbitMQ工作模式中的RPC通信模式详解》在RabbitMQ中,RPC模式通过消息队列实现远程调用功能,这篇文章给大家介绍RabbitMQ工作模式之RPC通信模式,感兴趣的朋友一起看看吧... 目录RPC通信模式概述工作流程代码案例引入依赖常量类编写客户端代码编写服务端代码RPC通信模式概述在R

C++ RabbitMq消息队列组件详解

《C++RabbitMq消息队列组件详解》:本文主要介绍C++RabbitMq消息队列组件的相关知识,本文给大家介绍的非常详细,对大家的学习或工作具有一定的参考借鉴价值,需要的朋友参考下吧... 目录1. RabbitMq介绍2. 安装RabbitMQ3. 安装 RabbitMQ 的 C++客户端库4. A

SpringCloud整合MQ实现消息总线服务方式

《SpringCloud整合MQ实现消息总线服务方式》:本文主要介绍SpringCloud整合MQ实现消息总线服务方式,具有很好的参考价值,希望对大家有所帮助,如有错误或未考虑完全的地方,望不吝赐... 目录一、背景介绍二、方案实践三、升级版总结一、背景介绍每当修改配置文件内容,如果需要客户端也同步更新,

macOS Sequoia 15.5 发布: 改进邮件和屏幕使用时间功能

《macOSSequoia15.5发布:改进邮件和屏幕使用时间功能》经过常规Beta测试后,新的macOSSequoia15.5现已公开发布,但重要的新功能将被保留到WWDC和... MACOS Sequoia 15.5 正式发布!本次更新为 Mac 用户带来了一系列功能强化、错误修复和安全性提升,进一步增

Maven 依赖发布与仓库治理的过程解析

《Maven依赖发布与仓库治理的过程解析》:本文主要介绍Maven依赖发布与仓库治理的过程解析,本文通过实例代码给大家介绍的非常详细,对大家的学习或工作具有一定的参考借鉴价值,需要的朋友参考下... 目录Maven 依赖发布与仓库治理引言第一章:distributionManagement配置的工程化实践1

一文带你搞懂Redis Stream的6种消息处理模式

《一文带你搞懂RedisStream的6种消息处理模式》Redis5.0版本引入的Stream数据类型,为Redis生态带来了强大而灵活的消息队列功能,本文将为大家详细介绍RedisStream的6... 目录1. 简单消费模式(Simple Consumption)基本概念核心命令实现示例使用场景优缺点2

Redis消息队列实现异步秒杀功能

《Redis消息队列实现异步秒杀功能》在高并发场景下,为了提高秒杀业务的性能,可将部分工作交给Redis处理,并通过异步方式执行,Redis提供了多种数据结构来实现消息队列,总结三种,本文详细介绍Re... 目录1 Redis消息队列1.1 List 结构1.2 Pub/Sub 模式1.3 Stream 结

使用Python构建一个Hexo博客发布工具

《使用Python构建一个Hexo博客发布工具》虽然Hexo的命令行工具非常强大,但对于日常的博客撰写和发布过程,我总觉得缺少一个直观的图形界面来简化操作,下面我们就来看看如何使用Python构建一个... 目录引言Hexo博客系统简介设计需求技术选择代码实现主框架界面设计核心功能实现1. 发布文章2. 加

售价599元起! 华为路由器X1/Pro发布 配置与区别一览

《售价599元起!华为路由器X1/Pro发布配置与区别一览》华为路由器X1/Pro发布,有朋友留言问华为路由X1和X1Pro怎么选择,关于这个问题,本期图文将对这二款路由器做了期参数对比,大家看... 华为路由 X1 系列已经正式发布并开启预售,将在 4 月 25 日 10:08 正式开售,两款产品分别为华

利用Python快速搭建Markdown笔记发布系统

《利用Python快速搭建Markdown笔记发布系统》这篇文章主要为大家详细介绍了使用Python生态的成熟工具,在30分钟内搭建一个支持Markdown渲染、分类标签、全文搜索的私有化知识发布系统... 目录引言:为什么要自建知识博客一、技术选型:极简主义开发栈二、系统架构设计三、核心代码实现(分步解析