rocketMQ proxy 延迟队列

rocketMQ proxy 延迟队列 Proxy 本身不做延迟队列的存储和调度 那是 Broker 端 ScheduleMessageService / TimerMessageStore 的职责Proxy 只负责 延迟消息的「发送端属性填充 延迟级别换算」以及消费端识别 。相关实现集中在 gRPC 发送链路和配置里。发送端延迟属性填充gRPC入口在 SendMessageActivity.fillDelayMessageProperty protectedvoidfillDelayMessageProperty(apache.rocketmq.v2.Messagemessage,org.apache.rocketmq.common.message.MessagemessageWithHeader){// 客户端在SystemProperties中指定投递时间, 1.判断是否为延迟消息if(message.getSystemProperties().hasDeliveryTimestamp()){TimestampdeliveryTimestampmessage.getSystemProperties().getDeliveryTimestamp();// 2/提取并转换时间戳秒纳秒转成毫秒级时间戳longdeliveryTimestampMsTimestamps.toMillis(deliveryTimestamp);// 3.校验延迟上限目标投递时间-当前时间超过最大限制抛异常validateDelayTime(deliveryTimestampMs);ProxyConfigconfigConfigurationManager.getProxyConfig();// 走经典延迟队列默认关闭if(config.isUseDelayLevel()){// 级别换算intdelayLevelconfig.computeDelayLevel(deliveryTimestampMs);// 写入属性 PROPERTY_DELAY_TIME_LEVEL 值是对应的级别数字MessageAccessor.putProperty(messageWithHeader,MessageConst.PROPERTY_DELAY_TIME_LEVEL,String.valueOf(delayLevel));}// 精确投递时间戳 供 RocketMQ 5 的定时消息TimerMessageStore 时间轮使用StringtimestampStringString.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 仅当 useDelayLeveltrue 时把时间换算成经典「延迟级别」写进去兼容老版延迟队列。2. 延迟级别换算配置在 ProxyConfig 字段 useDelayLevel 默认 false、 messageDelayLevel 默认 “1s 5s … 2h” 、 delayLevelTable 。parseDelayLevel 把字符串解析成 级别 → 毫秒 映射。computeDelayLevel 根据剩余时间算出对应的最小延迟级别。publicintcomputeDelayLevel(longtimeMillis){// 计算剩余延迟时间longintervalMillistimeMillis-System.currentTimeMillis();// 在延迟级别里找到第一个级别对应时长大于剩余时间的级别// delayLevelTable 由 parseDelayLevel 从配置字符串如 1s 5s 10s ... 2h 解析出来level 从 1 开始。ListMap.EntryInteger,LongsortedLevelsdelayLevelTable.entrySet().stream().sorted(Comparator.comparingLong(Map.Entry::getValue)).collect(Collectors.toList());for(Map.EntryInteger,Longentry:sortedLevels){if(entry.getValue()intervalMillis){returnentry.getKey();}}// 循环跑完都没命中 说明 intervalMillis 所有级别时长 也就是剩余延迟已经 超过了最大级别 此时只能返回最后一个最大级别。returnsortedLevels.get(sortedLevels.size()-1).getKey();}// proxy启动时就执行publicvoidparseDelayLevel(){this.delayLevelTablenewConcurrentSkipListMap();MapString,LongtimeUnitTablenewHashMap();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 2hStringlevelStringthis.getMessageDelayLevel();try{String[]levelArraylevelString.split( );for(inti0;ilevelArray.length;i){StringvaluelevelArray[i];Stringchvalue.substring(value.length()-1);// 时间单位映射LongtutimeUnitTable.get(ch);// 定义级别intleveli1;// 获取时长longnumLong.parseLong(value.substring(0,value.length()-1));longdelayTimeMillistu*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 仅 useDelayLeveltrue 时走经典延迟队列。级别换算 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 每个延迟级别一个队列。两套调度引擎 定时消息时间轮 TimerMessageStore store 模块TimerWheel TimerLog enqueue/dequeue 服务线程。doEnqueue 把消息按投递时间挂到时间轮槽位。到期后 convert 恢复原 topic/queueId通过 escapeBridge 重新写回正常队列。经典延迟队列 ScheduleMessageService broker 模块2.1 写入拦截 HookUtils.handleScheduleMessage 消息落盘前被 hook/** * 写入拦截 * param brokerController * param msg * return */publicstaticPutMessageResulthandleScheduleMessage(BrokerControllerbrokerController,finalMessageExtBrokerInnermsg){finalinttranTypeMessageSysFlag.getTransactionValue(msg.getSysFlag());if(tranTypeMessageSysFlag.TRANSACTION_NOT_TYPE||tranTypeMessageSysFlag.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);}PutMessageResulttransformRestransformTimerMessage(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 0level 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.deliverExecutorServiceThreadUtils.newScheduledThreadPool(this.maxDelayLevel,newThreadFactoryImpl(ScheduleMessageTimerThread_));//可选创建异步投递线程池。默认关闭。// 开启后投递结果的写盘处理交给独立的 handleExecutorService 异步执行扫描线程只负责「找到期消息」处理结果与扫描解耦提升吞吐if(this.enableAsyncDeliver){this.handleExecutorServiceThreadUtils.newScheduledThreadPool(this.maxDelayLevel,newThreadFactoryImpl(ScheduleMessageExecutorHandleThread_));}// 为每个级别起一个 DeliverDelayedMessageTimerTask 扫描对应 ConsumeQueuefor(Map.EntryInteger,Longentry:this.delayLevelTable.entrySet()){Integerlevelentry.getKey();// 级别1~18LongtimeDelayentry.getValue();// 时长1s~2h// 取该级别的起始消费进度Longoffsetthis.offsetTable.get(level);if(nulloffset){offset0L;// 无记录则从 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一个级别一个定时任务异步投递enableAsyncDelivertrue时投递结果处理交给独立handleExecutorService扫描与写盘解耦进度恢复load()恢复offsetTable重启后从上次 offset 续扫避免重复投递进度持久化scheduledPersistService定时persist()保证 offset 落盘级别隔离每个级别用独立DeliverDelayedMessageTimerTask 独立队列queueId level - 1互不干扰2.3 ConsumeQueue 的 tagsCode 存的是 投递时间戳 由 CommitLog 在 dispatch 时用 computeDeliverTimestamp 算出。2.4 到期后 messageTimeUp 清除延迟属性恢复原 topic/queueId重新投递。三、HandlePutResultTaskHandlePutResultTask 是 异步投递模式下的「结果处理器」 它定时轮询每个延迟级别的待处理队列 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){PutResultProcessresultProcessdeliverMessage(msgInner,msgId,offset,offsetPy,sizePy,false);// 执行 future.get() 阻塞直到投递返回结果PutMessageResultresultresultProcess.get();// 关键阻塞等待结果booleansendStatusresult!nullresult.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共同完成。