【RocketMQ 生产者消费者】- 延时消息原理解析-ScheduleMessageService
本文章基于 RocketMQ 4.9.3
1. 前言
- 【RocketMQ】- 源码系列目录
- 【RocketMQ 生产者消费者】- 同步、异步、单向发送消费消息
- 【RocketMQ 生产者和消费者】- 消费者启动源码
- 【RocketMQ 生产者和消费者】- 消费者重平衡(1)
- 【RocketMQ 生产者和消费者】- 消费者重平衡(2)- 分配策略
- 【RocketMQ 生产者和消费者】- 消费者重平衡(3)- 消费者 ID 对负载均衡的影响
- 【RocketMQ 生产者和消费者】- 消费者的订阅关系一致性
- 【RocketMQ 生产者和消费者】- 消费者发起消息拉取请求 PullMessageService
- 【RocketMQ 生产者和消费者】- broker 是如何处理消费者消息拉取的 Netty 请求的
- 【RocketMQ 生产者和消费者】- broker 处理消息拉取请求
- 【RocketMQ 生产者和消费者】- 消费者处理消息拉取结果
- 【RocketMQ 生产者和消费者】- 消费者处理消息拉取结果
- 【RocketMQ 生产者和消费者】- ConsumeMessageConcurrentlyService 并发消费消息
- 【RocketMQ 生产者和消费者】- ConsumeMessageOrderlyService 顺序消费消息
- 【RocketMQ 生产者和消费者】- sendMessageBack 发送重试消息
- 【RocketMQ 生产者和消费者】- 延时消息的使用
上一篇文章主要是看了延时消息的发送,这篇文章来看下延时消息的原理。
2. 延时服务初始化和启动

延时消息服务是在创建 DefaultMessageStore 的时候创建出来的,同时在 broker 启动的时候也会调用 handleScheduleMessageService 来启动 scheduleMessageService。
3. ScheduleMessageService#start 启动延时消息服务
/**
* 启动延时消息服务
*/
public void start() {
// 只能启动一次
if (started.compareAndSet(false, true)) {
// 从 ${user.home}/store/config/delayOffset.json 中加载延时 topic 的消费偏移量到 offsetTable 中
super.load();
// 1. 初始化延迟消息投递的线程池, 专门处理延时消息
this.deliverExecutorService = new ScheduledThreadPoolExecutor(this.maxDelayLevel, new ThreadFactoryImpl("ScheduleMessageTimerThread_"));
if (this.enableAsyncDeliver) {
// 如果可以异步投递, 初始化延迟消息异步投递的线程池, 默认不支持
this.handleExecutorService = new ScheduledThreadPoolExecutor(this.maxDelayLevel, new ThreadFactoryImpl("ScheduleMessageExecutorHandleThread_"));
}
// 2. 遍历所有延时级别, 对所有延时等级都创建一个调度任务 DeliverDelayedMessageTimerTask, 专门来处理这个延迟级别下面的延时消息
for (Map.Entry<Integer, Long> entry : this.delayLevelTable.entrySet()) {
// 延时等级
Integer level = entry.getKey();
// 这个延时等级下面的消息延时时间
Long timeDelay = entry.getValue();
// 消费者消费偏移量
Long offset = this.offsetTable.get(level);
if (null == offset) {
// 如果不存在就是从 0 开始, 也就是从头开始处理
offset = 0L;
}
// 为该等级的延迟队列构建一个调度任务,初始化是 1s 后执行
if (timeDelay != null) {
if (this.enableAsyncDeliver) {
this.handleExecutorService.schedule(new HandlePutResultTask(level), FIRST_DELAY_TIME, TimeUnit.MILLISECONDS);
}
// 这里面会传入延时等级和偏移量
this.deliverExecutorService.schedule(new DeliverDelayedMessageTimerTask(level, offset), FIRST_DELAY_TIME, TimeUnit.MILLISECONDS);
}
}
// 3. 添加一个定时调度任务, 专门用来每隔 10s 持久化 offsetTable 到 ${user.home}/store/config/delayOffset.json 中
this.deliverExecutorService.scheduleAtFixedRate(new Runnable() {
@Override
public void run() {
try {
if (started.get()) {
ScheduleMessageService.this.persist();
}
} catch (Throwable e) {
log.error("scheduleAtFixedRate flush exception", e);
}
}
// 添加完消息之后 10s 执行第一次任务, 后面每隔 10s 执行一次
}, 10000, this.defaultMessageStore.getMessageStoreConfig().getFlushDelayOffsetInterval(), TimeUnit.MILLISECONDS);
}
}
started 是延时消息服务状态,初始化启动的时候会设置为 false。
/**
* 延时消息服务状态
*/
private final AtomicBoolean started = new AtomicBoolean(false);
接下来通过 super.load() 从 ${user.home}/store/config/delayOffset.json 中加载延时 topic 的消费偏移量到 offsetTable 中,因为延时消息每一个等级都是一个队列,所以延时消息也会存储到这些队列里面,因此也会有偏移量,来看下这个文件的内容。

这里的内容就是 queueId: offset,这个 offset 跟普通消息是一样的,就是下一条要消费的 ConsumeQueue 索引下标,后面也会分析这个类的 load 方法。
接下来初始化延迟消息投递的线程池,专门处理延时消息,专门用于检测延时消息有没有到时间,如果到时间了就把这条消息投递到真实队列里面,然后这里还涉及一个 handleExecutorService 异步投递线程池,后面也会说的。
接下来遍历所有延时级别,对所有延时等级都创建一个调度任务 DeliverDelayedMessageTimerTask,专门来处理这个延迟级别下面的延时消息,注意这里是每一个延时等级都会有自己的 DeliverDelayedMessageTimerTask 去处理不同的 queueId。可以看到就构建出来的 DeliverDelayedMessageTimerTask 里面会传入 level 和 timeDelay,level 是用来获取真实的 queueId 的,延时等级从 1 开始计算,queueId 则是从 0 开始计算,因此 queueId = level - 1,而 timeDelay 是用来计算当前的消息有没有到期了,如果到期就投递到真实的 topic。
最后添加一个定时任务,添加一个定时调度任务,专门用来每隔 10s 持久化 offsetTable 到 ${user.home}/store/config/delayOffset.json 中。
3.1 load 加载延时消息信息
/**
* 加载延时消息的数据
* @return
*/
@Override
public boolean load() {
// 从 ${user.home}/store/config/delayOffset.json 中加载延时 topic 的消费偏移量
boolean result = super.load();
// 解析延时等级到 delayLevelTable 集合中
result = result && this.parseDelayLevel();
// 矫正每个延时等级的偏移量
result = result && this.correctDelayOffset();
// 返回是否加载成功
return result;
}
加载信息分为三个步骤:
- 从 ${user.home}/store/config/delayOffset.json 中加载延时 topic 的消费偏移量。
- 解析延时等级到 delayLevelTable 集合中。
- 矫正每个延时等级的偏移量,纠正应该是说写 offsetTable 是由定时任务来完成的,由于是每隔 10s 写一次,所以可能会导致偏移量在短时间(10s)内的不一致。
3.2 parseDelayLevel 解析延时等级
解析延时等级, 将[延时等级 -> 延时时间]的关系存储到 delayLevelTable 中。
/**
* 解析延时等级, 将[延时等级 -> 延时时间]的关系存储到 delayLevelTable 中
* @return
*/
public boolean parseDelayLevel() {
HashMap<String, Long> timeUnitTable = new HashMap<String, Long>();
// 时间表, 单位是毫秒
timeUnitTable.put("s", 1000L);
timeUnitTable.put("m", 1000L * 60);
timeUnitTable.put("h", 1000L * 60 * 60);
timeUnitTable.put("d", 1000L * 60 * 60 * 24);
// 延时等级 "1s 5s 10s 30s 1m 2m 3m 4m 5m 6m 7m 8m 9m 10m 20m 30m 1h 2h"
String levelString = this.defaultMessageStore.getMessageStoreConfig().getMessageDelayLevel();
try {
String[] levelArray = levelString.split(" ");
for (int i = 0; i < levelArray.length; i++) {
// 获取等级, 比如 2m
String value = levelArray[i];
// 获取时间单位 m
String ch = value.substring(value.length() - 1);
// 获取这个单位下面的毫秒数, m 就是 60 * 1000
Long tu = timeUnitTable.get(ch);
// 延时等级, 从 1 开始
int level = i + 1;
// 记录最大延时等级
if (level > this.maxDelayLevel) {
this.maxDelayLevel = level;
}
// 获取 2m 前面的 2
long num = Long.parseLong(value.substring(0, value.length() - 1));
// 计算延时时间, 2 * 60 * 1000
long delayTimeMillis = tu * num;
// 添加到集合中
this.delayLevelTable.put(level, delayTimeMillis);
if (this.enableAsyncDeliver) {
this.deliverPendingTable.put(level, new LinkedBlockingQueue<>());
}
}
} catch (Exception e) {
log.error("parseDelayLevel exception", e);
log.info("levelString String = {}", levelString);
return false;
}
return true;
}
首先来看下 messageDelayLevel,这个就是 RocketMQ 提供的延时等级,默认就是上一篇文章说的 18 个延时等级。
/**
* 延时等级, 最小就是 1s
*/
private String messageDelayLevel = "1s 5s 10s 30s 1m 2m 3m 4m 5m 6m 7m 8m 9m 10m 20m 30m 1h 2h";
接下来遍历这个延时等级,然后获取对应的时间单位统一转成毫秒,比如 1m 就对应了 60 * 1000 ms,1h 就对应了 60 * 60 * 1000 ms,最后将这些时间统计放到 delayLevelTable 中,如果说支持 延时投递,那么每一个等级都会设置一个阻塞队列。要注意一些就是 level = i + 1,也就是说 level 是从 1 开始计算的。
3.3 correctDelayOffset 纠正偏移量
/**
* 纠正延时偏移量
* @return
*/
public boolean correctDelayOffset() {
try {
// 遍历延时集合的所有延时等级
for (int delayLevel : delayLevelTable.keySet()) {
// 获取这个延时级别的消息队列, 队列 ID 是延时级别 - 1
ConsumeQueue cq =
ScheduleMessageService.this.defaultMessageStore.findConsumeQueue(TopicValidator.RMQ_SYS_SCHEDULE_TOPIC,
delayLevel2QueueId(delayLevel));
// 获取延时队列消费偏移量
Long currentDelayOffset = offsetTable.get(delayLevel);
if (currentDelayOffset == null || cq == null) {
// 如果不存在就不用管
continue;
}
long correctDelayOffset = currentDelayOffset;
// 偏移量最小的有效 ConsumeQueue 消息的索引
long cqMinOffset = cq.getMinOffsetInQueue();
// 偏移量最大的有效 ConsumeQueue 消息的索引
long cqMaxOffset = cq.getMaxOffsetInQueue();
// 如果当前延时队列的最小值比 ConsumeQueue 中的消息最小索引都要小
if (currentDelayOffset < cqMinOffset) {
// 修正 correctDelayOffset 为 cqMinOffset
correctDelayOffset = cqMinOffset;
log.error("schedule CQ offset invalid. offset={}, cqMinOffset={}, cqMaxOffset={}, queueId={}",
currentDelayOffset, cqMinOffset, cqMaxOffset, cq.getQueueId());
}
// 如果当前延时队列的最小值比 ConsumeQueue 中的消息最大索引都要大
if (currentDelayOffset > cqMaxOffset) {
// 修正 correctDelayOffset 为 cqMaxOffset
correctDelayOffset = cqMaxOffset;
log.error("schedule CQ offset invalid. offset={}, cqMinOffset={}, cqMaxOffset={}, queueId={}",
currentDelayOffset, cqMinOffset, cqMaxOffset, cq.getQueueId());
}
// 最后如果这两个不相同, 说明修正了
if (correctDelayOffset != currentDelayOffset) {
log.error("correct delay offset [ delayLevel {} ] from {} to {}", delayLevel, currentDelayOffset, correctDelayOffset);
// 添加到 offsetTable 集合中
offsetTable.put(delayLevel, correctDelayOffset);
}
}
} catch (Exception e) {
log.error("correctDelayOffset exception", e);
return false;
}
return true;
}
每一个延时等级对应一个队列,所以遍历延时集合的所有延时等级,然后获取这个延时级别的消息队列,队列 ID 是延时级别 - 1,然后获取这个队列的消费偏移量,如果不存在就直接返回不用管了,因为这个延时等级下还没有消费者发送延时消息。
获取到对应的消息队列之后,获取偏移量最小的有效 ConsumeQueue 消息的索引下标以及偏移量最大的有效 ConsumeQueue 消息的索引下标,然后判断如果当前延时队列的偏移量最小值比 ConsumeQueue 中的消息最小索引都要小,那么修正最小值,如果当前延时队列的最小值比 ConsumeQueue 中的消息最大索引都要大,那么修正最大值。
修正之后再把偏移量添加到 offsetTable 中,比如上一次服务关闭的时候一些偏移量没来得及刷新到本地文件中,这时候 offsetTable 记录的偏移量就会比对应的 ConsumeQueue 记录的最大偏移量要小,这种情况下就会更新最大值。
4. DeliverDelayedMessageTimerTask 处理延时队列的消息
这个就是真正处理延时队列的请求,里面包含了两个属性:delayLevel 和 offset,代表延时等级和消费偏移量,然后来看下里面的 run 方法。
@Override
public void run() {
try {
// 服务没有停止就不断执行
if (isStarted()) {
this.executeOnTimeup();
}
} catch (Exception e) {
// XXX: warn and notify me
log.error("ScheduleMessageService, executeOnTimeup exception", e);
this.scheduleNextTimerTask(this.offset, DELAY_FOR_A_PERIOD);
}
}
isStarted 意思是判断延时服务有没有停止,如果没有停止就不断去里面检测消息有没有到期。
public boolean isStarted() {
return started.get();
}
4.1 executeOnTimeup 处理投递的延时消息 - 核心逻辑
这里就是处理延时消息的核心逻辑,下面来看下里面的逻辑,首先就是根据延时 topic 和 延时队列 id 从 consumeQueueTable 中找出要写入的 ConsumeQueue,没找到就新建一个,其中查询的 topic 是 SCHEDULE_TOPIC_XXXX,queueId 是 delayLevel - 1,delayLevel2QueueId 就是通过 delayLevel - 1 算出队列 id 的。
// 1. 根据延时 topic 和 延时队列 id 从 consumeQueueTable 中找出要写入的 ConsumeQueue, 没找到就新建一个
ConsumeQueue cq =
ScheduleMessageService.this.defaultMessageStore.findConsumeQueue(TopicValidator.RMQ_SYS_SCHEDULE_TOPIC,
delayLevel2QueueId(delayLevel));
如果最终还是找不到队列,那么延时 100ms 之后再次执行,所以这里也可以看到为什么用一个线程任务来处理对应的延时队列,就是因为如果当前任务结束了,就会重新提交一个新的定时任务,并不会出现这个延时队列没有线程处理的情况。
// 2. 这里也找不到(并发冲突?)
if (cq == null) {
// 延时 100ms 后再次执行
this.scheduleNextTimerTask(this.offset, DELAY_FOR_A_WHILE);
return;
}
获取到队列之后,根据 offset 开始截取一段 Buffer,注意这里截取到的 buffer 就是从 offset 开始的一段 buffer,里面可能包含多条消息。
SelectMappedBufferResult bufferCQ = cq.getIndexBuffer(this.offset);
if (bufferCQ == null) {
// 没找到缓存 buffer, 看看是不是拉取的位置不对
long resetOffset;
// 如果说 ConsumeQueue 中有效消息的最小偏移量比传进来的 offset 要大, 说明拉取的点位太靠前了
if ((resetOffset = cq.getMinOffsetInQueue()) > this.offset) {
log.error("schedule CQ offset invalid. offset={}, cqMinOffset={}, queueId={}",
this.offset, resetOffset, cq.getQueueId());
// 如果说 ConsumeQueue 中有效消息的最大偏移量比传进来的 offset 要小, 说明拉取的点位太靠后了
} else if ((resetOffset = cq.getMaxOffsetInQueue()) < this.offset) {
log.error("schedule CQ offset invalid. offset={}, cqMaxOffset={}, queueId={}",
this.offset, resetOffset, cq.getQueueId());
} else {
// 这里就是正常情况
resetOffset = this.offset;
}
// 新建一个 DeliverDelayedMessageTimerTask 丢到线程池 deliverExecutorService 中去处理, 延时 100ms
this.scheduleNextTimerTask(resetOffset, DELAY_FOR_A_WHILE);
return;
}
// 下一次消费的 offset
long nextOffset = this.offset;
如果没有拉取到消息,也就是返回的 bufferCQ 为 null,这里返回 null 的原因比较多,比如这个偏移量找不到 MappedFile,又或者这个偏移量比队列的最小偏移量要小,如果获取不到消息最终返回的 SelectMappedBufferResult 不为空,但是由于 readPosition == pos,返回的 buffer 里 limit = 0,会获取不到消息。这种情况下也是一样修改偏移量,然后重新提交一个任务,延时 100ms 再提交。
接下来遍历所有消息,开始处理,接下来就是具体处理消息的逻辑。
int i = 0;
ConsumeQueueExt.CqExtUnit cqExtUnit = new ConsumeQueueExt.CqExtUnit();
// 4. 处理 buffer 中的数据, 遍历这些延时消息看看是不是到期了, 如果到期了就将这些消息投递回原始队列
for (; i < bufferCQ.getSize() && isStarted(); i += ConsumeQueue.CQ_STORE_UNIT_SIZE) {
...
}
首先获取 ConsumeQueue 三件套,物理偏移量、消息大小、延时消息的 tagsCode,这里要注意延时消息这里的 tagsCode 就是延时消息的发送时间,也就是什么时候发送到真实 topic,当消息添加到 CommitLog 之后会通过消息重放服务将 tagsCode 替换成延时消息的发送时间,然后构建 ConsumeQueue 索引,所以这里拉到的消息的 tagsCode 就是时间戳了。
// 获取第 i 条消息的物理偏移量
long offsetPy = bufferCQ.getByteBuffer().getLong();
// 获取第 i 条消息的大小
int sizePy = bufferCQ.getByteBuffer().getInt();
// 获取第 i 条消息的 tagsCode, 注意如果是延时消息这里的 tagsCode 就是延时消息的发送时间, 也就是什么时候发送到真实 topic,
// 当消息添加到 CommitLog 之后会通过消息重放服务将 tagsCode 替换成延时消息的发送时间, 然后构建 ConsumeQueue 索引, 所以
// 这里拉到的消息的 tagsCode 就是时间戳了
long tagsCode = bufferCQ.getByteBuffer().getLong();
然后下面判断扩展文件地址,目前还没用过这个。
// 扩展文件地址, 不知道有什么用
if (cq.isExtAddr(tagsCode)) {
if (cq.getExt(tagsCode, cqExtUnit)) {
tagsCode = cqExtUnit.getTagsCode();
} else {
//can't find ext content.So re compute tags code.
log.error("[BUG] can't find consume queue extend file content!addr={}, offsetPy={}, sizePy={}",
tagsCode, offsetPy, sizePy);
long msgStoreTime = defaultMessageStore.getCommitLog().pickupStoreTimestamp(offsetPy, sizePy);
tagsCode = computeDeliverTimestamp(delayLevel, msgStoreTime);
}
}
接下来判断下这条消息有没有到投递时间,如果没有,说明后面的消息也不需要处理了,延迟 100ms 之后接着提交一个 DeliverDelayedMessageTimerTask 任务去扫描延时消息,本次线程任务处理结束。至于为什么可以直接结束,因为这个队列的延时等级都是一样的,而且消息是顺序添加,如果当前 ConsumeQueue 处理到第 i 条消息还没到期,那么后面的消息也不会到期。
// 当前时间戳
long now = System.currentTimeMillis();
// 校验当前时间是不是到投递时间了, 返回投递时间
long deliverTimestamp = this.correctDeliverTimestamp(now, tagsCode);
// 设置下一个消费偏移量
nextOffset = offset + (i / ConsumeQueue.CQ_STORE_UNIT_SIZE);
// 剩余多少时间才执行
long countdown = deliverTimestamp - now;
if (countdown > 0) {
// 如果还没到投递时间, 延迟 100ms 之后接着提交一个 DeliverDelayedMessageTimerTask 任务去扫描延时消息, 本次处理结束
this.scheduleNextTimerTask(nextOffset, DELAY_FOR_A_WHILE);
return;
}
那如果 countdown < 0,说明到期了,这时候根据消息偏移量从 CommitLog 中找到这条消息,如果找不到就处理下一条不处理这条了,这里为什么不是再重新提交一个定时任务,是因为 ConsumeQueue 是在写入 CommitLog 之后通过消息重放生成的,所以 ConsumeQueue 索引找不到具体的消息就不管了。
// 这里就是到投递时间了, 根据消息偏移量从 CommitLog 中找到这条消息
MessageExt msgExt = ScheduleMessageService.this.defaultMessageStore.lookMessageByOffset(offsetPy, sizePy);
if (msgExt == null) {
// 如果没找到, 继续下一条消息
continue;
}
接下来构建内部消息对象, 这时候构建出来的消息的 topic 就是真实的 topic 了, 而 queueId 也是原始消息的 queueId,然后这里会判断如果是半事务消息,直接丢弃这条延时消息,可以看到 RocketMQ 的延时消息是 不支持事务的。
// 构建内部消息对象, 这时候构建出来的消息的 topic 就是真实的 topic 了, 而 queueId 也是原始消息的 queueId
MessageExtBrokerInner msgInner = ScheduleMessageService.this.messageTimeup(msgExt);
if (TopicValidator.RMQ_SYS_TRANS_HALF_TOPIC.equals(msgInner.getTopic())) {
// 延时半消息, 有 BUG, 延时消息是不支持事务消息的, 丢弃这条消息
log.error("[BUG] the real topic of schedule msg is {}, discard the msg. msg={}",
msgInner.getTopic(), msgInner);
continue;
}
下面就是具体的消息投递逻辑。
// 消息投递
boolean deliverSuc;
if (ScheduleMessageService.this.enableAsyncDeliver) {
// 异步投递消息, 默认是不支持的
deliverSuc = this.asyncDeliver(msgInner, msgExt.getMsgId(), offset, offsetPy, sizePy);
} else {
// 同步投递消息, offset 应该是传的 nextOffset 才对, 不应该把最开始的 offset 返回, 4.9.4 已修复
deliverSuc = this.syncDeliver(msgInner, msgExt.getMsgId(), offset, offsetPy, sizePy);
}
如果投递失败,重新提交一个 DeliverDelayedMessageTimerTask 延时 100ms 之后再次执行, 本地调度结束。
if (!deliverSuc) {
// 投递失败, 重新提交一个 DeliverDelayedMessageTimerTask 延时 100ms 之后再次执行, 本地调度结束
this.scheduleNextTimerTask(nextOffset, DELAY_FOR_A_WHILE);
return;
}
nextOffset 指向下一条要处理的消息。
// 处理下一条消息
nextOffset = this.offset + (i / ConsumeQueue.CQ_STORE_UNIT_SIZE);
最后如果处理完消息了, 新建一个 DeliverDelayedMessageTimerTask 延迟 100ms 之后再次执行,DELAY_FOR_A_WHILE 就是 100ms。
// 遍历结束, 新建一个 DeliverDelayedMessageTimerTask 延迟 100ms 之后再次执行, 可以说这个方法就是每 100ms 扫描一次队列
// 里面的消息找到需要投递的延时消息投递到真实队列中
this.scheduleNextTimerTask(nextOffset, DELAY_FOR_A_WHILE);
4.2 scheduleNextTimerTask 重新投递延时队列消息
/**
* 新建一个 DeliverDelayedMessageTimerTask 投递到 deliverExecutorService 线程池中
* @param offset 消息消费偏移量
* @param delay 延时时间
*/
public void scheduleNextTimerTask(long offset, long delay) {
ScheduleMessageService.this.deliverExecutorService.schedule(new DeliverDelayedMessageTimerTask(
this.delayLevel, offset), delay, TimeUnit.MILLISECONDS);
}
4.3 computeDeliverTimestamp 获取延时消息真正投递的时间
/**
* 获取延时消息真正投递的时间
* @param delayLevel 延时等级
* @param storeTimestamp 存到延时队列里面的时间
* @return
*/
public long computeDeliverTimestamp(final int delayLevel, final long storeTimestamp) {
// 根据延时等级获取延时时间
Long time = this.delayLevelTable.get(delayLevel);
if (time != null) {
// 存储到延时队列的时间 + 延时时间
return time + storeTimestamp;
}
// 如果没有设置延时等级,那么默认 1s 后投递到真实的 topic 队列中
return storeTimestamp + 1000;
}
这个方法就是用于计算延时消息真正的投递时间,可以看到参数是 storeTimestamp,这个就是消息的存储时间,在上一篇文章 CommitLog 调用 asyncPutMessage 添加消息的时候设置的。因此真实投递时间 = 消息存储时间 + 延时时间。
4.4 correctDeliverTimestamp 校准投递时间
/**
* 校准投递时间, deliverTimestamp 是延时消息要投递到真实 topic 的时间, now 是当前时间
* @return
*/
private long correctDeliverTimestamp(final long now, final long deliverTimestamp) {
// 投递的时间
long result = deliverTimestamp;
// 当前时间 + 延时时间
long maxTimestamp = now + ScheduleMessageService.this.delayLevelTable.get(this.delayLevel);
if (deliverTimestamp > maxTimestamp) {
// 如果投递时间比当前时间 + 延时时间要大, 返回结果就是当前时间
result = now;
}
// 这里确保投递时间比当前时间 + 延时时间要小, 因为投递时间的计算方式是: 消息存储时间 + 延时时间, 按我的理解正常都会返回
// deliverTimestamp, 并不会 > maxTimestamp, 当然如果真的出现这种情况返回 now, 外层会立刻投递到真实队列中
return result;
}
deliverTimestamp 是真实投递的时间,now 是当前时间,这里面的逻辑如下。

正常是不会出现上面 deliverTimestamp > maxTimestamp 的,因为当前时间肯定比存储时间要大,这种情况可能就只有消息存错了才会发生?比如延时等级为 2 的写到了延时等级为 3 的里面,这时候直接返回 now,上层 countdown 算出来的就是 0 了,马上开始投递。
4.5 syncDeliver 同步投递消息
/**
* 同步投递消息
* @param msgInner 要投递的消息对象
* @param msgId 消息 ID
* @param offset 消息逻辑偏移量
* @param offsetPy 消息物理偏移量
* @param sizePy 消息大小
* @return
*/
private boolean syncDeliver(MessageExtBrokerInner msgInner, String msgId, long offset, long offsetPy,
int sizePy) {
// 投递消息, 实际上就是调用 MessageStore#putMessage 方法将消息写入 CommitLog
PutResultProcess resultProcess = deliverMessage(msgInner, msgId, offset, offsetPy, sizePy, false);
// 同步等待投递结果
PutMessageResult result = resultProcess.get();
// 是否添加成功
boolean sendStatus = result != null && result.getPutMessageStatus() == PutMessageStatus.PUT_OK;
if (sendStatus) {
// 如果投递成功, 那么就更新下偏移量为处理成功的延时消息 offset + 1
// 实际上这里更新的 getNextOffset 就来自参数的 offset + 1, 但是在外层传入的时候应该传的是当前消息的偏移量而不是最开始的 offset
// 所以这里是有一个 bug 的, 4.9.4 已修复
ScheduleMessageService.this.updateOffset(this.delayLevel, resultProcess.getNextOffset());
}
return sendStatus;
}
同步投递消息分为两步,首先就是调用 MessageStore#putMessage 方法将消息写入 CommitLog,当消息写入成功,那么就更新下偏移量为处理成功的延时消息 offset + 1。但是这里外层传入的 offset 是这个任务 DeliverDelayedMessageTimerTask 的 offset,这种情况下直接 updateOffset 会导致所有消息写入的都是一个偏移量。
可以看下 deliverMessage 这个方法,这个方法就是直接将消息写入 CommitLog,这个 writeMessageStore 就是在创建出 ScheduleMessageService 时设置进来的,就是 DefaultMessageStore。
/**
* 消息投递
* @param msgInner 要投递的消息对象
* @param msgId 消息 ID
* @param offset 消息逻辑偏移量
* @param offsetPy 消息物理偏移量
* @param sizePy 消息大小
* @param autoResend 消息投递出异常的时候是否自动重投递, 同步投递是 false, 异步投递是 true
* @return
*/
private PutResultProcess deliverMessage(MessageExtBrokerInner msgInner, String msgId, long offset,
long offsetPy, int sizePy, boolean autoResend) {
// 异步添加消息到 CommitLog 中
CompletableFuture<PutMessageResult> future =
ScheduleMessageService.this.writeMessageStore.asyncPutMessage(msgInner);
// 返回 PutResultProcess 对象, 这个对象中对返回结果做了包装
return new PutResultProcess()
.setTopic(msgInner.getTopic())
.setDelayLevel(this.delayLevel)
.setOffset(offset)
.setPhysicOffset(offsetPy)
.setPhysicSize(sizePy)
.setMsgId(msgId)
.setAutoResend(autoResend)
// 设置返回结果
.setFuture(future)
// 在这里面设置 future 的返回结果处理逻辑
.thenProcess();
}
然后返回的 PutResultProcess 设置的 offset 就是参数传过来的 offset,syncDeliver 里面调用 deliverMessage 写入 CommitLog 之后返回 PutResultProcess,getNextOffset 其实就是 offset + 1 了,但是外层传入的 offset 是写死的,应该要传入 nextOffset 才对,在 4.9.4 已修复。
4.6 messageTimeup 还原延时消息为原始消息
这里方法就是将延时消息还原为原始的消息,主要还是 tagsCode、queueId、topic 这几个参数,下面可以直接看注释,写的比较清楚。
/**
* 还原延时消息为原始的消息
* @param msgExt 延时消息
* @return 还原后的真实消息
*/
private MessageExtBrokerInner messageTimeup(MessageExt msgExt) {
// 创建 MessageExtBrokerInner
MessageExtBrokerInner msgInner = new MessageExtBrokerInner();
// 设置消息体
msgInner.setBody(msgExt.getBody());
msgInner.setFlag(msgExt.getFlag());
// 设置消息属性
MessageAccessor.setProperties(msgInner, msgExt.getProperties());
// topic 过滤类型
TopicFilterType topicFilterType = MessageExt.parseTopicFilterType(msgInner.getSysFlag());
// 实际上就是直接对 TAGS 属性求 hashCode
long tagsCodeValue =
MessageExtBrokerInner.tagsString2tagsCode(topicFilterType, msgInner.getTags());
msgInner.setTagsCode(tagsCodeValue);
msgInner.setPropertiesString(MessageDecoder.messageProperties2String(msgExt.getProperties()));
// 设置一些基本属性
msgInner.setSysFlag(msgExt.getSysFlag());
msgInner.setBornTimestamp(msgExt.getBornTimestamp());
msgInner.setBornHost(msgExt.getBornHost());
msgInner.setStoreHost(msgExt.getStoreHost());
msgInner.setReconsumeTimes(msgExt.getReconsumeTimes());
// 不需要同步等待存储到 CommitLog 后才返回
msgInner.setWaitStoreMsgOK(false);
// 清除 PROPERTY_DELAY_TIME_LEVEL 这个属性, 已经是真实消息了, 不需要延时等级了
MessageAccessor.clearProperty(msgInner, MessageConst.PROPERTY_DELAY_TIME_LEVEL);
// 从 PROPERTY_REAL_TOPIC 中获取真实的 topic
msgInner.setTopic(msgInner.getProperty(MessageConst.PROPERTY_REAL_TOPIC));
// 从 PROPERTY_REAL_QUEUE_ID 中获取真实的队列 id
String queueIdStr = msgInner.getProperty(MessageConst.PROPERTY_REAL_QUEUE_ID);
int queueId = Integer.parseInt(queueIdStr);
msgInner.setQueueId(queueId);
// 返回消息
return msgInner;
}
5. HandlePutResultTask 异步处理投递任务的结果
上面第四小节是处理同步延时投递任务的逻辑,有同步就肯定有异步,至于说最后这个任务是通过哪种方式投递的,需要看 enableAsyncDeliver 这个参数,如果是 true,那么在 executeOnTimeup 中就用的异步投递,下面看下有关异步投递的逻辑。
5.1 启动初始化
首先就是 ScheduleMessageService 的时候如果 enableAsyncDeliver = true,这个配置项可以通过 broker.config 里面的 enableScheduleAsyncDeliver 设置为 true,可以看下构造器。
public ScheduleMessageService(final DefaultMessageStore defaultMessageStore) {
this.defaultMessageStore = defaultMessageStore;
this.writeMessageStore = defaultMessageStore;
if (defaultMessageStore != null) {
this.enableAsyncDeliver = defaultMessageStore.getMessageStoreConfig().isEnableScheduleAsyncDeliver();
}
}
然后通过 start 方法启动的时候首先如果允许异步投递,就会初始化延迟消息异步投递的线程池 handleExecutorService。
if (this.enableAsyncDeliver) {
// 如果可以异步投递, 初始化延迟消息异步投递的线程池, 默认不支持
this.handleExecutorService = new ScheduledThreadPoolExecutor(this.maxDelayLevel, new ThreadFactoryImpl("ScheduleMessageExecutorHandleThread_"));
}
然后每一个延时等级构建一个调度任务,初始化是 1s 后执行。
// 为该等级的延迟队列构建一个调度任务,初始化是 1s 后执行
if (timeDelay != null) {
if (this.enableAsyncDeliver) {
this.handleExecutorService.schedule(new HandlePutResultTask(level), FIRST_DELAY_TIME, TimeUnit.MILLISECONDS);
}
}
然后来看下下面的 HandlePutResultTask 的 run 方法,注意上面创建 HandlePutResultTask 的时候
5.2 消息投递到异步队列
在 DeliverDelayedMessageTimerTask 调用 executeOnTimeup 处理延时消息时会判断如果打开了异步投递开关,就走 asyncDeliver 去异步投递。
// 消息投递
boolean deliverSuc;
if (ScheduleMessageService.this.enableAsyncDeliver) {
// 异步投递消息, 默认是不支持的
deliverSuc = this.asyncDeliver(msgInner, msgExt.getMsgId(), offset, offsetPy, sizePy);
} else {
// 同步投递消息, offset 应该是传的 nextOffset 才对, 不应该把最开始的 offset 返回, 4.9.4 已修复
deliverSuc = this.syncDeliver(msgInner, msgExt.getMsgId(), offset, offsetPy, sizePy);
}
下面来看下 asyncDeliver 这个方法的逻辑。
/**
* 异步投递消息
* @param msgInner 要投递的消息对象
* @param msgId 消息 ID
* @param offset 消息逻辑偏移量
* @param offsetPy 消息物理偏移量
* @param sizePy 消息大小
* @return
*/
private boolean asyncDeliver(MessageExtBrokerInner msgInner, String msgId, long offset, long offsetPy,
int sizePy) {
// 获取这个延时等级下的异步投递队列
Queue<PutResultProcess> processesQueue = ScheduleMessageService.this.deliverPendingTable.get(this.delayLevel);
// 进行流控, 看看现在有多少个投递请求在等待处理
int currentPendingNum = processesQueue.size();
// 默认最多往投递队列里面放 2000 个请求结果, 考虑到线程处理不过来的原因
int maxPendingLimit = ScheduleMessageService.this.defaultMessageStore.getMessageStoreConfig()
.getScheduleAsyncDeliverMaxPendingLimit();
if (currentPendingNum > maxPendingLimit) {
// 流控失败, 直接返回 false
log.warn("Asynchronous deliver triggers flow control, " +
"currentPendingNum={}, maxPendingLimit={}", currentPendingNum, maxPendingLimit);
return false;
}
// 通过流控, 获取这个队列下面的第一条消息, 如果说第一条消息的重发次数超过了 3 次了, 就先不投递这条消息了, 延迟 100ms 之后再投递
PutResultProcess firstProcess = processesQueue.peek();
if (firstProcess != null && firstProcess.need2Blocked()) {
log.warn("Asynchronous deliver block. info={}", firstProcess.toString());
// 返回 false 之后延迟 100ms 再次处理
return false;
}
// 对于异步投递, autoResend 会设置成 true, 这就代表如果消息投递失败会自动重新投递
PutResultProcess resultProcess = deliverMessage(msgInner, msgId, offset, offsetPy, sizePy, true);
// 将结果添加到队列中等待 handleExecutorService 处理
processesQueue.add(resultProcess);
return true;
}
首先获取这个延时等级下的异步投递队列,然后进行流控,看看现在有多少个投递请求在等待处理,默认最多往投递队列里面放 2000 个请求结果,请求太多线程可能处理不过来,如果达到阈值,直接返回 false,代表流控失败,然后就返回 false,上层判断如果是 false,就会延时 100ms 之后再次提交一个 DeliverDelayedMessageTimerTask 任务再去处理。
如果流控判断成功了,获取这个队列下面的第一条消息,如果说第一条消息的重发次数超过了 3 次了,就先不投递这条消息了,延迟 100ms 之后再投递,这里的重试意思是如果投递过程中发生了异常就会判断是否需要重新投递,如果需要就会重试,然后投递次数 + 1。
如果上面的校验都没有问题,调用 deliverMessage 进行消息投递,然后将返回结果存到 processesQueue 队列中。这里就是异步投递了,可以看到对于投递结果并没有使用 get 方法去阻塞等待,而是丢到了 processesQueue 队列,那么这些请求结果是在哪处理的呢?就是 HandlePutResultTask 的 run 方法。
5.3 处理异步投递的结果
在看 run 方法的处理逻辑之前,我们先来看下 PutResultProcess 的 status 是怎么设置的,这个变量就是投递的结果,run 方法就是根据 status 来走不同的逻辑。
在 deliverMessage 方法通过 thenProcess 设置 future 的返回结果处理逻辑,status 默认是 ProcessStatus.RUNNING,表示消息正在运行中。
/**
* 这里是处理返回结果的逻辑
* @return
*/
public PutResultProcess thenProcess() {
this.future.thenAccept(result -> {
// 没有出异常, 就走 handleResult 方法去处理投递结果
this.handleResult(result);
});
this.future.exceptionally(e -> {
log.error("ScheduleMessageService put message exceptionally, info: {}",
PutResultProcess.this.toString(), e);
// 异常处理
onException();
return null;
});
return this;
}
thenProcess 就是设置处理返回结果的逻辑,如果投递过程没有出异常,就走 handleResult 方法去处理投递结果,如果出异常就用 onException 去处理结果,先来看下 handleResult 的逻辑。
/**
* 处理投递结果
* @param result
*/
private void handleResult(PutMessageResult result) {
if (result != null && result.getPutMessageStatus() == PutMessageStatus.PUT_OK) {
// 如果投递结果是 PUT_OK, 设置 status 状态并进行数据统计
onSuccess(result);
} else {
log.warn("ScheduleMessageService put message failed. info: {}.", result);
// 消息投递失败, 通过 onException 进行消息重投递或者直接跳过不管
onException();
}
}
这里会判断如果消息添加成功,调用 onSuccess 方法,否则就是 onException 处理异常情况,下面看下 onSuccess。
/**
* 成功投递之后调用这个 onSuccess 方法
* @param result
*/
public void onSuccess(PutMessageResult result) {
// 设置投递状态为 SUCCESS
this.status = ProcessStatus.SUCCESS;
// 进行消息统计
if (ScheduleMessageService.this.defaultMessageStore.getMessageStoreConfig().isEnableScheduleMessageStats()) {
ScheduleMessageService.this.defaultMessageStore.getBrokerStatsManager().incQueueGetNums(MixAll.SCHEDULE_CONSUMER_GROUP, TopicValidator.RMQ_SYS_SCHEDULE_TOPIC, delayLevel - 1, result.getAppendMessageResult().getMsgNum());
ScheduleMessageService.this.defaultMessageStore.getBrokerStatsManager().incQueueGetSize(MixAll.SCHEDULE_CONSUMER_GROUP, TopicValidator.RMQ_SYS_SCHEDULE_TOPIC, delayLevel - 1, result.getAppendMessageResult().getWroteBytes());
ScheduleMessageService.this.defaultMessageStore.getBrokerStatsManager().incGroupGetNums(MixAll.SCHEDULE_CONSUMER_GROUP, TopicValidator.RMQ_SYS_SCHEDULE_TOPIC, result.getAppendMessageResult().getMsgNum());
ScheduleMessageService.this.defaultMessageStore.getBrokerStatsManager().incGroupGetSize(MixAll.SCHEDULE_CONSUMER_GROUP, TopicValidator.RMQ_SYS_SCHEDULE_TOPIC, result.getAppendMessageResult().getWroteBytes());
ScheduleMessageService.this.defaultMessageStore.getBrokerStatsManager().incTopicPutNums(this.topic, result.getAppendMessageResult().getMsgNum(), 1);
ScheduleMessageService.this.defaultMessageStore.getBrokerStatsManager().incTopicPutSize(this.topic, result.getAppendMessageResult().getWroteBytes());
ScheduleMessageService.this.defaultMessageStore.getBrokerStatsManager().incBrokerPutNums(result.getAppendMessageResult().getMsgNum());
}
}
可以看到这里面都是统计一些数据,比如当前服务往 broker 添加了多少条消息,往对应的 topic 添加了多少消息 … 最重要的是将结果的 status 设置为 ProcessStatus.SUCCESS,然后再来看下 onException 的处理逻辑。
/**
* 如果消息投递过程中出异常了, 重新发送一次
*/
public void onException() {
log.warn("ScheduleMessageService onException, info: {}", this.toString());
if (this.autoResend) {
// 如果是自动重发送, 直接调用 resend 方法去重新投递消息
this.resend();
} else {
// 如果不是自动重发送, 设置消息的处理状态是 SKIP, 对于 SKIP 的处理就是将当前的 PutResultProcess 从阻塞队列中移除掉
this.status = ProcessStatus.SKIP;
}
}
如果说发生异常了,同时 autoResend 是 true,那么就直接调用 resend 方法去重新投递消息。如果不是自动重发送,设置消息的处理状态是 SKIP,run 方法中对于 SKIP 的处理就是将当前的 PutResultProcess 从阻塞队列中移除掉。
autoResend 是在 deliverMessage 中设置进去的,同步投递 syncDeliver 设置的是 false,异步投递 asyncDeliver 设置的是 true,所以可以看到同步投递的情况下如果投递失败是不会重投的,而异步投递则会重投,下面是消息重投的逻辑。
/**
* 消息重投
*/
private void resend() {
log.info("Resend message, info: {}", this.toString());
// Gradually increase the resend interval.
try {
// 设置消息重投的次数 + 1, 然后先睡眠一段时间
Thread.sleep(Math.min(this.resendCount++ * 100, 60 * 1000));
} catch (InterruptedException e) {
// 被中断就抛异常
e.printStackTrace();
}
try {
// 重新截取出消息
MessageExt msgExt = ScheduleMessageService.this.defaultMessageStore.lookMessageByOffset(this.physicOffset, this.physicSize);
if (msgExt == null) {
log.warn("ScheduleMessageService resend not found message. info: {}", this.toString());
// 找不到消息, 如果消息重投次数超过 6 次了, 设置状态为 SKIP, 后续就不再处理这个重投结果
this.status = need2Skip() ? ProcessStatus.SKIP : ProcessStatus.EXCEPTION;
return;
}
// 还原延时消息为原始的消息
MessageExtBrokerInner msgInner = ScheduleMessageService.this.messageTimeup(msgExt);
// 消息重投
PutMessageResult result = ScheduleMessageService.this.writeMessageStore.putMessage(msgInner);
// 处理投递结果
this.handleResult(result);
if (result != null && result.getPutMessageStatus() == PutMessageStatus.PUT_OK) {
log.info("Resend message success, info: {}", this.toString());
}
} catch (Exception e) {
this.status = ProcessStatus.EXCEPTION;
log.error("Resend message error, info: {}", this.toString(), e);
}
}
消息重投就是先睡眠 resendCount * 100 ms,然后再去获取出对应的延时消息,获取不到消息,判断下重投次数如果已经到达 6 次了,就设置本次投递结果状态为 SKIP,否则就是 EXCEPTION。
由于当前获取到的还是延时消息,就需要通过 messageTimeup 还原为原始的消息,之后调用 putMessage 同步添加消息到 CommitLog 中,不过这里虽然是同步阻塞添加,其实 putMessage#waitForPutResult 方法最多也只会阻塞 10s 来添加,因此这里并不是一定会无限阻塞下去,不过 10s 也挺久的了,除非出什么问题,不然添加消息不可能要这么久的,毕竟消息写入是写入到 MappedByteBuffer 中,性能比较高。获取到投递结果之后,调用 handleResult 再次去处理结果。
好了,上面的 onException 的逻辑也看了,到这里大家也大概清楚了 PutResultProcess.status 是怎么设置的,下面我们可以看下最后 run 方法的处理逻辑。
@Override
public void run() {
// 从 deliverPendingTable 集合中获取这个延时等级下面的投递结果集合
LinkedBlockingQueue<PutResultProcess> pendingQueue =
ScheduleMessageService.this.deliverPendingTable.get(this.delayLevel);
PutResultProcess putResultProcess;
while ((putResultProcess = pendingQueue.peek()) != null) {
try {
switch (putResultProcess.getStatus()) {
case SUCCESS:
// 投递成功, 更新延时等级的消费偏移量 offsetTable
ScheduleMessageService.this.updateOffset(this.delayLevel, putResultProcess.getNextOffset());
// 删除掉这个 PutResultProcess
pendingQueue.remove();
break;
case RUNNING:
// 默认状态就是 RUNNING, 如果投递没有完成前就会一直都是这个状态
break;
case EXCEPTION:
// 出异常, 如果服务不是启动状态, 直接退出
if (!isStarted()) {
log.warn("HandlePutResultTask shutdown, info={}", putResultProcess.toString());
return;
}
log.warn("putResultProcess error, info={}", putResultProcess.toString());
// 异步投递在这里面的逻辑就是通过 resend 重投消息
putResultProcess.onException();
break;
case SKIP:
// 如果是 SKIP 状态就直接跳过这条消息的结果处理
log.warn("putResultProcess skip, info={}", putResultProcess.toString());
pendingQueue.remove();
break;
}
} catch (Exception e) {
log.error("HandlePutResultTask exception. info={}", putResultProcess.toString(), e);
putResultProcess.onException();
}
}
if (isStarted()) {
// 没有消息投递或者说已经处理完所有投递结果了, 或者说消息还在投递中, 延迟 10ms 再次创建一个 HandlePutResultTask 去执行
ScheduleMessageService.this.handleExecutorService
.schedule(new HandlePutResultTask(this.delayLevel), DELAY_FOR_A_SLEEP, TimeUnit.MILLISECONDS);
}
}
首先是 RUNNING,如果是这种状态,说明消息还在投递中,还没投递完成,这种情况下直接 break 退出,然后延时 10ms 之后再次创建一个 HandlePutResultTask 去处理投递的结果。
然后是 SUCCESS,这个状态就表示原始消息投递成功了,本条延时消息处理完成,更新本地 offsetTable 对应的偏移量,然后删除掉这个 PutResultProcess,继续处理队列下一个 PutResultProcess。
其次是 EXCEPTION,在消息重投 resend 方法处理时,如果 获取不到消息且重投次数超过了 6 次,状态就是 SKIP,异步投递正常会不断重试,直到上面说的终止条件,因为异步投递设置了 autoResend 是 true,也就是会自动重投,所以投递过程发生异常什么的都会不断去重投,也就是调用 resend 方法去处理。
最后就是 SKIP,这个状态有两种情况,上面也说了一种就是重投过程中发现消息获取不到了,重试次数超过了 6 次,这种情况下就会将状态设置为 SKIP,还有一种就是 onException 中如果没有设置重发送(autoResend = false),就会将消息的处理状态设置为 SKIP,对于 SKIP 的处理就是将当前的 PutResultProcess 从阻塞队列中移除掉,不处理这条延时消息的投递结果。
6. 延时消息处理总流程
上面我们分析了延时消息处理服务的全部代码,下面再从整体视角来看下一条延时消息的发送流程。
首先创建出生产者,设置消息延时等级,然后发送消息。

写入 CommitLog 的时候会通过延时等级计算出队列 ID,然后将这个队列 ID 和消息需要发送的原始 topic 设置在属性 REAL_TOPIC 和 REAL_QID 中。

然后消息重放服务开始读取 CommitLog 里面的消息进行重放,并且在重放前调用 checkMessageAndReturnSize 去检查消息是否合法,在这个方法会将延时消息的 tagsCode 设置为实际的投递时间,也就是 消息存储时间 + 延时时间。


ConsumeQueue 在写入这条延时消息的索引时会写入 tagsCode,因此延时服务 ScheduleMessageService 获取 SCHEDULE_TOPIC_XXXX 下面的 ConsumeQueue 时能够直接用 tagsCode 来计算有没有到期。

接下来就是上面 1 - 5 小节的内容了,延时消息服务 ScheduleMessageService 启动任务定时扫描有没有 ConsumeQueue 索引需要处理,如果存在就判断延时消息是否到期,如果是就投递到真实队列里面,投递之前通过 messageTimeup 还原原始消息,重点是求出原来的 tagsCode,然后重新设置真实的队列 id 和 topic。

最后延时消息发送到真实的 topic 中,由对应的消费者拉取消息下来消费,消息的整体流程就结束了。

7. 小结
好了,这篇文章也是承接上两篇消息消费的流程,分析了延时消息的发送流程以及核心源码,对于延时消息,延时等级和 queueId 对应,一开始都会将消息发送到 SCHEDULE_TOPIC_XXX 这个 topic,后面延时服务在判断出消息到期之后会把消息发送回真实的 topic,要注意就是对于延时消息 tagsCode 会在构建 ConsumeQueue 索引的时候替换成这条延时消息的到期时间,然后 ScheduleMessageService 才可以在遍历 ConsumeQueue 索引时直接判断有没有到期,其真实的 topic 和 queueId 存在 REAL_TOPIC 和 REAL_QID 两个变量中。最后经过上面源码的分析,也能知道延时消息的到期时间计算是通过存储时间 + 延时时间来算出的,也就是 storeTimestamp + deliverTime,这个消息存储时间是在写入 CommitLog 会设置。
更多推荐




所有评论(0)