Skip to content
Go back

RocketMQ 顺序消费:从队列路由到线程绑定的完整链路

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)

先搞清楚:什么是有序?

“消息有序”有两个层次:

举个例子:

订单 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 做了两件事:

  1. 单队列加锁:向 Broker 申请该 MessageQueue 的分布式锁。同一时刻,同一个 Queue 只能被一个 Consumer 实例的一个线程消费。
  2. 同步消费:拉取的一批消息,按顺序逐个处理。前一条返回 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、映射到一个独立的锁对象 ObjectputIfAbsent 保证并发场景下所有线程拿到的是同一个锁实例。消费线程后续在 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 的处理:

  1. Rebalance 时释放锁:Consumer 收到 Rebalance 通知后,释放所有持有的 Queue 锁,停止消费。
  2. 锁抢占保护:新的 Consumer 在拿到 Rebalance 结果后,需要向 Broker 申请 Queue 的分布式锁。如果旧 Consumer 还没释放锁,新 Consumer 拿不到锁,不能开始消费。
  3. 锁超时释放:如果旧 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 顺序消费

维度MessageListenerConcurrentlyMessageListenerOrderly
线程模型多线程并发处理同一 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 时:

  1. 当前 Queue 的消费暂停。不会跳过这条消息去处理下一条——因为必须保证顺序。
  2. 等待一段时间后重试同一批消息(由 setSuspendCurrentQueueTimeMillis 控制)。
  3. 重试次数累计。达到 maxReconsumeTimes(默认 16 次)后,发往死信队列(DLQ)。

这就是顺序消费的代价:一条消息卡住,整个队列后面的所有消息都卡住。

🤔 思考穿插:为什么失败要返回 SUSPEND_CURRENT_QUEUE_A_MOMENT 暂停整队,而不是跳过这条继续?—— 因为顺序消费不允许跳过——后面的消息可能依赖这条的前序状态,跳过会让后续处理基于错误的上下文。所以只能「暂停 → 重试同一批 → 超限进 DLQ」,代价是一条卡住全队卡住。

生产实践建议

  1. Queue 数量提前规划:顺序消费的并行度 = Queue 数量。Queue 太少,吞吐量上不去;Queue 太多,每个 Queue 消息太少,浪费内存和文件句柄。一般设为 Consumer 节点数 x 核数的 2-4 倍。

  2. 业务标识选择要均匀:用 orderId.hashCode() % queueSizeorderId % queueSize 更好——如果 orderId 是连续自增的,后者会导致所有消息集中在某几个 Queue。

  3. 处理逻辑要有超时:顺序消费时一条卡住全队列卡住,每条消息的处理必须设置超时,超时后 SUSPEND 并告警。

  4. 监控队列积压差异:如果某个 Queue 积压严重而其他 Queue 空闲,说明该 Queue 分配了”热 key”(如某个大商家),需要检查路由算法。

  5. 避免 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 短暂暂停(通常几百毫秒)。


Share this post on:

Previous Post
RocketMQ延时消息的底层实现
Next Post
RocketMQ 消息领域模型:Topic、Tag、MessageQueue、Group 与 Offset 的层次设计