rocketmq延时消息实现
RocketMQ 的延时消息功能允许消息在指定的延迟时间之后才被投递给消费者,非常适合处理像订单超时自动关闭、延迟通知和定时任务触发这类场景
。其核心实现巧妙地利用了“延迟等级”和“内置调度服务”的机制。
🔍 核心实现原理
RocketMQ的延时消息实现不依赖于外部的定时任务系统(如Quartz),而是通过其内置的调度服务(ScheduleMessageService)和延迟等级(Delay Level)机制来完成的
。
- 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.
内部主题(SCHEDULE_TOPIC_XXXX):当你发送一条延时消息时,Broker并不会将它直接存入你指定的目标Topic,而是会先将其Topic属性改写为内部的
。每个延迟等级都对应这个内部主题的一个队列。SCHEDULE_TOPIC_XXXX,并根据设置的延迟等级决定存入哪个特定的队列(QueueId = delayLevel - 1) - 3.
调度服务(ScheduleMessageService):Broker内部有一个后台定时任务(ScheduleMessageService),它会每隔100毫秒扫描一次
。当发现消息的延迟时间已到期(SCHEDULE_TOPIC_XXXX中各个队列的消息当前时间 >= 消息存储时间 + 预设延迟时间),就会将这条消息从临时存储中取出,恢复其原始的Topic和属性,然后作为一条普通消息重新写入CommitLog并投递到目标Topic。此后,消费者就能像消费普通消息一样消费到它了。
⌚ 时间轮算法
RocketMQ的调度服务底层采用了时间轮(Time Wheel)算法来高效管理这些延迟任务
。时间轮是一种高效的批量调度定时任务的算法模型,它将时间分成多个槽位(slot),每个槽位代表一个基本的时间间隔。RocketMQ的18个延迟等级可以看作是18个不同的时间轮槽位。这种方式避免了为每条消息创建单独的定时器,大大降低了系统的开销,使得RocketMQ能够以O(1)的时间复杂度处理大量的延时消息。
更多推荐
所有评论(0)