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

相关文章

WordPress网创自动采集并发布插件

网创教程:WordPress插件网创自动采集并发布 阅读更新:随机添加文章的阅读数量,购买数量,喜欢数量。 使用插件注意事项 如果遇到404错误,请先检查并调整网站的伪静态设置,这是最常见的问题。需要定制化服务,请随时联系我。 本次更新内容 我们进行了多项更新和优化,主要包括: 界面设置:用户现在可以更便捷地设置文章分类和发布金额。代码优化:改进了采集和发布代码,提高了插件的稳定

AI赋能天气:微软研究院发布首个大规模大气基础模型Aurora

编者按:气候变化日益加剧,高温、洪水、干旱,频率和强度不断增加的全球极端天气给整个人类社会都带来了难以估计的影响。这给现有的天气预测模型提出了更高的要求——这些模型要更准确地预测极端天气变化,为政府、企业和公众提供更可靠的信息,以便做出及时的准备和响应。为了应对这一挑战,微软研究院开发了首个大规模大气基础模型 Aurora,其超高的预测准确率、效率及计算速度,实现了目前最先进天气预测系统性能的显著

常用MQ消息中间件Kafka、ZeroMQ和RabbitMQ对比及RabbitMQ详解

1、概述   在现代的分布式系统和实时数据处理领域,消息中间件扮演着关键的角色,用于解决应用程序之间的通信和数据传递的挑战。在众多的消息中间件解决方案中,Kafka、ZeroMQ和RabbitMQ 是备受关注和广泛应用的代表性系统。它们各自具有独特的特点和优势,适用于不同的应用场景和需求。   Kafka 是一个高性能、可扩展的分布式消息队列系统,被设计用于处理大规模的数据流和实时数据传输。它

消息认证码解析

1. 什么是消息认证码         消息认证码(Message Authentication Code)是一种确认完整性并进行认证的技术,取三个单词的首字母,简称为MAC。         消息认证码的输入包括任意长度的消息和一个发送者与接收者之间共享的密钥,它可以输出固定长度的数据,这个数据称为MAC值。         根据任意长度的消息输出固定长度的数据,这一点和单向散列函数很类似

RabbitMQ实践——临时队列

临时队列是一种自动删除队列。当这个队列被创建后,如果没有消费者监听,则会一直存在,还可以不断向其发布消息。但是一旦的消费者开始监听,然后断开监听后,它就会被自动删除。 新建自动删除队列 我们创建一个名字叫queue.auto.delete的临时队列 绑定 我们直接使用默认交换器,所以不用创建新的交换器,也不用建立绑定关系。 实验 发布消息 我们在后台管理页面的默认交换器下向这个队列

物联网系统运维——移动电商应用发布,Tomcat应用服务器,实验CentOS 7安装JDK与Tomcat,配置Tomcat Web管理界面

一.Tomcat应用服务器 1.Tomcat介绍 Tomcat是- -个免费的开源的Ser Ivet容器,它是Apache基金会的Jakarta 项目中的一个核心项目,由Apache, Sun和其他一 些公司及个人共同开发而成。Tomcat是一一个小型的轻量级应用服务器,在中小型系统和并发访问用户不是很多的场合下被普遍使用,是开发和调试JSP程序的首选。 在Tomcat中,应用程序的成部署很简

rocketmq问题汇总-如何将特定消息发送至特定queue,消费者从特定queue消费

业务描述 由于业务需要这样一种场景,将消息按照id(业务id)尾号发送到对应的queue中,并启动10个消费者(单jvm,10个消费者组),从对应的queue中集群消费,如下图1所示(假设有两个broker组成的集群):  producer如何实现 producer只需发送消息时调用如下方法即可 /*** 发送有序消息** @param messageMap 消息数据* @param

Spring 集成 RabbitMQ 与其概念,消息持久化,ACK机制

目录 RabbitMQ 概念exchange交换机机制 什么是交换机binding?Direct Exchange交换机Topic Exchange交换机Fanout Exchange交换机Header Exchange交换机RabbitMQ 的 Hello - Demo(springboot实现)RabbitMQ 的 Hello Demo(spring xml实现)RabbitMQ 在生产环境

Spring boot+RabbitMQ环境

消息队列在目前分布式系统下具备非常重要的地位,如下的场景是比较适合消息队列的: 跨系统的调用,异步性质的调用最佳。高并发问题,利用队列串行特点。订阅模式,数据被未知数量的消费者订阅,比如某种数据的变更会影响多个系统的数据,订单数据就是比较好理解的。 之前有一个场景是商品数据在修改后需要推送到elasticsearch中,由于修改产品的并发量以及数据量均不大,所以对于消息未做持久化,而且为了快速

SpringBoot中如何监听两个不同源的RabbitMQ消息队列

spring-boot如何配置监听两个不同的RabbitMQ 由于前段时间在公司开发过程中碰到了一个问题,需要同时监听两个不同的rabbitMq,但是之前没有同时监听两个RabbitMq的情况,因此在同事的帮助下,成功实现了监听多个MQ。下面我给大家一步一步讲解下,也为自己做个笔记; 详细步骤: 1. application.properties 文件配置: u.rabbitmq.ad