EMQX `$SYS` 保留消息过期机制:修复 StatefulSet 轮换后的陈旧节点标识问题
发布时间:2026/9/23 23:00:42 锦皓数字建站

后端物联网消息队列通信【免费下载链接】emqxThe most scalable and reliable MQTT broker for AI, IoT, IIoT and connected vehicles项目地址https://gitcode.com/gh_mirrors/em/emqx点击查看免费下载导读本文基于 EMQX 开源仓库的变更记录fix-16715深入解析一个针对$SYS系统主题保留消息的修复此前$SYS保留消息如 broker 节点标识、系统描述等主题在持久化时不带过期时间在 Kubernetes StatefulSet 滚动轮换rotation等场景下旧节点的身份主题会以“僵尸”状态残留在保留消息存储中并持续显示在 Dashboard 的节点视图中。修复后新发布的保留$SYS消息会自动携带Message-Expiry-Interval 36001 小时同时本文给出清理存量陈旧条目的完整命令并结合 emqx_sys.erl 与 emqx_retainer.erl 的源码说明该机制的实现原理与验证方式。读完本文你将掌握$SYS保留消息的过期语义、排查与清理陈旧节点主题的实战方法。一、问题背景$SYS保留消息为何会“滞留”1.1$SYS主题与保留消息EMQX 通过emqx_sysapps/emqx/src/emqx_sys.erl这个 gen_server 周期性向$SYS/...主题树发布系统信息包括节点 uptime、版本、sysdescr系统描述、broker 节点列表、stats 统计与 metrics 指标等。其中一部分消息带有retain保留标志例如publish(version, Version) - safe_publish(systop(version), #{retain true}, Version); publish(sysdescr, Descr) - safe_publish(systop(sysdescr), #{retain true}, Descr); publish(brokers, Nodes) - ... safe_publish($SYS/brokers, #{retain true}, Payload);即$SYS/brokers/...节点身份、broker 标识这类主题是保留消息当某个客户端例如 Dashboard 或监控采集器订阅该主题时会立即收到最近一次保留的消息内容。1.2 StatefulSet 轮换暴露出的“陈旧节点标识”在 Kubernetes 中以 StatefulSet 部署 EMQX 集群时Pod 会按序滚动更新rotation。节点 Pod 被销毁重建后会以新的emqxhostname身份加入集群而旧节点 Pod 已不复存在。但问题在于旧节点此前发布的保留$SYS消息如$SYS/brokers/emqxold-node/sysdescr被持久化在 retainer 存储中且没有设置任何过期时间。于是Dashboard 节点视图中会持续列出已经下线、不存在的节点标识订阅$SYS/brokers/#的客户端会反复收到这些陈旧数据由于这些条目永久不过期除非人工干预否则会一直残留在存储中。这正是本变更所修复的核心缺陷保留$SYS消息缺乏过期控制导致陈旧节点标识在 StatefulSet 轮换后长期可见。二、修复方案为保留$SYS消息统一注入 1 小时过期变更的核心结论非常简洁现在新发布的保留$SYS消息都会包含Message-Expiry-Interval 3600即 1 小时。Message-Expiry-Interval是 MQTT 5.0 规范中定义的报文属性表示消息在成为保留消息后允许保留的秒数retainer 模块会依据该属性计算消息的过期时间并在到期后将其清除。由此即使节点轮换后旧节点身份主题仍被发布为保留消息也会在 1 小时后自动过期清理不再永久滞留。三、源码实现剖析maybe_set_retained_expiry/13.1 常量定义在 emqx_sys.erl 中定义了本次引入的过期时长常量-define(RETAINED_SYS_MSG_EXPIRY_INTERVAL, 3600).3.2 注入逻辑所有$SYS系统消息最终都经由safe_publish/2,3统一发布emqx_sys.erlsafe_publish(Topic, Payload) - safe_publish(Topic, #{}, Payload). safe_publish(Topic, Flags, Payload) - Msg emqx_message:set_flags( maps:merge(#{sys true}, Flags), emqx_message:make(?SYS, Topic, iolist_to_binary(Payload)) ), emqx_broker:safe_publish(maybe_set_retained_expiry(Msg)).关键函数maybe_set_retained_expiry/1emqx_sys.erl在真正发布前对消息做“补丁”maybe_set_retained_expiry(Msg) - case emqx_message:get_flag(retain, Msg, false) of true - Props emqx_message:get_header(properties, Msg, #{}), case maps:is_key(Message-Expiry-Interval, Props) of true - Msg; false - emqx_message:set_header( properties, Props#{Message-Expiry-Interval ?RETAINED_SYS_MSG_EXPIRY_INTERVAL}, Msg ) end; false - Msg end.其行为可以概括为三条规则消息类型处理方式带 retain 标志、且未显式携带Message-Expiry-Interval自动注入Message-Expiry-Interval 3600带 retain 标志、但已显式携带该属性保留原值不覆盖例如显式设为0表示永不过期时尊重用户意图不带 retain 标志的普通$SYS消息不做任何修改从实现可见该补丁只影响保留消息普通$SYS实时消息的投递路径完全不受影响。3.3 单元测试印证源码内置了 EUnit 测试emqx_sys.erl直接验证上述三条规则maybe_set_retained_expiry_test() - Msg0 emqx_message:make(?SYS, $SYS/test, payload), Msg1 emqx_message:set_flag(retain, true, Msg0), Msg2 maybe_set_retained_expiry(Msg1), Props emqx_message:get_header(properties, Msg2, #{}), ?assertEqual(?RETAINED_SYS_MSG_EXPIRY_INTERVAL, maps:get(Message-Expiry-Interval, Props)). maybe_set_retained_expiry_non_retained_test() - ... %% 非保留消息properties 保持 undefined maybe_set_retained_expiry_preserve_existing_test() - ... %% 已有 Message-Expiry-Interval10 时结果仍为 10不被覆盖三个用例分别对应上表的三种分支可作为验证该行为的最小回归测试参考。四、过期时间如何被 Retainer 消费4.1get_expiry_time/1从消息属性换算绝对过期时间消息发布后retainer 负责持久化保留消息。在 emqx_retainer.erl 中get_expiry_time/1依据消息属性计算绝对过期时间get_expiry_time(#message{headers #{properties : #{Message-Expiry-Interval : 0}}}) - 0; get_expiry_time(#message{ headers #{properties : #{Message-Expiry-Interval : Interval}}, timestamp Ts }) - Ts Interval * 1000; get_expiry_time(#message{timestamp Ts}) - Interval emqx_conf:get([retainer, msg_expiry_interval]), case Interval of 0 - 0; _ - Ts Interval end.对照本次修复修复前$SYS保留消息没有Message-Expiry-Interval属性走第三个分支依赖retainer.msg_expiry_interval全局配置若该配置为0默认行为则expiry_time为0表示永不过期——这正是陈旧条目滞留的根源修复后$SYS保留消息自动携带Message-Expiry-Interval 3600get_expiry_time返回消息时间戳 3600 * 1000毫秒到期后由 retainer 的清理逻辑自动删除。expiry_time 0表示“永不过期”这一约定在 emqx_retainer_mnesia.erl如ExpiryTime : 0 orelse ExpiryTime Now的过期判断、#retained_message{}与#retained_index{}记录中expiry_time字段的持久化中贯穿使用读者可在该文件中进一步追踪过期条目的存储与回收细节。4.2 与emqx_message:update_expiry的关系在 emqx_message.erl 中update_expiry/1负责根据 zone 配置[mqtt, message_expiry_interval]对非保留消息的过期属性做换算/更新。它与$SYS补丁的分工是emqx_message:update_expiry/1处理普通消息的会话过期语义基于 zone 配置emqx_sys:maybe_set_retained_expiry/1专门针对$SYS保留消息注入固定 1 小时过期。两者相互独立、互不干扰共同完善了消息过期属性的覆盖范围。五、清理存量陈旧条目发布空保留消息本次修复只对新发布的保留$SYS消息生效。对于在修复生效前就已写入 retainer 存储、永不过期的陈旧$SYS条目需要手动清除。官方给出的方法是向该陈旧主题发布一条 payload 为空的保留消息以空消息覆盖旧内容。在运行中的 EMQX 节点上执行单行命令emqx eval emqx:publish(emqx_message:set_flag(retain, true, emqx_message:make(emqx_sys, $SYS/brokers/emqx127.0.0.1/sysdescr, ))).命令拆解片段作用emqx eval ...在 EMQX 节点内联执行 Erlang 表达式EMQX 自带的管理命令emqx_message:make(emqx_sys, Topic, )构造一条源为emqx_sys、payload 为空二进制的消息emqx_message:set_flag(retain, true, Msg)为消息打上 retain 标志使其按保留消息处理emqx:publish(Msg)发布消息入口定义于 emqx.erl消息经 broker 投递并进入 retainer 存储发布后该主题的保留内容被空消息覆盖即等效于“删除”该保留条目。请将命令中的主题替换为你实际要清理的陈旧$SYS/...主题例如emqx eval emqx:publish(emqx_message:set_flag(retain, true, emqx_message:make(emqx_sys, $SYS/brokers/emqxstale-node-0/sysdescr, ))).5.1 如何发现需要清理的陈旧主题排查时可以订阅或查询$SYS/brokers/#通配主题比对当前集群实际运行节点如通过 Dashboard 节点列表或emqx_ctl cluster status与保留消息中出现的节点标识凡是保留消息中出现、但已不在运行节点列表中的emqxhostname主题即为候选清理对象。5.2 运维建议升级后一次性巡检完成包含本次修复的 EMQX 版本升级后建议立即对存量$SYS/brokers/#保留主题做一次巡检与清理避免陈旧标识继续误导 Dashboard 与监控系统依赖自动过期而非手工修复上线后新产生的陈旧节点标识会在 1 小时内被 retainer 自动回收手工清理仅需针对存量数据执行一次勿随意调大该值RETAINED_SYS_MSG_EXPIRY_INTERVAL在源码中固定为3600不建议为追求“长期可见”而修改否则陈旧节点标识滞留窗口会重新变长。六、总结维度修复前修复后新发布保留$SYS消息无Message-Expiry-Interval可能永不过期自动注入Message-Expiry-Interval 36001 小时过期StatefulSet 轮换后的旧节点标识长期残留Dashboard 可见1 小时内自动过期清除显式携带过期属性的消息—尊重原值不覆盖存量陈旧条目需手动清理仍可通过发布空保留消息手动清理本次修复通过emqx_sys发布链路中的maybe_set_retained_expiry/1一处补丁配合 retainer 既有的get_expiry_time/1过期机制系统性解决了$SYS保留消息无限期滞留的问题对以 StatefulSet 方式滚动运维 EMQX 集群的场景尤为重要。相关实现与测试均可直接在仓库中查阅emqx_sys.erl补丁逻辑与单元测试、emqx_retainer.erl过期时间换算、emqx_retainer_mnesia.erl过期条目存储与回收。赞分享后端物联网消息队列通信【免费下载链接】emqxThe most scalable and reliable MQTT broker for AI, IoT, IIoT and connected vehicles项目地址https://gitcode.com/gh_mirrors/em/emqx点击查看免费下载相关推荐EMQX 保留消息投递限流重试机制修复解析从丢消息到指数退避恢复EMQX 保留消息投递限流重试机制修复解析从丢消息到指数退避恢复 导读 本篇文章基于 EMQX 开源仓库中的变更记录 changes/ee/fix 16553后端物联网消息队列通信pnpm Windows 升级后残留陈旧 pnpm.ps1 的修复bin 链接器删除过期 PowerShell shim 的完整解析pnpm Windows 升级后残留陈旧 pnpm.ps1 的修复bin 链接器删除过期 PowerShell shim 的完整解析 PowerShell 解包管理器开发工具CLIEMQX 会话恢复与接管场景下保留消息重复投递问题fix-16974修复解析EMQX 会话恢复与接管场景下保留消息重复投递问题fix 16974修复解析 导读 本文围绕 EMQX 变更记录 changes/ee/fix 16974.后端物联网消息队列通信创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
锦
锦皓数字建站
深耕本土企业品牌数字化升级,专注原创端正雅致商务官网,从视觉设计到稳定运维全程保驾护航。