news 2026/8/25 15:44:37

rocketMQ proxy 延迟队列

作者头像

张小明

前端开发工程师

1.2k 24
文章封面图
rocketMQ proxy 延迟队列

Proxy 本身不做延迟队列的存储和调度 (那是 Broker 端 ScheduleMessageService / TimerMessageStore 的职责),Proxy 只负责 延迟消息的「发送端属性填充 + 延迟级别换算」以及消费端识别 。相关实现集中在 gRPC 发送链路和配置里。

  1. 发送端:延迟属性填充(gRPC)
    入口在 SendMessageActivity.fillDelayMessageProperty :
protectedvoidfillDelayMessageProperty(apache.rocketmq.v2.Messagemessage,org.apache.rocketmq.common.message.MessagemessageWithHeader){// 客户端在SystemProperties中指定投递时间, 1.判断是否为延迟消息if(message.getSystemProperties().hasDeliveryTimestamp()){TimestampdeliveryTimestamp=message.getSystemProperties().getDeliveryTimestamp();// 2/提取并转换时间戳,秒+纳秒,转成毫秒级时间戳longdeliveryTimestampMs=Timestamps.toMillis(deliveryTimestamp);// 3.校验延迟上限,目标投递时间-当前时间,超过最大限制,抛异常validateDelayTime(deliveryTimestampMs);ProxyConfigconfig=ConfigurationManager.getProxyConfig();// 走经典延迟队列,默认关闭if(config.isUseDelayLevel()){// 级别换算intdelayLevel=config.computeDelayLevel(deliveryTimestampMs);// 写入属性 PROPERTY_DELAY_TIME_LEVEL ,值是对应的级别数字MessageAccessor.putProperty(messageWithHeader,MessageConst.PROPERTY_DELAY_TIME_LEVEL,String.valueOf(delayLevel));}// 精确投递时间戳, 供 RocketMQ 5 的定时消息(TimerMessageStore 时间轮)使用StringtimestampString=String.valueOf(deliveryTimestampMs);MessageAccessor.putProperty(messageWithHeader,MessageConst.PROPERTY_TIMER_DELIVER_MS,timestampString);}}
  • 客户端在 gRPC SystemProperties.deliveryTimestamp 指定投递时间。
  • Proxy 校验时间合法( validateDelayTime ,受 maxDelayTimeMills 限制)。
  • 写入两个属性:
    • PROPERTY_TIMER_DELIVER_MS :精确投递时间戳(RocketMQ 5 的定时消息)。
    • PROPERTY_DELAY_TIME_LEVEL :仅当 useDelayLevel=true 时,把时间换算成经典「延迟级别」写进去(兼容老版延迟队列)。

2. 延迟级别换算(配置)

在 ProxyConfig :

  • 字段: useDelayLevel (默认 false)、 messageDelayLevel (默认 “1s 5s … 2h” )、 delayLevelTable 。
  • parseDelayLevel :把字符串解析成 级别 → 毫秒 映射。
  • computeDelayLevel :根据剩余时间算出对应的最小延迟级别。
publicintcomputeDelayLevel(longtimeMillis){// 计算剩余延迟时间longintervalMillis=timeMillis-System.currentTimeMillis();// 在延迟级别里找到第一个级别对应时长大于剩余时间的级别// delayLevelTable 由 parseDelayLevel 从配置字符串(如 "1s 5s 10s ... 2h" )解析出来,level 从 1 开始。List<Map.Entry<Integer,Long>>sortedLevels=delayLevelTable.entrySet().stream().sorted(Comparator.comparingLong(Map.Entry::getValue)).collect(Collectors.toList());for(Map.Entry<Integer,Long>entry:sortedLevels){if(entry.getValue()>intervalMillis){returnentry.getKey();}}// 循环跑完都没命中 ,说明 intervalMillis >= 所有级别时长 ,也就是剩余延迟已经 超过了最大级别 ,此时只能返回最后一个(最大)级别。returnsortedLevels.get(sortedLevels.size()-1).getKey();}
// proxy启动时就执行publicvoidparseDelayLevel(){this.delayLevelTable=newConcurrentSkipListMap<>();Map<String,Long>timeUnitTable=newHashMap<>();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"StringlevelString=this.getMessageDelayLevel();try{String[]levelArray=levelString.split(" ");for(inti=0;i<levelArray.length;i++){Stringvalue=levelArray[i];Stringch=value.substring(value.length()-1);// 时间单位映射Longtu=timeUnitTable.get(ch);// 定义级别intlevel=i+1;// 获取时长longnum=Long.parseLong(value.substring(0,value.length()-1));longdelayTimeMillis=tu*num;this.delayLevelTable.put(level,delayTimeMillis);}}catch(Exceptione){log.error("parse delay level failed. messageDelayLevel:{}",messageDelayLevel,e);}}

3. 消费端:识别为 DELAY 类型

GrpcConverter 在把 MessageExt 转成 gRPC Message 时,判断带 PROPERTY_DELAY_TIME_LEVEL / PROPERTY_TIMER_DELIVER_MS / PROPERTY_TIMER_DELAY_SEC 任一属性的消息,标记为 MessageType.DELAY 。

4. 容易混淆的「不可见时间」

ChangeInvisibleTimeActivity 和 gRPC 的 ChangeInvisibleDurationActivity 处理的是 pop 消费模型的 invisible time(消息不可见时长) ,用于消费失败后延迟重投,这是消费侧的机制, 不属于传统的「延迟队列」 。

总结 :proxy 里没有延迟队列的调度引擎,只有「延迟消息发送时的属性填充与级别换算」这部分,集中在 SendMessageActivity 和 ProxyConfig ;真正的延迟投递由 Broker 完成

完整链路

完整链路梳理如下,分四段:发送 → Broker 调度 → 消费 → 失败重试。

一、生产者发送延迟消息(proxy 侧)

gRPC 入口 : SendMessageActivity.fillDelayMessageProperty

  • 客户端在 SystemProperties.deliveryTimestamp 里指定投递时间。
  • validateDelayTime 校验延迟上限。
  • 写入两个属性:
    • PROPERTY_TIMER_DELIVER_MS (精确投递时间戳,走定时消息时间轮);
    • PROPERTY_DELAY_TIME_LEVEL (仅 useDelayLevel=true 时,走经典延迟队列)。
      级别换算 : ProxyConfig 里的 parseDelayLevel 和 computeDelayLevel 。

Remoting 入口 : remoting/activity/SendMessageActivity 只是把老客户端自带 DELAY_TIME_LEVEL 的请求透传给 Broker。

最终发送 : MessagingProcessor.sendMessage → ProducerProcessor → MessageService.sendMessage (Cluster/Local 两套实现)→ 发到 Broker。

二、Broker 侧延迟存储与调度(真正的延迟引擎,不在 proxy)

  • 定时消息: transformTimerMessage :算出 deliverMs ,备份 PROPERTY_REAL_TOPIC / PROPERTY_REAL_QUEUE_ID ,把消息 topic 改成 TIMER_TOPIC ( rmq_sys_wheel_timer ),queueId 置 0。
  • 经典延迟消息: transformDelayLevelMessage :备份 real topic/queueId,把 topic 改成 RMQ_SYS_SCHEDULE_TOPIC , queueId = level - 1 (每个延迟级别一个队列)。
    两套调度引擎 :
  1. 定时消息(时间轮) : TimerMessageStore (store 模块)

    • TimerWheel + TimerLog + enqueue/dequeue 服务线程。
    • doEnqueue 把消息按投递时间挂到时间轮槽位。
    • 到期后 convert 恢复原 topic/queueId,通过 escapeBridge 重新写回正常队列。
  2. 经典延迟队列 : ScheduleMessageService (broker 模块)

    2.1 写入拦截 : HookUtils.handleScheduleMessage (消息落盘前被 hook)

/** * 写入拦截 * @param brokerController * @param msg * @return */publicstaticPutMessageResulthandleScheduleMessage(BrokerControllerbrokerController,finalMessageExtBrokerInnermsg){finalinttranType=MessageSysFlag.getTransactionValue(msg.getSysFlag());if(tranType==MessageSysFlag.TRANSACTION_NOT_TYPE||tranType==MessageSysFlag.TRANSACTION_COMMIT_TYPE){if(!isRolledTimerMessage(msg)){if(checkIfTimerMessage(msg)){if(!brokerController.getMessageStoreConfig().isTimerWheelEnable()){//wheel timer is not enabled, reject the messagereturnnewPutMessageResult(PutMessageStatus.WHEEL_TIMER_NOT_ENABLE,null);}PutMessageResulttransformRes=transformTimerMessage(brokerController,msg);if(null!=transformRes){returntransformRes;}}}// Delay Delivery 定时任务// getDelayTimeLevel() 内部读的就是 PROPERTY_DELAY_TIME_LEVEL 属性,所以只有经典延迟消息(而非时间轮定时消息)才会走到这里。if(msg.getDelayTimeLevel()>0){transformDelayLevelMessage(brokerController,msg);}}returnnull;}
/** * 把带延迟级别的消息「改头换面」写进经典延迟队列专用 topic( SCHEDULE_TOPIC ) ,同时备份真实 topic/queueId,供到期后还原投递。- * @param brokerController * @param msg */publicstaticvoidtransformDelayLevelMessage(BrokerControllerbrokerController,MessageExtBrokerInnermsg){// 级别上限保护if(msg.getDelayTimeLevel()>brokerController.getScheduleMessageService().getMaxDelayLevel()){msg.setDelayTimeLevel(brokerController.getScheduleMessageService().getMaxDelayLevel());}// Backup real topic, queueId,备份真实topic、队列idMessageAccessor.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()));msg.setTopic(TopicValidator.RMQ_SYS_SCHEDULE_TOPIC);// 每个延迟级别占用一个队列:level 1 → queueId 0,level 2 → queueId 1,…… level 18 → queueId 17。msg.setQueueId(ScheduleMessageService.delayLevel2QueueId(msg.getDelayTimeLevel()));}

2.2 每个 delay level 起一个 DeliverDelayedMessageTimerTask 遍历 SCHEDULE_TOPIC 的 ConsumeQueue。
经典延迟队列调度引擎的 启动入口

// broker启动时,会调用到这里// org.apache.rocketmq.broker.schedule.ScheduleMessageService// 经典延迟队列调度引擎的 启动入口publicvoidstart(){// CAS 防重入,用原子 CAS 保证 start() 只会真正初始化一次,重复调用直接跳过。if(started.compareAndSet(false,true)){// 加载消费进度:从磁盘恢复 offsetTable (每个延迟级别已经消费到哪个 offset),避免 broker 重启后从 0 开始重复扫描/投递。this.load();// 创建扫描线程池:用于执行每个级别的 DeliverDelayedMessageTimerTask 定时扫描任务// 线程池大小,maxDelayLevel (默认 18,即延迟级别数)。this.deliverExecutorService=ThreadUtils.newScheduledThreadPool(this.maxDelayLevel,newThreadFactoryImpl("ScheduleMessageTimerThread_"));//可选创建异步投递线程池。默认关闭。// 开启后,投递结果的写盘处理交给独立的 handleExecutorService 异步执行,扫描线程只负责「找到期消息」,处理结果与扫描解耦,提升吞吐if(this.enableAsyncDeliver){this.handleExecutorService=ThreadUtils.newScheduledThreadPool(this.maxDelayLevel,newThreadFactoryImpl("ScheduleMessageExecutorHandleThread_"));}// 为每个级别起一个 DeliverDelayedMessageTimerTask 扫描对应 ConsumeQueuefor(Map.Entry<Integer,Long>entry:this.delayLevelTable.entrySet()){Integerlevel=entry.getKey();// 级别:1~18LongtimeDelay=entry.getValue();// 时长:1s~2h// 取该级别的起始消费进度Longoffset=this.offsetTable.get(level);if(null==offset){offset=0L;// 无记录则从 0 开始}if(timeDelay!=null){// 异步投递开启,额外为每个级别调度一个HandlePutResultTask(处理投递结果)。if(this.enableAsyncDeliver){this.handleExecutorService.schedule(newHandlePutResultTask(level),FIRST_DELAY_TIME,TimeUnit.MILLISECONDS);}//为每个级别调度一个 DeliverDelayedMessageTimerTask ( 核心扫描任务 ),延迟 FIRST_DELAY_TIME (1000ms)后开始执行,之后按级别时长周期循环。this.deliverExecutorService.schedule(newDeliverDelayedMessageTimerTask(level,offset),FIRST_DELAY_TIME,TimeUnit.MILLISECONDS);}}// 定时持久化消费进度scheduledPersistService.scheduleAtFixedRate(()->{try{ScheduleMessageService.this.persist();}catch(Throwablee){log.error("scheduleAtFixedRate flush exception",e);}// 首次延迟10秒,之和每间隔一段时间执行一次持久化},10000,this.brokerController.getMessageStoreConfig().getFlushDelayOffsetInterval(),TimeUnit.MILLISECONDS);}}

这段代码做三件事:加载消费进度 → 为每个延迟级别启动定时扫描线程 → 启动进度持久化任务。

要点说明
线程模型扫描线程池deliverExecutorService大小 = 延迟级别数(默认 18),一个级别一个定时任务
异步投递enableAsyncDeliver=true时,投递结果处理交给独立handleExecutorService,扫描与写盘解耦
进度恢复load()恢复offsetTable,重启后从上次 offset 续扫,避免重复投递
进度持久化scheduledPersistService定时persist(),保证 offset 落盘
级别隔离每个级别用独立DeliverDelayedMessageTimerTask+ 独立队列(queueId = level - 1),互不干扰

2.3 ConsumeQueue 的 tagsCode 存的是 投递时间戳 (由 CommitLog 在 dispatch 时用 computeDeliverTimestamp 算出)。

2.4 到期后 messageTimeUp 清除延迟属性,恢复原 topic/queueId,重新投递。

三、HandlePutResultTask

HandlePutResultTask 是 异步投递模式下的「结果处理器」 :它定时轮询每个延迟级别的待处理队列 deliverPendingTable ,按投递结果的状态推进消费 offset、重发或丢弃。核心逻辑在 run 。
开启 enableAsyncDeliver 后,扫描线程 DeliverDelayedMessageTimerTask 找到到期消息时,按照级别把消息放到队列 deliverPendingTable 中,以「级别」为单位轮询 deliverPendingTable ,按 FIFO 顺序消费 PutResultProcess : SUCCESS 推进 offset、 RUNNING 停下等待、 EXCEPTION 递增退避重发、 SKIP 丢弃;配合重试上限与流控,在异步投递下保证「不丢、不重、可重试」的消费进度推进

如果投递过程中 Broker 宕机了,HandlePutResultTask 的重试机制能确保消息不丢失吗?

结论: 能保证「消息不丢失」,但不能保证「不重复投递」 ——这是典型的 at-least-once(至少一次)语义。而且关键在于,宕机场景下的不丢失保障 并不依赖 HandlePutResultTask 的 doResend ,而是依赖更底层的持久化机制。

原因:

  • 消息本体在 CommitLog 已持久化
  • offset 只在投递成功后推进
  • 重启后 load() 从磁盘恢复 offset

四、关闭 enableAsyncDeliver (默认)时走 同步投递处理流程

核心是「扫描线程内阻塞等待投递结果,成功后立即推进 offset」

privatebooleansyncDeliver(MessageExtBrokerInnermsgInner,StringmsgId,longoffset,longoffsetPy,intsizePy){PutResultProcessresultProcess=deliverMessage(msgInner,msgId,offset,offsetPy,sizePy,false);// 执行 future.get() 阻塞直到投递返回结果PutMessageResultresult=resultProcess.get();// 关键:阻塞等待结果booleansendStatus=result!=null&&result.getPutMessageStatus()==PutMessageStatus.PUT_OK;if(sendStatus){// 成功后立即推进 offsetScheduleMessageService.this.updateOffset(this.delayLevel,resultProcess.getNextOffset());}// 返回 false , executeOnTimeUp 里 scheduleNextTimerTask(nextOffset, DELAY_FOR_A_WHILE) 停止本次扫描,延迟一段时间后重新扫描。returnsendStatus;}

五、消费者消费(proxy 侧)

gRPC 入口 : ReceiveMessageActivity.receiveMessage

  • 计算 invisibleTime / pollingTime,调用 messagingProcessor.popMessage 。
  • pop 结果经过 PopMessageResultFilterImpl 过滤:
    • tag 不匹配 → NO_MATCH ,直接 ACK 掉;
    • reconsumeTimes >= maxAttempts → TO_DLQ ,转发死信;
    • 否则 → MATCH ,返回给客户端。
      核心消费 : ConsumerProcessor.popMessage 构造 PopMessageRequestHeader ,交给 MessageService.popMessage ,最终到 Broker 的 PopMessageProcessor 。Broker 在 pop 时会同时消费正常 topic 和重试 topic( %RETRY%group )。

返回给客户端的消息携带 receiptHandle( PROPERTY_POP_CK )。

四、失败重试

分 proxy 侧和 Broker 侧两处配合完成。

proxy 侧:

  • 正常 ACK : AckMessageActivity → ConsumerProcessor.ackMessage → Broker,消费成功后消息被删除。

  • 主动 NACK / 快速重试 : ChangeInvisibleDurationActivity (gRPC)/ ChangeInvisibleTimeActivity (Remoting)→ ConsumerProcessor.changeInvisibleTime ,把 invisible time 改小,让消息更快重新可见。

  • 自动续期(renew) : DefaultReceiptHandleManager

    • 定时线程 scheduleRenewTask 扫描即将过期的 receiptHandle,在 invisible time 快到前自动 renewMessage (内部走 changeInvisibleTime )。
    • 超过 renewMaxTimeMillis 或重试上限,则触发 STOP_RENEW (即 NACK)。
  • 死信 : PopMessageResultFilterImpl 判定重试次数超限后,走 forwardMessageToDeadLetterQueue 。
    Broker 侧:

  • PopReviveService.reviveRetry :消息 invisible time 过期且未被 ACK 时,把消息投到 retry topic, reconsumeTimes + 1 。

  • AckMessageProcessor 处理 ACK 请求。
    一句话总结 :proxy 负责「发送时填延迟属性」和「消费时 pop + 过滤 + ACK/NACK/renew 编排」;真正的延迟存储与调度在 Broker/store 的 ScheduleMessageService (经典延迟队列)和 TimerMessageStore (定时消息时间轮);失败重试由 proxy 的 ReceiptHandleManager (renew)+ changeInvisibleTime (NACK)与 Broker 的 PopReviveService (revive 到 retry topic)共同完成。

版权声明: 本文来自互联网用户投稿,该文观点仅代表作者本人,不代表本站立场。本站仅提供信息存储空间服务,不拥有所有权,不承担相关法律责任。如若内容造成侵权/违法违规/事实不符,请联系邮箱:809451989@qq.com进行投诉反馈,一经查实,立即删除!
网站建设 2026/8/25 15:24:30

机器人破人类百米纪录:四足机器人的运动控制与工业落地

2026年&#xff0c;机器人再次刷新公众对运动能力的认知&#xff1a;一台四足机器人以10米/秒级别的速度完成百米冲刺&#xff0c;将“机器人破人类百米纪录”从实验室话题推向工程现实。对政企采购决策者而言&#xff0c;真正值得关注的不是单次赛跑成绩&#xff0c;而是这类高…

作者头像 李华
网站建设 2026/8/25 15:15:25

共享内存呀

共享内存&#xff08;Shared Memory&#xff09;是 IPC 进程间通信 方式之一&#xff0c;原理&#xff1a;在内核开辟一块物理内存&#xff0c;多个进程将这块内存映射到自己的虚拟地址空间&#xff0c;直接读写内存&#xff0c;无需内核拷贝数据&#xff0c;速度最快&#xff…

作者头像 李华