RocketMQ 顺序消费:从队列路由到线程绑定的完整链路
一句话结论(30s)
RocketMQ 顺序消费的本质是「把并发串行化」,因为顺序保证不是某一个组件单独完成的,而是「生产取模路由 → Broker FIFO 存储 → 消费端分布式锁 + 单线程绑定」三环接力。关键设计是 MessageListenerOrderly 用队列锁把同一个 Queue 串行消费,并用 SUSPEND_CURRENT_QUEUE_A_MOMENT 在失败时暂停整队而不是跳过。核心权衡是顺序 = 串行化:要顺序就牺牲并行度,所以只对必须有序的业务(订单状态机、binlog 同步)做局部有序。
核心原理(2min)
主流程分三环:生产者用 orderId % queueSize 取模路由,保证同一业务标识始终落到同一个 MessageQueue;Broker 端每个 Queue 是 FIFO,CommitLog 顺序写天然保证「发送序 = 存储序 = 拉取序」;消费端最关键,MessageListenerOrderly 先向 Broker 申请该 Queue 的分布式锁(同刻只有一个实例的一个线程消费),再同步逐条处理(前一条返回 SUCCESS 才处理下一条)。关键机制有二:一是 Rebalance 时旧 Consumer 释放锁、新 Consumer 锁抢占、锁超时释放(默认 30 秒),保证任意时刻一个 Queue 最多一个线程消费;二是失败时返回 SUSPEND_CURRENT_QUEUE_A_MOMENT,暂停当前队列、等待后重试同一批消息、累计 16 次进 DLQ,代价是一条消息卡住会阻塞整队。
底层深入(5-10min)
先搞清楚:什么是有序?
“消息有序”有两个层次:
- 全局有序:整个 Topic 下所有消息按发送顺序消费。RocketMQ 通过”一个 Topic 只配一个队列”实现。代价:完全丧失并行度。
- 局部有序:同一个业务标识(如订单ID)的消息按顺序消费,不同标识之间可以并行。这是生产主流。
举个例子:
订单 1001: 创建 → 支付 → 发货 → 完成 (必须有序)
订单 1002: 创建 → 取消 (必须有序)
但 1001 和 1002 之间不需要有序,可以并行处理
> **🤔 思考穿插**:为什么生产上大多只要「局部有序」而不要「全局有序」?—— 因为全局有序要求整个 Topic 只配一个队列,完全丧失并行度、吞吐极低;而业务里真正需要严格顺序的只是同一实体(如同一订单)的状态流转,不同实体之间可以并行。所以「局部有序」是顺序和吞吐之间的最优平衡。
RocketMQ 如何实现局部有序?三步走。
第一步:生产者路由到同一个队列
RocketMQ 中,一个 Topic 有多个 MessageQueue。生产者发送消息时,需要确保同一业务标识的消息始终发到同一个 Queue。
producer.send(msg, new MessageQueueSelector() {
@Override
public MessageQueue select(List<MessageQueue> mqs, Message msg, Object arg) {
long orderId = (Long) arg; // 业务标识
int index = (int) (orderId % mqs.size()); // 取模路由
return mqs.get(index);
}
}, orderId); // arg 传入订单ID
关键:路由算法必须稳定。 用 orderId % queueSize 取模。如果队列数变了(扩容/缩容),路由结果就会变,同一 orderId 的消息会落入不同队列,顺序被打乱。
🤔 思考穿插:为什么顺序消费的路由算法必须「稳定」?—— 因为顺序的前提是「同一业务标识始终落到同一 Queue」,一旦队列数变了(扩容/缩容),取模结果就变、同一 orderId 的消息散到不同 Queue,顺序就断了。所以路由算法和队列数的稳定,是顺序消费的第一道闸。
第二步:Broker 端队列 FIFO
RocketMQ Broker 的每个 MessageQueue 是一个 FIFO 队列,写消息时 Append 到队尾,消费时从队头取。这个简单的设计保证了单个队列内消息的发送顺序 = 存储顺序 = 消费拉取顺序。
Queue-0: [创建] → [支付] → [发货] → [完成] ← 有序
Queue-1: [创建] → [取消] ← 有序
Broker 无需做任何额外的事情来实现顺序——CommitLog 的顺序写天然就是 FIFO。
第三步:消费端线程绑定
这是最关键的一步。 即使 Broker 保证了队列内有序,如果 Consumer 用多线程并发消费同一个队列,顺序仍然会乱:
Queue-0 拉取了:[创建, 支付, 发货]
线程1: 处理 "创建" (耗时 10ms)
线程2: 处理 "支付" (耗时 5ms) ← 先处理完,先提交 offset
线程3: 处理 "发货" (耗时 20ms)
结果:支付先提交,但创建还在处理中 —— 如果此时 Rebalance,创建丢失!
RocketMQ 的解决方案是 MessageListenerOrderly:
consumer.registerMessageListener(new MessageListenerOrderly() {
@Override
public ConsumeOrderlyStatus consumeMessage(List<MessageExt> msgs,
ConsumeOrderlyContext context) {
for (MessageExt msg : msgs) {
processOrder(msg); // 按顺序处理
}
return ConsumeOrderlyStatus.SUCCESS;
}
});
MessageListenerOrderly 做了两件事:
- 单队列加锁:向 Broker 申请该 MessageQueue 的分布式锁。同一时刻,同一个 Queue 只能被一个 Consumer 实例的一个线程消费。
- 同步消费:拉取的一批消息,按顺序逐个处理。前一条返回 SUCCESS 后,才处理下一条。
🤔 思考穿插:为什么消费端必须「单队列加锁 + 同步消费」,光靠 Broker FIFO 不够?—— 因为 Broker 只保证队列内有序,若 Consumer 用多线程并发消费同一队列,处理快慢不同会导致后到的先提交 offset,一旦 Rebalance 就会丢消息、乱序。所以要在消费端再串行化一次,才能守住「顺序」。
Queue-0 加锁 ──▶ 线程A: [创建] → SUCCESS → [支付] → SUCCESS → [发货] → SUCCESS
Queue-1 加锁 ──▶ 线程B: [创建] → SUCCESS → [取消] → SUCCESS
源码印证:队列锁与顺序取消息
上面说的「单队列加锁 + 同步消费」具体落在三处源码里:MessageQueueLock 负责队列锁、ProcessQueue 负责按 offset 顺序取消息、ConsumeMessageOrderlyService 负责加锁后的消费循环。
第一处:队列锁(MessageQueueLock.java)
public class MessageQueueLock {
private ConcurrentMap<MessageQueue, ConcurrentMap<Integer, Object>> mqLockTable =
new ConcurrentHashMap<>(32);
public Object fetchLockObject(final MessageQueue mq) {
return fetchLockObject(mq, -1);
}
public Object fetchLockObject(final MessageQueue mq, final int shardingKeyIndex) {
ConcurrentMap<Integer, Object> objMap = this.mqLockTable.get(mq);
if (null == objMap) {
objMap = new ConcurrentHashMap<>(32);
ConcurrentMap<Integer, Object> prevObjMap = this.mqLockTable.putIfAbsent(mq, objMap);
if (prevObjMap != null) {
objMap = prevObjMap;
}
}
Object lock = objMap.get(shardingKeyIndex);
if (null == lock) {
lock = new Object();
Object prevLock = objMap.putIfAbsent(shardingKeyIndex, lock);
if (prevLock != null) {
lock = prevLock;
}
}
return lock;
}
}
这个 mqLockTable 就是「队列锁」:它以 MessageQueue 为 key、映射到一个独立的锁对象 Object,putIfAbsent 保证并发场景下所有线程拿到的是同一个锁实例。消费线程后续在 synchronized (lock) 上排队,同一时刻只有一个线程能进入该 Queue 的消费逻辑,这是 JVM 进程内的单队列串行保证。跨实例的互斥则交由 Broker 端的分布式锁(isLocked)兜底。
第二处:顺序取消息(ProcessQueue.java)
private final ReadWriteLock treeMapLock = new ReentrantReadWriteLock();
private final TreeMap<Long, MessageExt> msgTreeMap = new TreeMap<>();
/**
* A subset of msgTreeMap, will only be used when orderly consume
*/
private final TreeMap<Long, MessageExt> consumingMsgOrderlyTreeMap = new TreeMap<>();
public List<MessageExt> takeMessages(final int batchSize) {
List<MessageExt> result = new ArrayList<>(batchSize);
final long now = System.currentTimeMillis();
try {
this.treeMapLock.writeLock().lockInterruptibly();
this.lastConsumeTimestamp = now;
try {
if (!this.msgTreeMap.isEmpty()) {
for (int i = 0; i < batchSize; i++) {
Map.Entry<Long, MessageExt> entry = this.msgTreeMap.pollFirstEntry();
if (entry != null) {
result.add(entry.getValue());
consumingMsgOrderlyTreeMap.put(entry.getKey(), entry.getValue());
} else {
break;
}
}
}
if (result.isEmpty()) {
consuming = false;
}
} finally {
this.treeMapLock.writeLock().unlock();
}
} catch (InterruptedException e) {
log.error("take Messages exception", e);
}
return result;
}
msgTreeMap 是以 queueOffset 为 key 的 TreeMap,天然按 offset 升序排列,pollFirstEntry() 每次从最小 offset 开始取,从数据结构层面保证了「先入先出」的顺序。取出的消息从 msgTreeMap 移入 consumingMsgOrderlyTreeMap(顺序消费专用的「消费中」暂存区),只有前一批全部 commit() 成功后才清空;若返回 SUSPEND_CURRENT_QUEUE_A_MOMENT,则由 makeMessageToConsumeAgain 把消息放回 msgTreeMap 等待重试。
第三处:加锁后的消费循环(ConsumeMessageOrderlyService.java)
final Object objLock = messageQueueLock.fetchLockObject(this.messageQueue);
synchronized (objLock) {
if (MessageModel.BROADCASTING.equals(ConsumeMessageOrderlyService.this.defaultMQPushConsumerImpl.messageModel())
|| this.processQueue.isLocked() && !this.processQueue.isLockExpired()) {
final long beginTime = System.currentTimeMillis();
for (boolean continueConsume = true; continueConsume; ) {
if (this.processQueue.isDropped()) {
log.warn("the message queue not be able to consume, because it's dropped. {}", this.messageQueue);
break;
}
if (MessageModel.CLUSTERING.equals(ConsumeMessageOrderlyService.this.defaultMQPushConsumerImpl.messageModel())
&& !this.processQueue.isLocked()) {
log.warn("the message queue not locked, so consume later, {}", this.messageQueue);
ConsumeMessageOrderlyService.this.tryLockLaterAndReconsume(this.messageQueue, this.processQueue, 10);
break;
}
if (MessageModel.CLUSTERING.equals(ConsumeMessageOrderlyService.this.defaultMQPushConsumerImpl.messageModel())
&& this.processQueue.isLockExpired()) {
log.warn("the message queue lock expired, so consume later, {}", this.messageQueue);
ConsumeMessageOrderlyService.this.tryLockLaterAndReconsume(this.messageQueue, this.processQueue, 10);
break;
}
long interval = System.currentTimeMillis() - beginTime;
if (interval > MAX_TIME_CONSUME_CONTINUOUSLY) {
ConsumeMessageOrderlyService.this.submitConsumeRequestLater(processQueue, messageQueue, 10);
break;
}
final int consumeBatchSize =
ConsumeMessageOrderlyService.this.defaultMQPushConsumer.getConsumeMessageBatchMaxSize();
List<MessageExt> msgs = this.processQueue.takeMessages(consumeBatchSize);
defaultMQPushConsumerImpl.resetRetryAndNamespace(msgs, defaultMQPushConsumer.getConsumerGroup());
if (!msgs.isEmpty()) {
final ConsumeOrderlyContext context = new ConsumeOrderlyContext(this.messageQueue);
// ... 执行消费 Hook、耗时统计等
status = messageListener.consumeMessage(Collections.unmodifiableList(msgs), context);
}
// ... 统计 consumeRT、执行 After Hook
continueConsume = ConsumeMessageOrderlyService.this.processConsumeResult(msgs, status, context, this);
}
}
}
整个消费循环被 synchronized (objLock) 包住,拿到队列锁后才进入;循环内先做三层防御:isLocked 判断是否还持有 Broker 分布式锁、isLockExpired 判断锁是否超时、MAX_TIME_CONSUME_CONTINUOUSLY 限制单次连续消费时长(默认 60 秒)。takeMessages 顺序取出一批后,直接同步调用 messageListener.consumeMessage(...) 交给业务,前一批的返回状态经 processConsumeResult 决定是 commit() 提交还是 SUSPEND 挂起重试。
Rebalance 对顺序消费的影响
Rebalance 是 RocketMQ 的消费者重平衡机制——当消费者上下线、扩缩容时,重新分配 Queue 与 Consumer 的对应关系。
问题场景:
时刻 1: Consumer-A 持有 Queue-0 的锁,正在处理 [支付]
时刻 2: Consumer-B 上线,触发 Rebalance
时刻 3: Queue-0 被分配给 Consumer-B
时刻 4: Consumer-B 开始从 Queue-0 拉取消息
问题:Consumer-A 可能还没提交 [支付] 的 offset,
Consumer-B 拉到了 [支付] 之后的 [发货],
但 [支付] 还没处理完!
RocketMQ 的处理:
- Rebalance 时释放锁:Consumer 收到 Rebalance 通知后,释放所有持有的 Queue 锁,停止消费。
- 锁抢占保护:新的 Consumer 在拿到 Rebalance 结果后,需要向 Broker 申请 Queue 的分布式锁。如果旧 Consumer 还没释放锁,新 Consumer 拿不到锁,不能开始消费。
- 锁超时释放:如果旧 Consumer 挂了(锁没主动释放),Broker 的锁有超时机制(默认 30 秒),超时后自动释放。
Rebalance 流程:
1. Consumer-A 检测到 Consumer-B 加入
2. Consumer-A 停止消费、释放 Queue-0 锁、提交当前 offset
3. Rebalance 完成,Queue-0 分配给 Consumer-B
4. Consumer-B 申请 Queue-0 锁成功,开始消费
这保证了:任意时刻,一个 Queue 最多只能被一个线程消费。即使 Rebalance 期间也不会有两个 Consumer 同时消费同一个 Queue。 但代价是 Rebalance 期间该 Queue 暂停消费(通常是几百毫秒)。
并发消费 vs 顺序消费
| 维度 | MessageListenerConcurrently | MessageListenerOrderly |
|---|---|---|
| 线程模型 | 多线程并发处理同一 Queue | 单线程独占一个 Queue |
| 吞吐量 | 高(充分利用多线程) | 低(受限于单线程处理速度) |
| 顺序保证 | 无 | 队列内严格有序 |
| 失败重试 | 当前批立即重试 | 挂起当前 Queue,逐条重试,直到成功才继续 |
| 适用场景 | 日志、通知、统计 | 订单状态机、数据库 binlog 同步 |
顺序消费的并行度优化:
// 不要这样:所有订单共享 4 个 Queue
// 订单量一大,每个 Queue 的消费速度成为瓶颈
// 应该这样:增加 Queue 数量,提高并行度
// Topic 配置:writeQueueNums=16, readQueueNums=16
// 16 个 Queue → 最多 16 个线程并行消费
// 同一个 orderId 仍然路由到同一个 Queue
顺序消费的可靠性:SUSPEND_CURRENT_QUEUE_A_MOMENT
这是 MessageListenerOrderly 中最重要的设计细节。
public ConsumeOrderlyStatus consumeMessage(List<MessageExt> msgs,
ConsumeOrderlyContext context) {
for (MessageExt msg : msgs) {
try {
processOrder(msg);
} catch (Exception e) {
// 返回 SUSPEND,而不是直接失败
context.setSuspendCurrentQueueTimeMillis(3000);
return ConsumeOrderlyStatus.SUSPEND_CURRENT_QUEUE_A_MOMENT;
}
}
return ConsumeOrderlyStatus.SUCCESS;
}
当返回 SUSPEND_CURRENT_QUEUE_A_MOMENT 时:
- 当前 Queue 的消费暂停。不会跳过这条消息去处理下一条——因为必须保证顺序。
- 等待一段时间后重试同一批消息(由
setSuspendCurrentQueueTimeMillis控制)。 - 重试次数累计。达到
maxReconsumeTimes(默认 16 次)后,发往死信队列(DLQ)。
这就是顺序消费的代价:一条消息卡住,整个队列后面的所有消息都卡住。
🤔 思考穿插:为什么失败要返回 SUSPEND_CURRENT_QUEUE_A_MOMENT 暂停整队,而不是跳过这条继续?—— 因为顺序消费不允许跳过——后面的消息可能依赖这条的前序状态,跳过会让后续处理基于错误的上下文。所以只能「暂停 → 重试同一批 → 超限进 DLQ」,代价是一条卡住全队卡住。
生产实践建议
-
Queue 数量提前规划:顺序消费的并行度 = Queue 数量。Queue 太少,吞吐量上不去;Queue 太多,每个 Queue 消息太少,浪费内存和文件句柄。一般设为 Consumer 节点数 x 核数的 2-4 倍。
-
业务标识选择要均匀:用
orderId.hashCode() % queueSize比orderId % queueSize更好——如果 orderId 是连续自增的,后者会导致所有消息集中在某几个 Queue。 -
处理逻辑要有超时:顺序消费时一条卡住全队列卡住,每条消息的处理必须设置超时,超时后
SUSPEND并告警。 -
监控队列积压差异:如果某个 Queue 积压严重而其他 Queue 空闲,说明该 Queue 分配了”热 key”(如某个大商家),需要检查路由算法。
-
避免 Rebalance 抖动:Consumer 频繁上下线会导致 Queue 锁频繁切换。确保 Consumer 稳定运行,不要轻易重启。
总结
RocketMQ 的顺序消费实现,本质上是三个环节的接力:
生产者路由(取模/哈希) → Broker FIFO 存储 → Consumer 线程绑定 + 分布式锁
每一环都不能出错。生产者路由决定了消息进入哪个 Queue,Broker FIFO 保证了队列内有序,Consumer 的锁 + 单线程保证了消费的有序性。
核心心法:顺序 = 串行化。要顺序就没有完全并行,要完全并行就没有顺序。 关键是识别哪些消息需要顺序、哪些不需要,把需要顺序的放一个 Queue、不需要的放其他 Queue,在”必须有序”的最小范围内串行。
章末提问
追问 1:RocketMQ 如何保证消息的顺序消费?分几步?
结论先行:三步接力——生产者取模路由到同一 Queue、Broker FIFO 存储、消费端队列锁 + 单线程同步消费。
因为:生产端用 orderId % queueSize 保证同一业务标识进同一 Queue;Broker 的 CommitLog 顺序写天然保证队列内「发送序 = 存储序 = 拉取序」;消费端 MessageListenerOrderly 先申请队列锁再单线程逐条处理,前一条 SUCCESS 才处理下一条。三环缺一不可。
追问 2:顺序消费和并发消费的本质区别是什么?各有什么代价?
结论先行:本质区别是「并发 vs 串行」——顺序消费单线程独占 Queue、保证队列内严格有序,但吞吐低;并发消费多线程并行、吞吐高,但不保证顺序。
因为:顺序消费用分布式锁 + 同步逐条消费把并发串行化,失败还要 SUSPEND_CURRENT_QUEUE_A_MOMENT 暂停整队,一条卡住阻塞整队;并发消费失败只重试当前批。所以「顺序 = 串行化 = 牺牲并行度」,只应对必须有序的业务(订单状态机、binlog 同步)做局部有序。
追问 3:Rebalance 期间,RocketMQ 如何避免两个 Consumer 同时消费同一个 Queue?
结论先行:靠三件事——Rebalance 时旧 Consumer 释放锁、新 Consumer 锁抢占、Broker 端锁超时(默认 30 秒)自动释放。
因为:新 Consumer 拿到 Rebalance 结果后要向 Broker 申请该 Queue 的分布式锁,旧 Consumer 未释放就拿不到锁、不能消费;若旧 Consumer 挂掉没主动释放,锁超时后自动释放。三管齐下保证任意时刻一个 Queue 最多一个线程消费,代价是 Rebalance 期间该 Queue 短暂暂停(通常几百毫秒)。