云计算百科
云计算领域专业知识百科平台

RocketMQ顺序消费的加锁机制

有开发者在使用RocketMQ的MessageListenerOrderly做顺序消费时,发现日志里同一个线程在同一毫秒内连续打印了多条消息:

[ConsumeMessageThread_10] msg is :Hello RocketMQ 75
[ConsumeMessageThread_10] msg is :Hello RocketMQ 79
[ConsumeMessageThread_10] msg is :Hello RocketMQ 83
[ConsumeMessageThread_10] msg is :Hello RocketMQ 87
[ConsumeMessageThread_10] msg is :Hello RocketMQ 91
[ConsumeMessageThread_10] msg is :Hello RocketMQ 95
[ConsumeMessageThread_10] msg is :Hello RocketMQ 99

消息编号75、79、83…99,递增顺序没问题,但同一个线程一口气消费了7条,看起来并不像「取一条、锁住、消费完、解锁、再取下一条」的模式。提问者怀疑这不算真正的顺序消费,觉得客户端在没有commit的情况下,会继续从队列拿数据放到线程池里并发消费。

从RocketMQ 4.9.8的源码来看,顺序消费的保证是真实的。同一个线程连续消费多条消息,不是并发,而是while循环在同一把锁内逐条处理。下面从源码角度尝试梳理一下这个加锁机制。

顺序消费的两层保护

顺序消费要解决一个问题:同一个队列的消息必须按顺序一条一条处理,不能有两个线程同时消费同一个队列。RocketMQ在客户端用了两层机制来保证这一点。在这里插入图片描述

第一层是ProcessQueue里的consuming标志。Pull线程从Broker拉回消息后,调用ProcessQueue的putMessage方法把消息存入内部的TreeMap。putMessage存完消息后会检查:如果当前没有消费任务在跑(consuming为false),就把consuming置为true,返回一个标记让上层提交一个新的消费任务到线程池。如果consuming已经是true,说明这个队列已经有消费任务在跑了,putMessage返回false,上层就不会重复提交。

这一层要达到的目的是:

同一个队列,同一时刻只有一个消费任务在跑。后续Pull拉到的新消息只是放进TreeMap里等着,不会触发新的消费任务提交。

第二层是MessageQueueLock为每个消息队列分配的独立锁。ConsumeMessageOrderlyService内部维护了一个MessageQueueLock,它用一个ConcurrentHashMap为每个消息队列分配一个独立的Object锁。消费任务执行时,第一步就是拿到这个队列对应的锁对象,然后进入synchronized块。

这一层是兜底。即使因为某些极端情况导致两个消费任务同时存在,它们也会在synchronized上排队,同一时刻只有一个线程能进入消费逻辑。

while循环:锁住队列连续消费

进入synchronized块之后,消费逻辑不是消费完一条消息就退出释放锁。ConsumeRequest的run方法里是一个while循环,每轮循环做三件事:

1.从ProcessQueue按偏移量顺序取出一批消息。

2.调用注册的监听器执行业务逻辑。

3.根据返回值决定是继续还是停止。

ProcessQueue内部用TreeMap存储消息,以消息偏移量为键,天然有序。取消息时从TreeMap头部依次取出最小的几个偏移量,保证了消费顺序。默认每次只取1条(由consumeMessageBatchMaxSize控制,默认值为1)。

如果队列里还有消息,continueConsume保持为true,循环继续。如果队列为空,consuming标志重置为false,循环退出,synchronized锁释放。

这就是提问者日志里看到的现象:同一个线程ConsumeMessageThread_10,在同一毫秒内连续消费了消息75、79、83…99。这些消息是在同一个线程、同一把synchronized锁内、按偏移量递增顺序逐条处理的。消费完75不释放锁,直接进入下一轮循环消费79,再下一轮消费83,依此类推,直到队列空了才退出循环释放锁。

容易被误解的一点是:被提交到线程池的不是单条消息,而是整个队列的消费任务。一个消费任务内部用while循环逐条消费这个队列的消息。提问者的日志里看到同一个线程连续消费多条消息,恰恰说明顺序消费在正常工作。

顺序消费锁的粒度是队列,不是消息。 锁住队列后,同一个线程在循环里连续消费多条,中间不释放锁。这和提问者想象的「每条消息单独加锁解锁」在结果上是一样的,都保证了顺序,只是锁的持有周期不同。

autoCommit(false)为什么不阻止循环继续

提问者的测试代码里做了两件事:设了context.setAutoCommit(false),以及返回ConsumeOrderlyStatus.SUCCESS。提问者期望「没有commit就不继续消费」,但实际行为是循环继续消费下一条。

autoCommit控制的是消费位点是否提交给Broker,不影响消费循环是否继续。

看ConsumeMessageOrderlyService的processConsumeResult方法,当autoCommit为false且状态为SUCCESS时,只做了TPS统计,没有调用commit(),continueConsume保持默认值true。循环照常继续取下一条消息。而autoCommit为true且状态为SUCCESS时,会调用ProcessQueue的commit()方法提交偏移量。两者的区别仅在于偏移量是否上报给Broker。

返回状态autoCommit=trueautoCommit=false
SUCCESS 自动提交偏移量,继续下一条 不提交偏移量,继续下一条
COMMIT 提交偏移量,继续下一条 提交偏移量,继续下一条
ROLLBACK 视为SUCCESS处理,打警告日志 消息回退到队列,暂停当前队列

返回SUCCESS时,不管autoCommit怎么设,循环都会继续取下一条消息。autoCommit只决定偏移量是否上报。

另外,autoCommit为false加上返回SUCCESS,消息被消费了但偏移量永远不会提交给Broker。消费者重启后,Broker不知道这些消息已经被消费过,会重新投递,导致重复消费。除非有明确的事务需求(比如需要根据消费结果决定commit还是rollback),不建议关闭autoCommit。

如果确实需要使用autoCommit=false的场景,正确做法是:业务处理成功后返回ConsumeOrderlyStatus.COMMIT而不是SUCCESS,这样偏移量才会被提交。

回答提问者的问题

提问者想要的是「拿出一条队列数据,锁住队列,消费完消息再解锁」的模式。

从源码来看,RocketMQ的原生客户端已经支持这个模式,只是实现方式和提问者的预想有出入。

提问者预想的是消息级别的加锁:取一条消息,加锁,消费完,解锁,再取下一条。每条消息都经历一次加锁和解锁。

实际的实现是队列级别的加锁:拿到队列锁之后,在一个while循环里连续取消息消费,直到队列空了才释放锁。整个过程中锁一直持有,不会在消息之间释放。在这里插入图片描述

两种方式都保证了消息按顺序消费,区别只在于锁的持有周期。RocketMQ选择队列级别的锁是出于性能考虑,避免每条消息都经历加锁解锁的开销。在消息量大的场景下,频繁加锁解锁的开销不可忽略。

提问者不需要自己写代码。他看到的日志行为就是顺序消费在正常工作:同一个线程按偏移量递增顺序消费,没有其他线程能插队。

小结

顺序消费的保证来自两层机制协同工作:ProcessQueue的consuming标志防止重复提交消费任务,synchronized锁保证同一队列同一时刻只有一个线程在消费。消费线程拿到锁后进入while循环,按偏移量递增顺序逐条处理消息,队列空了才释放锁。

autoCommit这个参数容易让人产生误解,以为它能控制消费是否继续。实际上它只管偏移量是否上报给Broker,和消费循环的继续与否没有关系。在需要用autoCommit=false的场景下,要用COMMIT而不是SUCCESS来触发偏移量提交。

实际项目中,顺序消费最容易出问题的地方往往不是客户端的加锁机制,而是消息是否被正确路由到了同一个队列。如果业务上需要保证顺序的消息分散在不同队列,客户端再怎么加锁也无济于事。发送端用同一个hash key把相关消息路由到同一个队列,是顺序消费的前提条件。

赞(0)
未经允许不得转载:网硕互联帮助中心 » RocketMQ顺序消费的加锁机制
分享到: 更多 (0)

评论 抢沙发

评论前必须登录!