后端物联网消息队列通信【免费下载链接】emqxThe most scalable and reliable MQTT broker for AI, IoT, IIoT and connected vehicles项目地址https://gitcode.com/gh_mirrors/em/emqx点击查看免费下载导读本文围绕 EMQX 变更记录 changes/ee/fix-16974.en.md 所描述的缺陷展开在 EMQX 6.1.1 中当客户端会话已订阅过包含保留消息retained message的主题过滤器之后发生会话接管takeover或会话恢复resume且未重新订阅该主题过滤器时客户端会再次收到此前已经收到的保留消息造成重复投递。该问题在 6.1.2 中通过恢复旧有行为得到修复。读完本文你将理解 EMQX 会话接管/恢复的完整链路、emqx_extsub订阅扩展机制在其中扮演的角色以及保留消息迭代为何会在恢复场景下被显式终止。一、问题陈述变更记录原文changes/ee/fix-16974.en.md全文如下In EMQX 6.1.1, when a session was subscribed to a topic filter containing retained messages and was later taken over or resumed without re-subscribing to the same topic filter, it would receive again the received messages. Now, the previous behavior is restored, meaning that, upon session resumption or takeover without explicit re-subscription, retained message iteration will cease.核心要点可拆解为三点缺陷版本EMQX 6.1.1 引入回归缺陷现象会话接管或恢复后未重新订阅却再次收到已投递过的保留消息修复结果恢复旧行为——会话恢复/接管且未显式重新订阅时保留消息迭代立即停止。该条目同时收录于版本变更记录 changes/6.1.2.en.md确认修复随 6.1.2 版本发布。二、背景机制会话接管与恢复在 EMQX 中的含义在理解缺陷前需要先厘清 EMQX 中两个容易混淆的概念会话接管takeover同一 ClientID 的新连接发起时旧连接上的会话被接管。接管过程分两阶段begin/end由 emqx_cm_takeover.erl 实现旧通道在handle_call({takeover, begin}, ...)与handle_call({takeover, end}, ...)见 emqx_channel.erl中完成会话交接。会话恢复resumeMQTT 5.0 会话延续Session Expiry Interval 大于 0场景下客户端断线重连并恢复既有会话。从内存会话实现 emqx_session_mem.erl 可以看到两者底层的订阅处理takeover/1将当前会话的所有订阅从 broker 上逐一取消订阅emqx_broker:unsubscribe/1resume/2再把全部订阅重新注册到 brokeremqx_broker:subscribe/1。也就是说恢复/接管过程中 broker 侧的订阅关系经历先拆除、再重建但会话的订阅表本身并未变化。修复前这个重建动作会连带触发保留消息的重新迭代这正是缺陷的根源。三、缺陷根因emqx_extsub机制与session.resumed钩子3.1 6.1.1 引入的架构变化EMQX 6.1.x 引入了emqx_extsubexternal subscription扩展订阅机制保留消息等功能的订阅不再直接挂在 broker 订阅路径上而是通过扩展订阅处理器handler注册到统一的注册表中。保留消息模块在加载钩子时完成注册emqx_retainer.erlload_hooks() - ok emqx_extsub_handler_registry:register(emqx_retainer_extsub_handler, #{ handle_generic_messages false, multi_topic false, ignore_resubscribe false }), ...同时emqx_extsub模块在启动时挂接了session.resumed钩子emqx_extsub.erl并在会话恢复时以resume类型触发订阅回调emqx_extsub.erlon_session_resumed(ClientInfo, #{subscriptions : Subs} SessionInfo) - ... on_subscribed(resume, ClientInfo, Subs).3.2 缺陷链条在引入emqx_extsub之前保留消息的投递并不挂钩session.resumed因此会话恢复/接管时保留消息迭代会自然停止。而 6.1.1 将保留消息迁移到扩展订阅机制后session.resumed钩子会在恢复时以resume订阅类型重新触发保留消息处理器若处理器按普通订阅一样重新开启迭代就会把已投递过的保留消息再次发送给客户端——这就是 6.1.1 回归缺陷的直接原因。3.3 修复resume分支显式忽略修复点在保留消息的扩展订阅处理器 emqx_retainer_extsub_handler.erl。handle_subscribe/4对三种情况分派handle_subscribe(_SubscribeType, _SubscribeCtx, _Handler, #share{}) - ignore; handle_subscribe( _SubscribeType, #{clientinfo : #{protocol : Protocol}}, _Handler, _TopicFilter ) when Protocol / mqtt - ignore; handle_subscribe(resume _SubscribeType, _SubscribeCtx, _Handler, _TopicFilter) - %% Note: the implementation prior to using emqx_extsub did **not** resume iteration %% of retained messages when resuming a session. The previous implementation without %% extsub did not hook into session.resumed, hence iteration stopped when %% resuming/taking over. For now, we adopt the same behavior. Resuming iteration %% would thus be an improvement. ignore; handle_subscribe(subscribe _SubscribeType, SubscribeCtx, Handler, TopicFilter) - ...其中第三个子句resume类型直接返回ignore不再启动或延续保留消息的迭代。源码注释明确说明了设计意图恢复旧实现的行为——恢复/接管时停止保留消息迭代同时注释也客观指出未来若想恢复迭代这将是一个改进点improvement即当前实现并非最终理想态而是兼容性优先的选择。对比subscribe分支emqx_retainer_extsub_handler.erl可以看到只有真正的新订阅subscribe类型才会根据 MQTT 订阅选项 RHRetain Handling决定是否投递保留消息handle_subscribe(subscribe _SubscribeType, SubscribeCtx, Handler, TopicFilter) - IsNew ..., #{subopts : #{rh : RH} SubOpts} SubscribeCtx, case RH 0 orelse (RH 1 andalso IsNew) of true - subscribe(Handler, SubscribeCtx, TopicFilter, SubOpts); false - ignore end.3.4 配套逻辑接管/断开时强制终止迭代除了入口分支处理器还在handle_save_subopts/3emqx_retainer_extsub_handler.erl中保证接管/断开持久化发生时迭代被强制终止handle_save_subopts(#h{cursor ?cursor(_)} Handler0, _Context, _SubOpts) - Res #{delivered Handler0#h.delivered}, %% Make the handler stop, since take over/disconnect persistence is ongoing. Handler Handler0#h{cursor ?done}, ensure_cursor_deleted(Handler), {ok, Handler, Res};将游标置为?done并删除游标记录使迭代停止同时保存已投递计数delivered。这一设计与resume分支的ignore相互配合接管时停止迭代、恢复时不重新开启共同保证未显式重新订阅则不重复投递的语义。四、修复后的行为语义修复后保留消息的投递遵循以下规则结合 MQTT 规范与源码实现场景保留消息行为新订阅subscribe非共享、MQTT 协议按订阅选项 RH 投递RH0或RH1且订阅为新订阅时投递RH2不投递会话恢复/接管resume未显式重新订阅不投递保留消息迭代停止共享订阅#share{}忽略不处理保留消息非 MQTT 协议忽略这里需要特别说明 RH 选项语义MQTT 5.0 的 Retain Handling 订阅选项中RH1表示仅当订阅为新订阅时发送保留消息RH2表示不发送保留消息。修复前的 6.1.1 缺陷本质上破坏了这一语义——即使客户端没有重新订阅订阅并非新的session.resumed钩子仍然触发了保留消息投递等价于把RH1错误地当成RH0处理。从注册表侧看resume订阅类型还会在recreate/3中被用于按会话保存的订阅重建处理器emqx_extsub_handler_registry.erl因此上述resume分支的ignore同样覆盖会话恢复时重建扩展订阅处理器的路径确保整条恢复链路行为一致。五、影响范围与验证方式5.1 影响范围影响版本EMQX 6.1.1含该回归的版本修复版本EMQX 6.1.2见 changes/6.1.2.en.md影响场景使用持久会话clean_start0 或 Session Expiry Interval 0、订阅了含有保留消息的主题过滤器、并经历断线重连或同 ClientID 多连接接管的客户端。5.2 验证思路可依据仓库内测试体系与源码进行验证例如在 6.1.1 与 6.1.2 上分别构造订阅含保留消息的主题 → 接收保留消息 → 断开 → 以相同会话恢复/接管连接的用例断言恢复后收到的消息数量修复后不应再次收到已投递的保留消息直接单测emqx_retainer_extsub_handler:handle_subscribe(resume, ...)断言其返回ignore且不产生任何新的迭代任务结合handle_save_subopts/3验证接管过程中游标被置?done、已投递计数被持久化保存。六、总结fix-16974 是一次典型的架构迁移引入回归、以兼容旧语义方式修复的缺陷修复根因EMQX 6.1.1 将保留消息投递迁移到emqx_extsub扩展订阅机制后session.resumed钩子会在会话恢复/接管时以resume订阅类型重新触发保留消息处理器导致已投递的保留消息被再次发送修复在 emqx_retainer_extsub_handler.erl 的resume分支显式返回ignore恢复恢复/接管时保留消息迭代停止的旧行为并与handle_save_subopts/3的强制终止逻辑形成完整闭环取舍源码注释明确指出将来在会话恢复时恢复保留消息迭代仍可作为改进方向本次修复优先保证行为兼容与消息不重复。对于实际部署 EMQX 6.1.1 且依赖持久会话 保留消息组合的用户建议升级到 6.1.2 或更高版本以获得该修复。赞分享后端物联网消息队列通信【免费下载链接】emqxThe most scalable and reliable MQTT broker for AI, IoT, IIoT and connected vehicles项目地址https://gitcode.com/gh_mirrors/em/emqx点击查看免费下载相关推荐EMQX 会话接管后保留消息迭代恢复机制解析如何减少通配符订阅中的重复投递EMQX 会话接管后保留消息迭代恢复机制解析如何减少通配符订阅中的重复投递 导读 本文围绕 EMQX 开源仓库中的变更记录 feat 16637 https:后端物联网消息队列通信EMQX 保留消息投递限流重试机制修复解析从丢消息到指数退避恢复EMQX 保留消息投递限流重试机制修复解析从丢消息到指数退避恢复 导读 本篇文章基于 EMQX 开源仓库中的变更记录 changes/ee/fix 16553后端物联网消息队列通信EMQX 会话接管场景下 connected_at 与 disconnected_at 时间戳乱序问题的修复解析EMQX 会话接管场景下 connected_at 与 disconnected_at 时间戳乱序问题的修复解析 导读 在 EMQX 中当一个客户端以相同 C后端物联网消息队列通信创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考