RocketMQ源码 — 九、 RocketMQ延时消息

上一节消息重试里面提到了重试的消息可以被延时消费,其实除此之外,用户发送的消息也可以指定延时时间(更准确的说是延时等级),然后在指定延时时间之后投递消息,然后被consumer消费。阿里云的ons还支持定时消息,而且延时消息是直接指定延时时间,其实阿里云的延时消息也是定时消息的另一种表述方式,都是通过设置消息被投递的时间来实现的,但是Apache RocketMQ在版本4.2.0中尚不支持指定时间的延时,只能通过配置延时等级和延时等级对应的时间来实现延时。

一个延时消息被发出到消费成功经历以下几个过程:

  1. 设置消息的延时级别delayLevel
  2. producer发送消息
  3. broker收到消息在准备将消息写入存储的时候,判断是延时消息则更改Message的topic为延时消息队列的topic,也就是将消息投递到延时消息队列
  4. 有定时任务从延时队列中读取消息,拿到消息后判断是否达到延时时间,如果到了则修改topic为原始topic。并将消息投递到原始topic的队列
  5. consumer像消费其他消息一样从broker拉取消息进行消费

注意:批量消息是不支持延时消息的

tips:下文中说到的延时队列可以理解为一个ConsumeQueue

producer发送延时消息

在producer中发送消息的时候,设置Message的delayLevel

1// org.apache.rocketmq.common.message.Message 2public void setDelayTimeLevel(int level) { 3 this.putProperty(MessageConst.PROPERTY_DELAY_TIME_LEVEL, String.valueOf(level)); 4}

调用上面的方法设置延时等级的时候,会向message添加"DELAY"属性,后面broker处理延时消息就是依赖该属性进行特别的处理。

接下来发送消息的流程和正常发送消息的流程基本一致,只是会将该消息标记为延时消息类型

1// org.apache.rocketmq.client.impl.producer.DefaultMQProducerImpl#sendKernelImpl 2if (msg.getProperty("__STARTDELIVERTIME") != null || msg.getProperty(MessageConst.PROPERTY_DELAY_TIME_LEVEL) != null) { 3 context.setMsgType(MessageType.Delay_Msg); 4}

broker处理延时消息

broker收到延时消息和正常消息在前置的处理流程是一致的,对于延时消息的特殊处理体现在将消息写入存储(内存或文件)的时候

1// org.apache.rocketmq.store.CommitLog#putMessage 2public PutMessageResult putMessage(final MessageExtBrokerInner msg) { 3 // 省略中间代码... 4 StoreStatsService storeStatsService = this.defaultMessageStore.getStoreStatsService(); 5 6 // 拿到原始topic和对应的queueId 7 String topic = msg.getTopic(); 8 int queueId = msg.getQueueId(); 9 10 final int tranType = MessageSysFlag.getTransactionValue(msg.getSysFlag()); 11 // 非事务消息和事务的commit消息才会进一步判断delayLevel 12 if (tranType == MessageSysFlag.TRANSACTION_NOT_TYPE 13 || tranType == MessageSysFlag.TRANSACTION_COMMIT_TYPE) { 14 // Delay Delivery 15 if (msg.getDelayTimeLevel() > 0) { 16 // 纠正设置过大的level,就是delayLevel设置都大于延时时间等级的最大级 17 if (msg.getDelayTimeLevel() > this.defaultMessageStore.getScheduleMessageService().getMaxDelayLevel()) { 18 msg.setDelayTimeLevel(this.defaultMessageStore.getScheduleMessageService().getMaxDelayLevel()); 19 } 20 21 // 设置为延时队列的topic 22 topic = ScheduleMessageService.SCHEDULE_TOPIC; 23 // 每一个延时等级一个queue,queueId = delayLevel - 1 24 queueId = ScheduleMessageService.delayLevel2QueueId(msg.getDelayTimeLevel()); 25 26 // Backup real topic, queueId 27 // 备份原始的topic和queueId 28 MessageAccessor.putProperty(msg, MessageConst.PROPERTY_REAL_TOPIC, msg.getTopic()); 29 MessageAccessor.putProperty(msg, MessageConst.PROPERTY_REAL_QUEUE_ID, String.valueOf(msg.getQueueId())); 30 // 更新properties 31 msg.setPropertiesString(MessageDecoder.messageProperties2String(msg.getProperties())); 32 33 msg.setTopic(topic); 34 msg.setQueueId(queueId); 35 } 36 } 37 // 省略中间代码... 38}

上面的SCHEDULE_TOPIC是:

public static final String SCHEDULE_TOPIC = "SCHEDULE_TOPIC_XXXX";

这个topic是一个特殊的topic,和正常的topic不同的地方是:

  1. 不会创建TopicConfig,因为也不需要consumer直接消费这个topic下的消息
  2. 不会将topic注册到namesrv
  3. 这个topic的队列个数和延时等级的个数是相同的

后面消息写入的过程和普通的又是一致的。

上面将消息写入延时队列中了,接下来就是处理延时队列中的消息,然后重新发送回原始topic的队列中。

在此之前先说明下至今还有疑问的一个个概念——delayLevel。这个概念和我们接下要需要用到的的类org.apache.rocketmq.store.schedule.ScheduleMessageService有关,这个类的字段delayLevelTable里面保存了具体的延时等级

private final ConcurrentMap<Integer /* level */, Long/* delay timeMillis */> delayLevelTable = new ConcurrentHashMap<Integer, Long>(32);

看下这个字段的初始化过程

1// org.apache.rocketmq.store.schedule.ScheduleMessageService#parseDelayLevel 2public boolean parseDelayLevel() { 3 HashMap<String, Long> timeUnitTable = new HashMap<String, Long>(); 4 // 每个延时等级延时时间的单位对应的ms数 5 timeUnitTable.put("s", 1000L); 6 timeUnitTable.put("m", 1000L * 60); 7 timeUnitTable.put("h", 1000L * 60 * 60); 8 timeUnitTable.put("d", 1000L * 60 * 60 * 24); 9 10 // 延时等级在MessageStoreConfig中配置 11 // private String messageDelayLevel = "1s 5s 10s 30s 1m 2m 3m 4m 5m 6m 7m 8m 9m 10m 20m 30m 1h 2h"; 12 String levelString = this.defaultMessageStore.getMessageStoreConfig().getMessageDelayLevel(); 13 try { 14 // 根据空格将配置分隔出每个等级 15 String[] levelArray = levelString.split(" "); 16 for (int i = 0; i < levelArray.length; i++) { 17 String value = levelArray[i]; 18 String ch = value.substring(value.length() - 1); 19 // 时间单位对应的ms数 20 Long tu = timeUnitTable.get(ch); 21 22 // 延时等级从1开始 23 int level = i + 1; 24 if (level > this.maxDelayLevel) { 25 // 找出最大的延时等级 26 this.maxDelayLevel = level; 27 } 28 long num = Long.parseLong(value.substring(0, value.length() - 1)); 29 long delayTimeMillis = tu * num; 30 this.delayLevelTable.put(level, delayTimeMillis); 31 // 省略部分代码... 32}

上面这个load方法在broker启动的时候DefaultMessageStore会调用来初始化延时等级。

接下来就应该解决怎么处理延时消息队列中的消息的问题了。处理延时消息的服务是:ScheduleMessageService。

还是broker启动的时候DefaultMessageStore会调用org.apache.rocketmq.store.schedule.ScheduleMessageService#start来启动处理延时消息队列的服务:

1public void start() { 2 3 for (Map.Entry<Integer, Long> entry : this.delayLevelTable.entrySet()) { 4 Integer level = entry.getKey(); 5 Long timeDelay = entry.getValue(); 6 // 记录队列的处理进度 7 Long offset = this.offsetTable.get(level); 8 if (null == offset) { 9 offset = 0L; 10 } 11 12 if (timeDelay != null) { 13 // 每个延时队列启动一个定时任务来处理该队列的延时消息 14 this.timer.schedule(new DeliverDelayedMessageTimerTask(level, offset), FIRST_DELAY_TIME); 15 } 16 } 17 18 this.timer.scheduleAtFixedRate(new TimerTask() { 19 20 @Override 21 public void run() { 22 try { 23 // 持久化offsetTable(保存了每个延时队列对应的处理进度offset) 24 ScheduleMessageService.this.persist(); 25 } catch (Throwable e) { 26 log.error("scheduleAtFixedRate flush exception", e); 27 } 28 } 29 }, 10000, this.defaultMessageStore.getMessageStoreConfig().getFlushDelayOffsetInterval()); 30}

DeliverDelayedMessageTimerTask是一个TimerTask,启动以后不断处理延时队列中的消息,直到出现异常则终止该线程重新启动一个新的TimerTask

1public void executeOnTimeup() { 2 // 找到该延时等级对应的ConsumeQueue 3 ConsumeQueue cq = 4 ScheduleMessageService.this.defaultMessageStore.findConsumeQueue(SCHEDULE_TOPIC, 5 delayLevel2QueueId(delayLevel)); 6 // 记录异常情况下一次启动TimerTask开始处理的offset 7 long failScheduleOffset = offset; 8 9 if (cq != null) { 10 // 找到offset所处的MappedFile中offset后面的buffer 11 SelectMappedBufferResult bufferCQ = cq.getIndexBuffer(this.offset); 12 if (bufferCQ != null) { 13 try { 14 long nextOffset = offset; 15 int i = 0; 16 ConsumeQueueExt.CqExtUnit cqExtUnit = new ConsumeQueueExt.CqExtUnit(); 17 for (; i < bufferCQ.getSize(); i += ConsumeQueue.CQ_STORE_UNIT_SIZE) { 18 // 下面三个字段信息是ConsumeQueue物理存储的信息 19 long offsetPy = bufferCQ.getByteBuffer().getLong(); 20 int sizePy = bufferCQ.getByteBuffer().getInt(); 21 // 注意这个tagCode,不再是普通的tag的hashCode,而是该延时消息到期的时间 22 long tagsCode = bufferCQ.getByteBuffer().getLong(); 23 // 省略中间代码.... 24 long now = System.currentTimeMillis(); 25 // 计算应该投递该消息的时间,如果已经超时则立即投递 26 long deliverTimestamp = this.correctDeliverTimestamp(now, tagsCode); 27 // 计算下一个消息的开始位置,用来寻找下一个消息位置(如果有的话) 28 nextOffset = offset + (i / ConsumeQueue.CQ_STORE_UNIT_SIZE); 29 // 判断延时消息是否到期 30 long countdown = deliverTimestamp - now; 31 32 if (countdown <= 0) { 33 MessageExt msgExt = 34 ScheduleMessageService.this.defaultMessageStore.lookMessageByOffset( 35 offsetPy, sizePy); 36 37 if (msgExt != null) { 38 try { 39 // 将消息恢复到原始消息的格式,恢复topic、queueId、tagCode等,清除属性"DELAY" 40 MessageExtBrokerInner msgInner = this.messageTimeup(msgExt); 41 PutMessageResult putMessageResult = 42 ScheduleMessageService.this.defaultMessageStore 43 .putMessage(msgInner); 44 45 if (putMessageResult != null 46 && putMessageResult.getPutMessageStatus() == PutMessageStatus.PUT_OK) { 47 // 投递成功,处理下一个 48 continue; 49 } else { 50 // XXX: warn and notify me 51 log.error( 52 "ScheduleMessageService, a message time up, but reput it failed, topic: {} msgId {}", 53 msgExt.getTopic(), msgExt.getMsgId()); 54 // 投递失败,结束当前task,重新启动TimerTask,从下一个消息开始处理,也就是说当前消息丢弃 55 // 更新offsetTable中当前队列的offset为下一个消息的offset 56 ScheduleMessageService.this.timer.schedule( 57 new DeliverDelayedMessageTimerTask(this.delayLevel, 58 nextOffset), DELAY_FOR_A_PERIOD); 59 ScheduleMessageService.this.updateOffset(this.delayLevel, 60 nextOffset); 61 return; 62 } 63 } catch (Exception e) { 64 // 重新投递期间出现任何异常,结束当前task,重新启动TimerTask,从当前消息开始重试 65 /* 66 * XXX: warn and notify me 67 */ 68 log.error( 69 "ScheduleMessageService, messageTimeup execute error, drop it. msgExt=" 70 + msgExt + ", nextOffset=" + nextOffset + ",offsetPy=" 71 + offsetPy + ",sizePy=" + sizePy, e); 72 } 73 } 74 } else { 75 ScheduleMessageService.this.timer.schedule( 76 new DeliverDelayedMessageTimerTask(this.delayLevel, nextOffset), 77 countdown); 78 ScheduleMessageService.this.updateOffset(this.delayLevel, nextOffset); 79 return; 80 } 81 } // end of for 82 // 处理完当前MappedFile中的消息后,重新启动TimerTask,从下一个消息开始处理 83 // 更新offsetTable中当前队列的offset为下一个消息的offset 84 nextOffset = offset + (i / ConsumeQueue.CQ_STORE_UNIT_SIZE); 85 ScheduleMessageService.this.timer.schedule(new DeliverDelayedMessageTimerTask( 86 this.delayLevel, nextOffset), DELAY_FOR_A_WHILE); 87 ScheduleMessageService.this.updateOffset(this.delayLevel, nextOffset); 88 return; 89 } finally { 90 91 bufferCQ.release(); 92 } 93 } // end of if (bufferCQ != null) 94 else { 95 // 如果根据offsetTable中的offset没有找到对应的消息(可能被删除了),则按照当前ConsumeQueue的最小offset开始处理 96 long cqMinOffset = cq.getMinOffsetInQueue(); 97 if (offset < cqMinOffset) { 98 failScheduleOffset = cqMinOffset; 99 log.error("schedule CQ offset invalid. offset=" + offset + ", cqMinOffset=" 100 + cqMinOffset + ", queueId=" + cq.getQueueId()); 101 } 102 } 103 } // end of if (cq != null) 104 105 ScheduleMessageService.this.timer.schedule(new DeliverDelayedMessageTimerTask(this.delayLevel, 106 failScheduleOffset), DELAY_FOR_A_WHILE); 107}

对于上面的tagCode做一下特别说明,延时消息的tagCode和普通消息不一样:

  • 延时消息的tagCode:存储的是消息到期的时间
  • 非延时消息的tagCode:tags字符串的hashCode

对延时消息的tagCode的特别处理是在下面这个方法中完成的,也就是在build ConsumeQueue信息的时候

org.apache.rocketmq.store.CommitLog#checkMessageAndReturnSize(java.nio.ByteBuffer, boolean, boolean)

总结

以上就是RocketMQ延时消息的实现方式,上面没有详说的是重试消息的延时是怎么实现的,其实就是在consumer将延时消息发送回broker的时候设置了(用户可以自己设置,如果没有自己设置默认是0)delayLevel,到了broker处理重试消息的时候如果delayLevel是0(也就是说是默认的延时等级)的时候会在原来的基础上加3,后面的处理就和上面说的延时消息一样了,存储的时候将消息投递到延时队列,等待延时到期后再重新投递到原始topic队列中等到consumer消费。

点赞
收藏

评论区

加载中...

相关推荐

MySQL:[Err] 1292 - Incorrect datetime value: ‘0000-00-00 00:00:00‘ for column ‘CREATE_TIME‘ at row 1

文章目录问题用navicat导入数据时,报错:原因这是因为当前的MySQL不支持datetime为0的情况。解决修改sql\mode:sql\mode:SQLMode定义了MySQL应支持的SQL语法、数据校验等,这样可以更容易地在不同的环境中使用MySQL。全局s

Oracle 分组与拼接字符串同时使用

SELECTT.,ROWNUMIDFROM(SELECTT.EMPLID,T.NAME,T.BU,T.REALDEPART,T.FORMATDATE,SUM(T.S0)S0,MAX(UPDATETIME)CREATETIME,LISTAGG(TOCHAR(

MySQL部分从库上面因为大量的临时表tmp_table造成慢查询

背景描述Time:20190124T00:08:14.70572408:00User@Host:@Id:Schema:sentrymetaLast_errno:0Killed:0Query_time:0.315758Lock_

皕杰报表之UUID

​在我们用皕杰报表工具设计填报报表时,如何在新增行里自动增加id呢?能新增整数排序id吗?目前可以在新增行里自动增加id,但只能用uuid函数增加UUID编码,不能新增整数排序id。uuid函数说明:获取一个UUID,可以在填报表中用来创建数据ID语法:uuid()或uuid(sep)参数说明:sep布尔值,生成的uuid中是否包含分隔符'',缺省为

手写Java HashMap源码

HashMap的使用教程HashMap的使用教程HashMap的使用教程HashMap的使用教程HashMap的使用教程22

2020年前端实用代码段,为你的工作保驾护航

有空的时候,自己总结了几个代码段,在开发中也经常使用,谢谢。1、使用解构获取json数据let jsonData  id: 1,status: "OK",data: 'a', 'b';let  id, status, data: number   jsonData;console.log(id, status, number )