RabbitMQ工作队列之公平分发消息与消息应答(ACK)

上篇文章中,我们讲了工作队列轮询的分发模式,该模式无论有多少个消费者,不管每个消费者处理消息的效率,都会将所有消息平均的分发给每一个消费者,也就是说,大家最后各自消费的消息数量都是一样多的。由此也就引发我们今天要介绍的公平分发模式。

消息应答(ACK)

消息丢失

我们之前的所有代码,如果消息队列将消息分发给消费者,那么就会从队列中删除,如果在我们处理任务的过程中,处理失败或者服务器宕机,那么这条消息肯定得不到执行,就会出现丢失。

我们所设想的如果任务在处理的过程中,如果服务器宕机等原因造成消息未被正常消费,那么必须分发给其他的消费者再次进行消费,这样及时服务器宕机也不会丢失任何的消息了。

ACK

所以ACK,就是消息应答机制,我们之前写的代码都是开启了自动应答,所以如果我们的消息没被正常消费,就会丢失。

要想确保消息不丢失,就必须将ACK自动应答关闭掉,在我们处理消息的流程中,如果消息正常被处理,那么最后进行手动应答,告诉队列我们正常消费了消息。

超时

RabbitMQ它是没有我们平常所见到的超时时间限制的,只要当消费者服务宕机,消息才会被重新分发,哪怕处理这条消息需要花费很长的时间。

公平分发模式

缺陷

我们提供多个消费者,目的就是为了提高系统的性能,提升系统处理任务的速度,如果将消息平均的分发给每个消费者,那么处理消息快的服务是不是会空闲下来,而处理慢的服务可能会阻塞等待处理,这样的场景是我们不愿意看到的。所以有了今天要说的分发模式,公平分发

能者多劳

所谓的公平分发,其实用能者多劳描述更为贴切,根据名字就可以知道,谁有能力处理更多的任务,那么就交给谁处理,防止消息的挤压。

那么想要实现公平分发,那么必须要将自动应答改为手动应答。这是公平分发的前提。

代理

消息生产者

1public class Send { 2 3 public static final String QUEUE_NAME = "test_word_queue"; 4 5 public static void main(String[] args) throws IOException, TimeoutException, InterruptedException { 6 7 // 获取连接 8 Connection connection = MQConnectUtil.getConnection(); 9 10 // 创建通道 11 Channel channel = connection.createChannel(); 12 13 // 声明队列 14 channel.queueDeclare(QUEUE_NAME, false, false, false, null); 15 16 for (int i = 0; i < 10; i++) { 17 18 String msg = "消息:" + i; 19 20 channel.basicPublish("", QUEUE_NAME, null, msg.getBytes()); 21 22 Thread.sleep(i * 20); 23 24 System.out.println(msg); 25 } 26 27 channel.close(); 28 connection.close(); 29 } 30} 31

消费者1

我们在消费者中设置了channel.basicQos(1);这样一个参数,这个意思就是表示,此消费者每次最多只接收一条消息进行处理,只有将消息处理结束,手动应答之后,下一条消息才会被分发进来。

1public class Consumer1 { 2 3 public static final String QUEUE_NAME = "test_word_queue"; 4 5 public static void main(String[] args) throws Exception { 6 7 // 获取连接 8 Connection connection = MQConnectUtil.getConnection(); 9 10 // 创建频道 11 Channel channel = connection.createChannel(); 12 13 // 一次仅接受一条未经确认的消息 14 channel.basicQos(1); 15 16 // 队列声明 17 channel.queueDeclare(QUEUE_NAME, false, false, false, null); 18 19 // 定义消费者 20 DefaultConsumer consumer = new DefaultConsumer(channel) { 21 @SneakyThrows 22 @Override 23 public void handleDelivery(String consumerTag, Envelope envelope, AMQP.BasicProperties properties, byte[] body) { 24 String msg = new String(body, StandardCharsets.UTF_8); 25 26 System.out.println("消费者[1]-内容:" + msg); 27 28 Thread.sleep(2 * 1000); 29 30 // 手动回执消息 31 channel.basicAck(envelope.getDeliveryTag(), false); 32 } 33 }; 34 35 // 监听队列,将自动应答方式改为false,关闭自动应答机制 36 boolean autoAck = false; 37 channel.basicConsume(QUEUE_NAME, autoAck, consumer); 38 } 39} 40

消费者2

1public class Consumer2 { 2 3 public static final String QUEUE_NAME = "test_word_queue"; 4 5 public static void main(String[] args) throws Exception { 6 7 // 获取连接 8 Connection connection = MQConnectUtil.getConnection(); 9 10 // 创建频道 11 Channel channel = connection.createChannel(); 12 13 // 队列声明 14 channel.queueDeclare(QUEUE_NAME, false, false, false, null); 15 16 channel.basicQos(1); 17 18 // 定义消费者 19 DefaultConsumer consumer = new DefaultConsumer(channel) { 20 @SneakyThrows 21 @Override 22 public void handleDelivery(String consumerTag, Envelope envelope, AMQP.BasicProperties properties, byte[] body) { 23 String msg = new String(body, StandardCharsets.UTF_8); 24 25 System.out.println("消费者[2]-内容:" + msg); 26 27 Thread.sleep(1000); 28 29 // 手动回执消息 30 channel.basicAck(envelope.getDeliveryTag(), false); 31 } 32 }; 33 34 // 监听队列,需要将自动应答方式改为false 35 boolean autoAck = false; 36 channel.basicConsume(QUEUE_NAME, autoAck, consumer); 37 } 38} 39

消费结果

那么结果就会像我们之前预想的那样,由于消费者2消费消息花费的时间比消费者1更少,所以消费者2处理的消息的数量要比消费者1处理的消息的数量要多。这里我就不贴图了,大家可以敲代码进行尝试。


今天的文章到这里就结束了,下篇呢,会给介绍介绍另外一种模式,发布订阅模式

<p style="text-align:center;font-weight:bold;color:#0e88eb;font-size:20px">日拱一卒,功不唐捐</p> <p style="text-align:center;font-weight:bold;color:#773098;font-size:16px">更多内容请关注:</p>

点赞
收藏

评论区

加载中...

相关推荐

MySQL:[Err] 1292 - Incorrect datetime value: ‘0000-00-00 00:00:00‘ for column ‘CREATE_TIME‘ at row 1

文章目录问题用navicat导入数据时,报错:原因这是因为当前的MySQL不支持datetime为0的情况。解决修改sql\mode:sql\mode:SQLMode定义了MySQL应支持的SQL语法、数据校验等,这样可以更容易地在不同的环境中使用MySQL。全局s

Oracle 分组与拼接字符串同时使用

SELECTT.,ROWNUMIDFROM(SELECTT.EMPLID,T.NAME,T.BU,T.REALDEPART,T.FORMATDATE,SUM(T.S0)S0,MAX(UPDATETIME)CREATETIME,LISTAGG(TOCHAR(

MySQL部分从库上面因为大量的临时表tmp_table造成慢查询

背景描述Time:20190124T00:08:14.70572408:00User@Host:@Id:Schema:sentrymetaLast_errno:0Killed:0Query_time:0.315758Lock_

皕杰报表之UUID

​在我们用皕杰报表工具设计填报报表时,如何在新增行里自动增加id呢?能新增整数排序id吗?目前可以在新增行里自动增加id,但只能用uuid函数增加UUID编码,不能新增整数排序id。uuid函数说明:获取一个UUID,可以在填报表中用来创建数据ID语法:uuid()或uuid(sep)参数说明:sep布尔值,生成的uuid中是否包含分隔符'',缺省为

手写Java HashMap源码

HashMap的使用教程HashMap的使用教程HashMap的使用教程HashMap的使用教程HashMap的使用教程22

RabbitMQ如何高效的消费消息

在上篇介绍了如何简单的发送一个消息队列之后,我们本篇来看下RabbitMQ的另外一种模式,工作队列。什么是工作队列我们上篇文章说的是,一个生产者生产了消息被一个消费者消费了,如下图!(https://usergoldcdn.xitu.io/2020/5/15/1721768c1b303014?w1824&h55