Skip to content
Go back

RocketMQ延时消息的底层实现

RocketMQ 延时消息:不是”真的等了 30 分钟”

一句话结论(30s)

RocketMQ 延时消息底层不是「真的等 30 分钟」,而是 Broker 拦截后存入内部 Topic、后台定时轮询到点再恢复投递。因为给每条消息单独计时成本高,所以按 delayLevel 分 Queue 存到 SCHEDULE_TOPIC_XXXX,ScheduleMessageService 每 1s 轮询到期消息;代价是约 1 秒的精度偏差。

核心原理(2min)

Broker 收到带 delayTimeLevel 的消息,检查属性后把 Topic 替换为 SCHEDULE_TOPIC_XXXX、QueueId 设为 delayLevel - 1 写入 CommitLog;ScheduleMessageService 为每个 delayLevel 启动定时任务每 1s 轮询,到期后从 ConsumeQueue 读出消息、恢复 REAL_TOPIC 属性、重新 putMessage 到真实 Topic。关键机制:扫描间隔 1s 带来 ±1s 偏差(最坏消息刚好在扫描刚结束时到期要等下一轮);5.x 用 TimerWheel 实现任意时间延时,精度 1s 但每条额外存约 50 bytes 定时器元数据。

底层深入(5-10min)

表面现象

// 发一条延时消息
message.setDelayTimeLevel(16);  // Level 16 = 30 分钟
producer.send(message);

看起来是”消息发出去,30 分钟后 Consumer 收到”。底层完全不是这样。

🤔 思考穿插:为什么「延时消息」看起来是「等 30 分钟」,底层却完全不是「真的等」?—— 因为消息一旦写入就是持久化的静态数据,Broker 不会为每条消息挂一个线程去「倒计时」;它只是把消息暂存到内部 Topic,用后台定时线程「到点才投递」。所以「延时」是「延迟投递」,不是「延迟存储」。

底层流程

第一步:Broker 拦截

Broker 收到延时消息后,检查 delayTimeLevel 属性,不按真实 Topic 投递,而是:

// 源码:CommitLog.putMessage()
String topic = msg.getProperty(MessageConst.PROPERTY_REAL_TOPIC);
int queueId = msg.getProperty(MessageConst.PROPERTY_REAL_QUEUE_ID);

// 替换为内部延时 Topic
msg.setTopic(TopicValidator.RMQ_SYS_SCHEDULE_TOPIC);  // = "SCHEDULE_TOPIC_XXXX"
msg.setQueueId(delayLevel - 1);  // Level 16 → queueId = 15

第二步:ScheduleMessageService 定时扫描

// 源码:ScheduleMessageService.start()
for (Map.Entry<Integer, Long> entry : delayLevelTable.entrySet()) {
    Integer level = entry.getKey();
    Long delayTimeMillis = entry.getValue();
    this.deliverExecutorService.schedule(
        new DeliverDelayedMessageTimerTask(level, delayTimeMillis),
        FIRST_DELAY_TIME,  // = 1000ms
        TimeUnit.MILLISECONDS
    );
}

每个 delayLevel 启动一个定时任务,每 1 秒轮询一次到期消息。

第三步:到期后恢复投递

// 源码:ScheduleMessageService.executeOnTimeup()
// 1. 从 ConsumeQueue 读出消息
// 2. 恢复真实 Topic
msg.putUserProperty(MessageConst.PROPERTY_REAL_TOPIC, originalTopic);
// 3. 重新 putMessage 到 CommitLog
this.messageStore.putMessage(msgInner);

为什么有约 1 秒的偏差

FIRST_DELAY_TIME = 1000ms——每个 delayLevel 的扫描间隔是 1 秒。最坏情况:消息刚好在一轮扫描刚结束时到期,需要等下一轮(+1 秒)才被投递。

消息到期时间:  12:00:00.100
次扫描:      12:00:00.000  → 还没到期
次扫描:      12:00:01.000  → 到期了!投递!
                          ↑ 约 900ms 偏差

对比:原来定时轮询的偏差是 ±5 分钟,延时消息的偏差是 ±1 秒,精度提升了 300 倍。

🤔 思考穿插:为什么轮询会带来 ±1 秒偏差,而不是精确到点投递?—— 因为扫描是周期性的,消息到期后要等下一轮扫描才被发现;最坏情况恰好错过一轮,就要多等一个完整扫描间隔。所以精度由「扫描间隔」决定,轮询天然存在「量化误差」。

RocketMQ 5.x 的改进

5.x 版本引入了任意时间延时(TimerWheel 实现):

🤔 思考穿插:5.x 用 TimerWheel 支持任意时间延时,为什么还要付出「每条多存 ~50 bytes 元数据」的代价?—— 因为任意时间意味着要精确记录每条消息的到期时刻,不再是「第几级」这种粗粒度信息;时间轮需要额外的定时器元数据来定位槽位。所以「精度/灵活度提升」是用「存储开销」换来的。

项目中的应用:订单超时处理

旧方案:@Scheduled 每 5 分钟全表扫描
  → 日均 288 次扫描,230 次无效
  → 超时精度 ±5 分钟

新方案:RocketMQ Level 16(30 分钟)延时消息
  → 订单创建时发一条延时消息,30 分钟后自动触发
  → 超时精度 ±1 秒
  → 日均无效扫描:0

章末提问

追问 1:RocketMQ 延时消息底层为什么不是「真的等 30 分钟」?

结论先行:因为它不是给每条消息挂定时器倒计时,而是 Broker 拦截后存入内部 Topic(SCHEDULE_TOPIC_XXXX),由后台定时线程到点再恢复投递。

因为:给每条消息单独计时成本太高;实际做法是按 delayLevel 分 Queue 存进内部 Topic,ScheduleMessageService 周期轮询到期消息,再恢复 REAL_TOPIC 属性、重新 putMessage 到真实 Topic。所以「延时」本质是「延迟投递」,而不是「延迟存储」。

追问 2:约 1 秒的精度偏差是怎么产生的?

结论先行:来自轮询扫描的量化误差——消息到期后要等下一轮扫描才被投递。

因为:扫描间隔是 1 秒(FIRST_DELAY_TIME = 1000ms),最坏情况消息刚好在一轮扫描刚结束时到期,就得等下一轮(约 +1s)才被投递。所以偏差约 ±1 秒;相比原来定时轮询的 ±5 分钟,精度提升了 300 倍。

追问 3:4.x 的 18 级延时和 5.x 的任意时间延时,本质区别是什么?

结论先行:4.x 是固定 18 级的「时间分桶」扫表、只支持预定义级别;5.x 是 TimerWheel 时间轮、支持任意时间延时且精度 1s,但每条多存约 50 bytes 元数据。

因为:18 级是粗粒度分桶,调度开销极低但有级别上限;时间轮按到期时刻把消息分配到对应槽位、指针周期推进批量投递,精度和灵活性更高,代价是存储和实现更复杂。二者是「空间/精度换调度开销」的不同落点。


Share this post on:

Previous Post
Spring AOP——JDK动态代理和CGLIB的核心区别
Next Post
RocketMQ 顺序消费:从队列路由到线程绑定的完整链路