Skip to content
Go back

RocketMQ 延时消息原理:18 级延时与 SCHEDULE_TOPIC_XXXX 揭秘

RocketMQ 延时消息原理:18 级延时与隐藏的 SCHEDULE_TOPIC_XXXX

一句话结论(30s)

RocketMQ 延时消息用「时间分桶」替代「独立定时器」,是空间换时间、批量换独立的架构取舍。因为支持任意时间延时需要为每条消息维护一个定时器(百万消息百万定时器),而固定 18 级只需 18 个桶、定时线程只扫 18 个桶,调度开销极低;代价是只支持预定义级别、精度为秒级。

核心原理(2min)

延时消息原样写入 CommitLog,同时把 REAL_TOPIC/REAL_QID/DELAY 属性写入 SCHEDULE_TOPIC_XXXX 的对应 Queue(queueId = delayLevel - 1);ScheduleMessageService 为每个级别启动独立定时任务、每 5s 扫描一次,比较 storeTimestamp + delayMs 与当前时间,到期则投递到真实 Topic。关键机制:利用 ConsumeQueue 的时间顺序,扫到第一条未到期即 break;投递后原消息不删除,靠 offsetTable 记录进度防重复;5.0 用时间轮支持任意时间延时。

底层深入(5-10min)

延时消息:最常用的功能之一

// 发送一条延时消息——30 分钟后检查订单是否支付
Message msg = new Message("OrderTopic", "CHECK", body);
msg.setDelayTimeLevel(16);  // 第 16 级 = 30 分钟延迟
producer.send(msg);

延时消息在电商中无处不在:下单 30 分钟未支付自动取消、支付后 1 小时未确认收货自动完成、3 天后自动好评。

但 RocketMQ 不支持任意时间延时。 它只支持 18 个预定义的延时级别。这是最让新手困惑的点。

为什么只支持 18 级?

这个问题本质上是在问:为什么不用定时器(ScheduledExecutorService)?

答案:如果支持任意时间延时,每条消息都需要一个独立的定时任务。假设有 100 万条延时消息分布在不同的延时时间(1 分钟、3 分钟、7 分钟、13 分钟…),就需要维护 100 万个定时器——这对 JVM 内存和调度器的压力是灾难性的。

RocketMQ 的解法是时间分桶

不存"30 分 15 秒后投递",而是存"第 16 级(30 分钟)"

所有 30 分钟的延时消息都落在同一个时间桶里
定时线程只需要检查 18 个桶(每个延时级别一个),而不是 100 万个独立时间

这不是 RocketMQ 的开源限制,而是架构取舍——用有限的延时级别换取极低的调度开销。

🤔 思考穿插:为什么「任意时间延时」看似简单,RocketMQ 却用 18 级固定桶替代?—— 因为任意时间意味着每条消息要一个独立定时器,百万消息就要百万定时器,内存和调度开销是灾难性的;时间分桶把「按消息计时」变成「按桶扫描」,定时线程只扫 18 个桶。所以它是「用空间/精度换调度开销」的取舍。

18 级延时时间表

// MessageStoreConfig.java
private String messageDelayLevel = "1s 5s 10s 30s 1m 2m 3m 4m 5m 6m 7m 8m 9m 10m 20m 30m 1h 2h";

级别 1:  1s        级别 7:  3m        级别 13: 9m
级别 2:  5s        级别 8:  4m        级别 14: 10m
级别 3:  10s       级别 9:  5m        级别 15: 20m
级别 4:  30s       级别 10: 6m        级别 16: 30m
级别 5:  1m        级别 11: 7m        级别 17: 1h
级别 6:  2m        级别 12: 8m        级别 18: 2h

为什么到 2h 就停了? 因为 2 小时以上的延时消息可靠性难以保证——时间太长,Broker 重启、文件过期、磁盘回收等异常概率大幅上升。如果真的需要更长时间延时,应该用定时任务 + DB + MQ 组合方案,而不是依赖 MQ 本身。

🤔 思考穿插:为什么延时级别到 2h 就停,不干脆支持到 1 天?—— 因为时间越长,Broker 重启、文件过期、磁盘回收等异常概率越高,延时消息的可靠性越难保证;超过 2h 的「延时」本质已经退化成「定时任务」问题,交给 DB + 定时任务更稳。所以 2h 上限是可靠性和功能边界的权衡。

可以修改这个配置来定制延时级别:

# broker.properties
messageDelayLevel = 1s 5s 10s 30s 1m 2m 3m 4m 5m 6m 7m 8m 9m 10m 20m 30m 1h 2h 4h 6h 12h 1d
# 最多支持扩展,但不建议超过 24 小时

延时消息的完整生命周期

第一步:写入 CommitLog(伪装成普通消息)

生产者发送延时消息时,RocketMQ 做了三件事:

生产者发送:Topic=OrderTopic, delayLevel=16(30min)

Broker 处理:
1. 消息原样写入 CommitLog(ConsumeQueue 也是按 OrderTopic 落盘)
2. 但在 CommitLog 的消息属性中塞入:
   - REAL_TOPIC = "OrderTopic"      ← 原始 Topic
   - REAL_QID   = 0                 ← 原始 Queue ID
   - DELAY      = 16                ← 延时级别
3. 写入 "SCHEDULE_TOPIC_XXXX" 的 ConsumeQueue(按延时级别对应的 QueueId)

第二步:SCHEDULE_TOPIC_XXXX——延时消息的中转站

SCHEDULE_TOPIC_XXXX 是 RocketMQ 内部的一个特殊 Topic,用于暂存延时消息。它的 Queue 有特殊含义:

SCHEDULE_TOPIC_XXXX:
  Queue-0:  延时级别 1(1s)
  Queue-1:  延时级别 2(5s)
  Queue-2:  延时级别 3(10s)
  ...
  Queue-17: 延时级别 18(2h)

每个 Queue 对应一个延时级别!

延时消息写入时,根据 delayLevel 写入对应的 Queue:

// CommitLogDispatcherBuildConsumeQueue(简化)
if (msg.getDelayTimeLevel() > 0) {
    // 延时消息 → 写入 SCHEDULE_TOPIC_XXXX 的对应 Queue
    int delayLevel = msg.getDelayTimeLevel();
    int queueId = delayLevel - 1;  // delayLevel 从 1 开始,QueueId 从 0 开始
    
    // 构建 ConsumeQueue 条目存入 SCHEDULE_TOPIC_XXXX/queueId/
    consomeQueue.put(SCHEDULE_TOPIC_XXXX, queueId, commitLogOffset, size, tagsCode);
}
存储结构:
/store/consumequeue/
  └── SCHEDULE_TOPIC_XXXX/
      ├── 0/   ← 延时 1s  的消息的 ConsumeQueue
      ├── 1/   ← 延时 5s  的消息的 ConsumeQueue
      ├── 2/   ← 延时 10s 的消息的 ConsumeQueue
      ├── ...
      └── 17/  ← 延时 2h  的消息的 ConsumeQueue

第三步:定时调度线程扫描

ScheduleMessageService 是延时消息的核心调度器。它在后台运行,定期检查各个延时级别是否有到期消息。

// ScheduleMessageService(简化核心逻辑)
public class ScheduleMessageService {
    // 每个延时级别对应一个延时任务
    private final ConcurrentHashMap<Integer, Long> delayLevelTable;

    public void start() {
        // 为每个延时级别创建定时任务
        for (int level = 1; level <= maxDelayLevel; level++) {
            // 第一个任务:延时对应的间隔后启动
            // 后续:每隔 5 秒扫描一次
            scheduledExecutorService.scheduleAtFixedRate(
                new DeliverDelayedMessageTimerTask(level),
                getDelayTime(level),   // 首次延迟
                5000,                  // 之后每 5 秒执行一次
                TimeUnit.MILLISECONDS
            );
        }
    }

    class DeliverDelayedMessageTimerTask implements Runnable {
        private int delayLevel;

        @Override
        public void run() {
            // 从 SCHEDULE_TOPIC_XXXX 的对应 Queue 中拉取到期消息
            long offset = offsetTable.get(delayLevel);
            ConsumeQueue cq = findConsumeQueue(SCHEDULE_TOPIC_XXXX, delayLevel - 1);
            SelectMappedBufferResult buffer = cq.getIndexBuffer(offset);

            for (ConsumeQueue entry : cq.iterateFrom(offset)) {
                // 从 CommitLog 读取原始消息
                MessageExt msg = commitLog.getMessage(entry.getCommitLogOffset());

                // 计算消息在 CommitLog 中的存储时间
                long storeTimestamp = msg.getStoreTimestamp();
                long delayMs = getDelayTime(delayLevel) * 1000L;
                long expectedDeliveryTime = storeTimestamp + delayMs;

                if (System.currentTimeMillis() >= expectedDeliveryTime) {
                    // 到期了 → 投递到原始 Topic
                    deliverToOriginalTopic(msg);
                    offset++;
                } else {
                    break;  // 还没到期,后面的更不可能到期
                }
            }

            offsetTable.put(delayLevel, offset);
        }
    }
}

关键设计:

  1. 每个级别独立的任务:延时 1s 的级别每 1s 扫一次(其实也改成了 5s),延时 2h 的级别每 5s 扫一次。
  2. 扫描频率并非严格精确:默认每 5 秒扫描一次。这意味着延时时间是”近似”的,误差在 5 秒以内。
  3. 利用 ConsumeQueue 的时间顺序:同一个延时级别的消息在 ConsumeQueue 中按时间排序(因为写入顺序 = 时间顺序),所以扫描到第一条未到期的就可以 break。

💭 思考:为什么扫到第一条未到期的消息就可以 break,而不是继续扫完整个桶?—— 因为同一延时级别的消息「写入顺序 = 到期时间顺序」,ConsumeQueue 天然按时间排序;既然这条还没到期,它后面的消息写入更晚、到期也更晚,扫下去纯属浪费。这个 break 依赖的是「时间分桶 + 顺序追加」共同保证的有序性,是把「按时间遍历」从 O(全量) 降成 O(已到期) 的关键。

第四步:投递到原始 Topic

到期消息从 SCHEDULE_TOPIC_XXXX “搬”到原始 Topic:

private void deliverToOriginalTopic(MessageExt msgExt) {
    // 从消息属性中恢复原始路由信息
    String realTopic = msgExt.getProperty("REAL_TOPIC");
    int realQueueId = Integer.parseInt(msgExt.getProperty("REAL_QID"));

    // 构造一条新消息(清除延时属性,恢复原始 Topic)
    MessageExt newMsg = new MessageExt();
    newMsg.setTopic(realTopic);
    newMsg.setQueueId(realQueueId);
    newMsg.setBody(msgExt.getBody());
    // ... 拷贝其他属性

    // 写入原始 Topic 的 CommitLog 和 ConsumeQueue
    // (正常情况下由 CommitLog 分发逻辑处理)
    messageStore.putMessage(newMsg);
}

注意:投递后原始延时消息不被删除——它仍然在 SCHEDULE_TOPIC_XXXX 的 ConsumeQueue 中。但 Broker 通过 offsetTable 记录消费进度,下次从 offset 之后开始扫描,不会重复投递。

🤔 思考穿插:到期投递后为什么不删掉 SCHEDULE_TOPIC_XXXX 里的原消息,还要靠 offsetTable 记进度?—— 因为 ConsumeQueue 是只追加的顺序文件,删除中间一条会破坏其顺序结构、代价极高;用 offsetTable 跳过已投递的消息,既保证不重复投递,又保留了原始记录方便排查。所以「不删 + 记进度」是追加写存储下的标准做法。


延时消息的误差来源

误差来源描述典型误差量
扫描间隔定时任务每 5s 扫描一次最多 5s
消息写入耗时CommitLog 写入时间 vs 存储时间戳的精度< 1ms
投递耗时从 SCHEDULE_TOPIC 写到原始 Topic 的 CommitLog 写入< 1ms
ConsumeQueue 构建延迟投递后,消费者需要等 ConsumeQueue 异步构建1-5ms

综合来看,延时精度在 “秒级”,不适合需要精确到毫秒的场景(如高频交易)。


生产实践建议

1. 不要依赖延时消息的精确性

// ❌ 不靠谱——5s 扫描间隔意味着可能延时 30 分 0 秒到 30 分 5 秒之间到达
msg.setDelayTimeLevel(16);  // 期望 30 分钟

// ✅ 在消费端检查实际时间
if (order.getCreatedAt() + Duration.ofMinutes(30).isAfter(Instant.now())) {
    // 还没到 30 分钟 → 忽略(或重新发一条延时消息)
    return;
}

2. 延时消息数量和级别规划

// 不要这样——100 万条消息全部延时 30 分钟
// SCHEDULE_TOPIC_XXXX/queue-15(30min)会积压 100 万条
// 每 5 秒扫描一次,扫描 100 万条需要很久

// 尽量分散到不同延时级别

3. 超过 2 小时的延时需求

// 方案:定时任务 + MQ
// 每小时扫描一次数据库,找出需要处理的过期订单
// 找到后再发 MQ 消息,即时消费

@Scheduled(cron = "0 0 * * * *")  // 每小时执行
public void processExpiredOrders() {
    List<Order> expired = orderDao.findExpired(30, TimeUnit.MINUTES);
    for (Order order : expired) {
        mqProducer.send(buildMessage(order));  // 发送即时消息
    }
}

4. 修改延时级别需要重启

# 修改后需要重启 Broker 才能生效
messageDelayLevel = 1s 5s 10s 30s 1m 2m 3m 4m 5m 6m 7m 8m 9m 10m 20m 30m 1h 2h

并且不能删除已有的级别——如果之前有消息设了级别 16(30 分钟),而你把 16 改成了 40 分钟,那些已经写入的消息仍然会按旧的 ScheduleMessageService 配置处理,可能导致行为不一致。


RocketMQ 5.0 的改进:任意时间延时

RocketMQ 5.0 开始支持任意时间的定时消息

// RocketMQ 5.0+
Message msg = new Message("OrderTopic", "CHECK", body);
// 设置精确的投递时间
msg.setDeliverTimeMs(System.currentTimeMillis() + 1800_000);  // 30 分钟后的毫秒时间戳
producer.send(msg);

实现原理:

💭 思考:为什么 5.0 要用时间轮替代 18 级分桶,时间轮凭什么能支持任意时间?—— 因为时间轮把时间轴切成 N 个槽,每条消息只记「投递时间戳」,指针周期性推进到哪个槽就投递哪个槽,消息数和定时器数解耦了——不再「一条消息一个定时器」,而是「一个时间槽一批消息」。代价是实现更复杂、内存更大,所以 4.x 才用固定 18 级这种更省事的方案;时间轮是「精度换复杂度」的升级。

如果生产环境还是 RocketMQ 4.x,只能使用 18 级延时。升级到 5.0 后可以用任意时间延时。

总结

RocketMQ 延时消息的实现可以浓缩为四步:

1. 延时消息 → 写入 SCHEDULE_TOPIC_XXXX(按 delayLevel 分 Queue)
2. ScheduleMessageService → 18 个延时级别的独立定时扫描线程
3. 扫描 ConsumeQueue → 检查消息存储时间 + delay 是否 ≤ 当前时间
4. 到期消息 → 投递到原始 Topic(REAL_TOPIC)

核心心法:延时不是对每条消息单独计时,而是把相同延时时间的消息放入同一个桶,定时批量扫描桶。 这是典型的”用空间换时间,用批量换独立”的设计思想。

章末提问

追问 1:RocketMQ 延时消息为什么不支持任意时间,18 级是怎么来的?

结论先行:因为任意时间需要为每条消息维护一个独立定时器,百万消息百万定时器撑不住,所以用 18 级固定时间分桶替代。

因为:18 个延时级别 = 18 个桶 = 18 个定时任务,调度开销是 O(级别数) 而非 O(消息数);每条消息只记录「第几级」,到点后由定时线程批量扫描投递。代价是只能选预定义级别、精度秒级——这是「空间换时间、批量换独立」的取舍。

追问 2:延时消息投递出去后,原来的消息去哪儿了?会不会重复投递?

结论先行:原消息不删除、仍留在 SCHEDULE_TOPIC_XXXX,靠 offsetTable 记录投递进度,所以不会重复投递。

因为:ConsumeQueue 是只追加的顺序文件,删除中间一条会破坏顺序结构、代价极高;每次扫描都从 offsetTable 记录的 offset 之后继续,投递过的消息不会再被扫到,于是「不删 + 记进度」既保留了原始记录,又保证了不重复投递。

追问 3:延时消息的精度为什么是秒级?如果要毫秒级怎么办?

结论先行:秒级精度来自定时扫描间隔(4.x 默认每 5s 扫一次,5.0 时间轮也仍是 1s 精度);要毫秒级就不能用 RocketMQ 延时消息。

因为:扫描是「批量轮询桶」而非「每条消息单独计时」,消息到期后要等下一轮扫描才被投递,误差最大等于扫描间隔;需要毫秒级精确触发的场景(如高频交易)应改用定时任务或其他专门组件,而不是依赖 MQ 的延时能力。


Share this post on:

Previous Post
RocketMQ架构——为什么CommitLog所有Topic共用一份文件
Next Post
RocketMQ 刷盘机制:同步刷盘、异步刷盘与顺序写磁盘