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);
}
}
}
关键设计:
- 每个级别独立的任务:延时 1s 的级别每 1s 扫一次(其实也改成了 5s),延时 2h 的级别每 5s 扫一次。
- 扫描频率并非严格精确:默认每 5 秒扫描一次。这意味着延时时间是”近似”的,误差在 5 秒以内。
- 利用 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);
实现原理:
- Broker 使用**时间轮(TimingWheel)**替代了固定 18 级的扫表机制。时间轮是 Netty 的
HashedWheelTimer变体——N 个 bucket,每个 bucket 覆盖一个时间槽,消息按照投递时间分配到对应槽位,时间轮的指针周期性推进,到期槽位中的消息批量投递。 - 实现更复杂,内存开销更大,但提供了无限的时间精度。
💭 思考:为什么 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 的延时能力。