rocketmq顺序消息源码解析

发布时间:2026/8/25 11:12:46
rocketmq顺序消息源码解析 支持消费者按照发送消息的先后顺序获取消息从而实现有业务场景中的顺序处理。顺序消息的顺序关系通过消息组messageGroup判定和识别发送顺序消息需要为每条消息设置归属的消息组相同消息组的多条消息之间遵循先进先出的顺序关系。如何保证顺序消息局部有序生产者生产消息时相关消息都放在1个队列消费者使用单线程进行消费全局有序消息队列为topic设置一个消息队列生产者使用1个生产者单线程发送数据消费者使用单线程进行消费代码示例底层处理流程锁续期、失效分布式锁队列锁ConcurrentMapString/* group */, ConcurrentHashMapMessageQueue, LockEntry mqLockTable 最大存活时间默认60秒。如果当前时间和锁的最后更新时间大于60秒则锁过期。在锁续期或者重平衡时会判断锁是否过期是的话则覆盖新的数据clientId)重新加锁交互流程顺序消息重试顺序消息优先本地重试暂停当前队列消费等一段时间后重试失败后才发往重试队列。并发消费直接发往重试队列消息发送与普通消息的区别在于普通消息发送时从所有broker的队列集合中 轮询选择一个队列而顺序队列可以提供用户自定义消息队列选择器从NameServer 分配的顺序 broker集合中选择一个队列1.rocketMq提供了一个 MessageQueueSelector 消息队列选择器用于自定义选择队列的逻辑2.触发队列选择器顺序消息消费源码定时任务对消息队列加锁消费者在启动时判断是否顺序消费是则调用ConsumeMessageOrderlyService的start方法定期对该消费者负责的消息队列进行加锁队列与消费者绑定关系存储在本地——processQueueTablebroker端加锁消息队列负载在消费者重新分配消息队列过程中更新当前消费者持有的消息队列1新增队列加锁2移除队列消息拉取提交消息到线程池消息消费消费重试