发布时间:2026/8/25 15:47:06
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共同完成。

相关新闻

2026/8/25 15:47:06

第21届全国大学生智能车竞赛软件盲盒任务

第21届全国大学生智能车竞赛总决赛赛道设计与部署 【软件盲盒任务】 在 第21届全国大学生智能汽车竞赛全国总决赛 中,不再设置硬件盲盒任务, 针对每个赛题组设置软件盲盒任务。 盲盒任务在竞赛领队会上进行公布。 如果没有完成盲盒任务, 每个…

2026/8/25 15:47:06

【SRC】EDU实战篇1:弱口令之新站挖掘与横向资产排查

文章目录 思路 案例一(1个) 案例二(1个) 案例三(1个) 案例四(4个) 站点1 站点2 站点3 站点4 案例五(1个) 案例六(1个) 总结 ⚠️本博文所涉安全渗透测试技术、方法及案例,仅用于网络安全技术研究与合规性交流,旨在提升读者的安全防护意识与技术能力。任何个人或组…

2026/8/25 18:13:01

FPGA实现后调试实战:ILA、VIO与增量编译技术解析

1. 项目概述:为什么实现后的调试是FPGA开发的“深水区”刚接触FPGA开发的朋友,往往把大部分精力放在写代码、跑仿真上,觉得综合实现通过、比特流生成成功,项目就大功告成了。但真正在一线摸爬滚打过的人都知道,把设计下…

2026/8/25 18:13:01

C++字符与字符串处理:从编码原理到高效实践

1. 项目概述:从字符到字符串的C表达艺术在C的世界里,字符和字符串的处理是构建任何程序,无论是底层系统、游戏逻辑还是数据处理工具,都绕不开的基石。很多初学者,甚至一些有经验的开发者,常常会陷入一些看似…

2026/8/25 18:13:01

2013-2020年中国地面一氧化碳数据集(日/月/年)

简介 一氧化碳(CO)是一种无色、无味、无臭的气体,由一分子碳和一分子氧组成。它是一种有毒气体,在高浓度下容易危及人体健康,可以影响心脏、中枢神经系统和呼吸系统的正常功能。一氧化碳是一种主要的空气污染物之一&a…

2026/8/25 18:13:01

SpringBoot启动后自动打开浏览器:原理、实现与生产级配置

1. 项目概述:为什么需要启动后自动打开浏览器?做SpringBoot开发的朋友,估计都经历过这个场景:本地调试时,项目启动成功了,控制台打印出“Tomcat started on port(s): 8080 (http)”或者看到熟悉的Spring Lo…

2026/8/25 18:13:01

ABAP读取SMW0 Excel模板并写入数据的自动化方案

1. 项目概述与核心价值在SAP ERP的日常运维和开发中,我们经常遇到一个场景:业务用户需要定期生成格式固定、但数据动态变化的报表,比如月度销售分析、库存盘点表或者财务对账单。这些报表往往要求使用公司统一的Excel模板,包含特定…

2026/8/25 1:04:19

[光学原理与应用-521]:对光的错误理解与纠偏

首先光是一种能量的载体和形态,宏观上观察到的光是由无数个微观的光量子组成的,每个光子在产生的瞬间,其在真空的空间中以确定不变的速度沿着一个初始的方向一直向前,在微观层面,每个光量子的运动轨迹是以波函数所展现…

2026/8/25 11:48:27

SIP通话转接原理与REFER方法实战解析

1. 通话转接不是“挂断再拨号”,而是SIP会话的动态重定向你有没有遇到过这样的场景:客服坐席A正在和客户通电话,突然需要把这通对话无缝转给专家坐席B,客户完全感知不到中间的断连——既没听到忙音,也没被要求重新拨号…

2026/8/25 16:56:43

Kolla-ansible单节点OpenStack部署实战:从环境准备到排坑指南

1. 为什么选择Kolla-ansible来部署单节点OpenStack?如果你正在寻找一种能把OpenStack从“概念”快速变成“可用的实验环境”的方法,那么Kolla-ansible几乎是当前最主流、最省心的选择。我见过太多人卡在手动编译依赖、配置服务、处理版本冲突的泥潭里&am…

2026/8/25 0:04:14

三步把QQ空间历史说说导出到本地:GetQzonehistory 极简指南

三步把QQ空间历史说说导出到本地:GetQzonehistory 极简指南 【免费下载链接】GetQzonehistory 获取QQ空间发布的历史说说 项目地址: https://gitcode.com/GitHub_Trending/ge/GetQzonehistory Meta Description:GetQzonehistory 是一个QQ空间历史说…

2026/8/25 0:04:14

洛谷 P7912:[CSP-J 2021 T4] 小熊的果篮 ← 双向链表

【题目来源】 https://www.luogu.com.cn/problem/P7912 【题目描述】 小熊的水果店里摆放着一排 n 个水果。每个水果只可能是苹果或桔子,从左到右依次用正整数 1,2,…,n 编号。连续排在一起的同一种水果称为一个“块”。小熊要把这一排水果挑到若干个果篮里&#x…

2026/8/24 13:42:17

实测才敢推 AI论文网站 2026最新测评与推荐

2026年真正好用的AI论文网站,核心看生成的论文质量、低AI味、格式正确、学术适配四大指标。综合实测,千笔AI、ThouPen、豆包、DeepSeek、Grammarly 是当前最值得推荐的梯队,覆盖从免费到付费、从中文到英文、从文科到理工的全场景需求。一、综…

2026/8/24 18:13:48

2026必备!AI论文网站测评:最新推荐与深度对比

2026年真正好用的AI论文网站,核心看生成的论文质量、低AI味、格式正确、学术适配四大指标。综合实测,千笔AI、ThouPen、豆包、DeepSeek、Grammarly 是当前最值得推荐的梯队,覆盖从免费到付费、从中文到英文、从文科到理工的全场景需求。 一、…

2026/8/25 1:08:14

摆脱论文困扰!盘点2026年全网爆红的的AI论文写作工具

一天写完毕业论文在2026年已不再是天方夜谭。2026年最炸裂、实测能大幅提速的AI论文写作工具,覆盖选题构思、文献整理、内容生成、格式排版等核心场景,真正帮你高效搞定论文难题。 一、全流程王者:一站式搞定论文全链路(一天定稿首…