RocketMQ 延迟消息
RocketMQ 消费者启动流程
什么是延迟消息
RocketMQ 延迟消息是指,生产者发送消息给消费者消息,消费者需要等待一段时间后才能消费到。
使用场景
用户下单之后,15分钟未支付,对支付账单进行提醒或者关单处理。
RocketMQ 开源版本的消息不支持任意时间精度,只支持5s 10s 1m
等等。
Broker 如何处理延迟消息
消息投递如下:
- 生产者发送一个延迟消息到一个
topic
中 - Broker 判断是个延迟消息后,将消息暂存
- Broker 通过延迟服务, 先检查消息是否过期,如果到期将消息投递到目标
topic
- 消费者消费
topic
中的投递延迟消息。
开源RocketMQ 的消息不支持任意精度,默认支持 18个 level:
messageDelayLevel=1s 5s 10s 30s 1m 2m 3m 4m 5m 6m 7m 8m 9m 10m 20m 30m 1h 2h
Broker 在启动的时候,会创建一个内部 topic:“SCHEDULE_TOPIC_XXXX” 根据延迟 level 数量,创建对应数量的 队列。 也就是说 18 level 对应了18 个队列。
具体可以在 代码TopicConfigManager.java 中 看到:
private static final int SCHEDULE_TOPIC_QUEUE_NUM = 18;
要注意的是,Broker 一般是集群模式
部署,也就是说,每个Broker 都会有18个队列。
TopicConfigManager#TopicConfigManager(BrokerController brokerController)
生产者消息延迟发送
代码示例如下:
Message msg=new Message();
msg.setTopic("TopicA");
msg.setTags("Tag");
msg.setBody("this is a delay message".getBytes());
//设置延迟level为5,对应延迟1分钟
msg.setDelayTimeLevel(5);
producer.send(msg);
Broker 存储延迟消息
上一篇文章已经谈到,Broker
收到消费者消息后,会进行消息存储,然后再转发到消费队列(ConsumerQueue),然后再推给消费者。
其实一旦消息转发到
存储延迟消息的流程也类似
- 确定延迟消息投递到topic 哪个队列。存储生产者写入的消息时,将消息转发到 ConsumeQueue 中,消费者就能消费到。 延迟消息不能立即消息到,于是将 topic 名称修改为
SCHEDULE_TOPIC_XXX
,并根据延迟消息级别,确定投递到哪个队列上。同时还会将原来消息要发送到的目标 topic 和队列记录投递到哪个队列。
代码在CommitLog#asyncPutMessage 中
设置延迟消息的投递队列信息代码如下:
// Delay Delivery
if (msg.getDelayTimeLevel() > 0) {
// 如果设置的级别超过了最大级别,重置延迟级别
if (msg.getDelayTimeLevel() > this.defaultMessageStore.getScheduleMessageService().getMaxDelayLevel()) {
msg.setDelayTimeLevel(this.defaultMessageStore.getScheduleMessageService().getMaxDelayLevel());
}
// 计算延迟消息应该投递到 SCHEDULE_TOPIC_XXXX 到哪个队列。
topic = TopicValidator.RMQ_SYS_SCHEDULE_TOPIC;
int queueId = ScheduleMessageService.delayLevel2QueueId(msg.getDelayTimeLevel());
// Backup real topic, queueId
// 记录原始 topic ,queueid,方便后期投递到目标 topic
MessageAccessor.putProperty(msg, MessageConst.PROPERTY_REAL_TOPIC, msg.getTopic());
MessageAccessor.putProperty(msg, MessageConst.PROPERTY_REAL_QUEUE_ID, String.valueOf(msg.getQueueId()));
msg.setPropertiesString(MessageDecoder.messageProperties2String(msg.getProperties()));
// 更新消息投递目标为 SCHEDULE_TOPIC_XXX,queueId
msg.setTopic(topic);
msg.setQueueId(queueId);
}
消息转发
消息转发过程其实中会对延迟消息做一些特殊处理
CommitLog中的消息转发到CosumeQueue中是异步进行的。在转发过程中,会对延迟消息进行特殊处理,主要是计算这条延迟消息需要在什么时候进行投递。