RocketMQ 的延时消息功能允许消息在指定的延迟时间之后才被投递给消费者,非常适合处理像​​订单超时自动关闭​​、​​延迟通知​​和​​定时任务触发​​这类场景

。其核心实现巧妙地利用了“​​延迟等级​​”和“​​内置调度服务​​”的机制。

🔍 核心实现原理

RocketMQ的延时消息实现不依赖于外部的定时任务系统(如Quartz),而是通过其​​内置的调度服务(ScheduleMessageService)和延迟等级(Delay Level)机制​​来完成的

  1. 1.

    ​延迟等级(Delay Level)​​:RocketMQ预定义了18个延迟级别(1到18),每个级别对应一个固定的延迟时间

    。默认的延迟时间配置如下表所示:

    延迟等级

    延迟时间

    延迟等级

    延迟时间

    1

    1秒

    10

    6分钟

    2

    5秒

    11

    7分钟

    3

    10秒

    12

    8分钟

    4

    30秒

    13

    9分钟

    5

    1分钟

    14

    10分钟

    6

    2分钟

    15

    20分钟

    7

    3分钟

    16

    30分钟

    8

    4分钟

    17

    1小时

    9

    5分钟

    18

    2小时

    ⚠️ 注意:在开源版RocketMQ 5.0之前,通常只能使用这些固定的延迟等级,不能自定义任意时间(如7秒或1天)

    。但RocketMQ 5.x版本支持设置具体的延迟时间戳或毫秒数,提供了更大的灵活性
  2. 2.

    ​内部主题(SCHEDULE_TOPIC_XXXX)​​:当你发送一条延时消息时,Broker并不会将它直接存入你指定的目标Topic,而是会先将其​​Topic属性改写为内部的 SCHEDULE_TOPIC_XXXX​,并根据设置的延迟等级决定存入哪个特定的队列(QueueId = delayLevel - 1)

    。每个延迟等级都对应这个内部主题的一个队列。
  3. 3.

    ​调度服务(ScheduleMessageService)​​:Broker内部有一个​​后台定时任务(ScheduleMessageService)​​,它会​​每隔100毫秒​​扫描一次SCHEDULE_TOPIC_XXXX中各个队列的消息

    。当发现消息的延迟时间已到期(当前时间 >= 消息存储时间 + 预设延迟时间),就会将这条消息从临时存储中取出,​​恢复其原始的Topic和属性​​,然后作为一条普通消息重新写入CommitLog并投递到目标Topic。此后,消费者就能像消费普通消息一样消费到它了。

⌚ 时间轮算法

RocketMQ的调度服务底层采用了​​时间轮(Time Wheel)算法​​来高效管理这些延迟任务

。时间轮是一种高效的批量调度定时任务的算法模型,它将时间分成多个槽位(slot),每个槽位代表一个基本的时间间隔。RocketMQ的18个延迟等级可以看作是18个不同的时间轮槽位。这种方式避免了为每条消息创建单独的定时器,大大降低了系统的开销,使得RocketMQ能够以​​O(1)​​的时间复杂度处理大量的延时消息。

Logo

北京人形旗下天工造物具身智能开源社区,聚焦具身天工与慧思开物两大平台

更多推荐