ARTICLE DETAIL

资讯详情

深耕郑州网站建设与运营推广的一线实战洞察。

EMQX MQTT 桥接 `$queue/` 订阅消息丢失修复:Subscription-Identifier 在 MQ 投递链路中的完整解析

EMQX MQTT 桥接 `$queue/` 订阅消息丢失修复:Subscription-Identifier 在 MQ 投递链路中的完整解析 后端物联网消息队列通信【免费下载链接】emqxThe most scalable and reliable MQTT broker for AI, IoT, IIoT and connected vehicles项目地址https://gitcode.com/gh_mirrors/em/emqx点击查看免费下载本文基于仓库变更记录 changes/ee/fix-16999.en.md 展开深入剖析 EMQX 在启用 Message QueueMQ特性的远端 broker 上MQTT 桥接MQTT Bridge作为 MQTT Source 从$queue/主题订阅时无法收到消息的问题根因与修复方案。读者读完本文后将掌握 MQTT v5 Subscription-Identifier 属性在 EMQX 订阅、投递、桥接路由三个环节的完整流转机制理解为什么$queue/队列订阅强制要求 MQTT v5 连接并能通过源码与测试用例验证该修复的可靠性。背景为什么 MQTT 桥接需要 Subscription-IdentifierEMQX 的 MQTT 桥接emqx_bridge_mqtt在作为 MQTT Sourceingress运行时会通过一条 MQTT 客户端连接订阅远端 broker 的主题把收到的消息导入本地 broker 或触发规则引擎处理。当远端 broker 同时开启了 EMQX 的 Message QueueMQ特性时客户端需要订阅形如$queue/queue-name/topic-filter的主题才能消费队列中的消息。MQTT v5 协议中客户端可以在 SUBSCRIBE 包的属性中携带一个或多个Subscription-Identifier订阅标识符服务端在向该订阅投递消息时会把对应的订阅标识符原样写入 PUBLISH 包的属性中。这让订阅方无需在客户端维护主题过滤器到业务处理的映射即可根据 PUBLISH 包属性快速识别这条消息来自哪一条订阅。EMQX MQTT 桥接的 ingress 正是依赖这一机制完成消息路由的桥接启动时为每条 ingress 订阅分配唯一的 Subscription-Identifier订阅时在 SUBSCRIBE 属性中带上该标识符收到消息后优先根据 PUBLISH 包中的 Subscription-Identifier 查索引定位对应的通道配置channel config只有找不到时才退化为按主题匹配。如果远端 broker 在投递 MQ 队列消息时没有在 PUBLISH 包中携带 Subscription-Identifieringress 就无法通过标识符定位到正确的通道消息随之丢失——这正是本次修复要解决的核心问题。问题描述与根因定位根据 changes/ee/fix-16999.en.md 的记载Fixed an issue where MQTT source failed to receive messages from$queue/subscriptions when the remote broker has the Message Queue (mq) feature enabled.问题现象远端 broker 开启 MQ 特性时MQTT 桥接 Source 从$queue/订阅收不到消息。根因MQ 消息投递时PUBLISH 包中没有包含 MQTT v5 Subscription-Identifier 属性而 MQTT 桥接 ingress 正是依赖该属性来路由来自队列订阅的消息。从源码层面可以还原这条因果链桥接 ingress 在订阅远端主题时会携带Subscription-Identifier属性见 emqx_bridge_mqtt_ingress.erlsubscribe_properties(#{subscription_id : SubscriptionId}) - #{Subscription-Identifier SubscriptionId}; subscribe_properties(_Ingress) - #{}.消息到达 ingress 后先按 Subscription-Identifier 查找通道配置找不到才按主题回退匹配见 emqx_bridge_mqtt_ingress.erlfind_channel_configs(Topic, Props, SubscriptionIdToHandlerIndex, TopicToHandlerIndex) - case find_channel_configs_by_subscription_identifier(Props, SubscriptionIdToHandlerIndex) of [] - find_channel_configs_by_topic(Topic, TopicToHandlerIndex); ChannelConfigs - ChannelConfigs end.修复前远端 broker 的 MQ 投递链路在生成 PUBLISH 包时未注入Subscription-Identifier属性导致第 2 步永远命中按标识符查找的空结果消息被静默丢弃。修复方案MQ 投递链路补全 Subscription-Identifier修复的实质是在 MQ 的消息投递路径上补齐 Subscription-Identifier 的注入同时完善桥接 ingress 侧的订阅协商与强制校验逻辑。整个修复可以拆解为三个层面。1. Broker 侧订阅解析与投递注入MQTT v5 通用能力EMQX 本身作为 MQTT broker 时对 Subscription-Identifier 的支持分为订阅与投递两步这也是远端 broker 正确投递队列消息的前提。订阅阶段emqx_channel在处理 SUBSCRIBE 包时把属性中的Subscription-Identifier提取为subid存入订阅选项见 emqx_channel.erlenrich_subopts_subid(TopicFilters, #{sub_props : #{Subscription-Identifier : SubId}}) - [{Topic, SubOpts#{subid SubId}} || {Topic, SubOpts} - TopicFilters]; enrich_subopts_subid(TopicFilters, _State) - TopicFilters.同时EMQX 在 CONNACK 属性中声明Subscription-Identifier-Available 1见 emqx_channel.erl明确告知客户端本服务端支持该特性。投递阶段emqx_session的enrich_message/3在把消息投递给订阅者时若该订阅带有subid则在消息头 properties 中注入Subscription-Identifier见 emqx_session.erl%% Add Subscription-Identifier if its specified for the %% subscription: Headers case SubId of undefined - Headers0; _ - case Headers0 of #{properties : Properties} - Headers0#{properties : Properties#{Subscription-Identifier SubId}}; #{} - Headers0#{properties #{Subscription-Identifier SubId}} end end,该逻辑有对应的单元测试佐证见 emqx_session.erl例如enrich_subid1_test/0验证subid 42的订阅会把Subscription-Identifier : 42写入消息属性。2. MQ 消费链路队列消息经 broker 投递时保留订阅标识符MQ 消息从入队到出队再进入 broker 投递管线的路径大致如下emqx_mq模块挂载message.publishhook见 emqx_mq.erl拦截发往 MQ 主题的消息通过publish_to_queue/2调用emqx_mq_message_db:insert/2写入消息数据库见 emqx_mq.erlMQ 消费者侧把出队消息通过emqx_mq_sub:messages/3送回订阅端见 emqx_mq_sub.erl跨节点时走 RPC见 emqx_mq_sub_proto_v1.erl消息最终经扩展订阅处理器emqx_mq_extsub_handler进入 broker 的投递流程见 emqx_mq_extsub_handler.erl。MQ 投递链路此前在消息进入 broker 投递管线前对订阅标识符的保留不完整导致最终生成的 PUBLISH 包缺失该属性。修复后队列消息在被投递给$queue/订阅者时会走与普通订阅一致的enrich_message/3注入逻辑从而在 PUBLISH 包中携带订阅时的 Subscription-Identifier。桥接 ingress 收到后即可按标识符完成路由。3. 桥接 ingress 侧订阅协商与强制校验桥接侧在本次修复中同步完善了与 Subscription-Identifier 相关的订阅行为涉及 emqx_bridge_mqtt_ingress.erl 与 emqx_bridge_mqtt_connector.erl 两个模块。订阅标识符分配每条 ingress 订阅会被分配一个唯一标识符最大取值为?MAX_SUBSCRIPTION_ID定义为268435455即 MQTT v5 中 Subscription-Identifier 的最大合法值 2^28−1见 emqx_bridge_mqtt_ingress.erl。分配逻辑会跳过已被占用的编号编号耗尽时报no_available_subscription_id见 emqx_bridge_mqtt_ingress.erl。降级重试协商订阅时若远端 broker 返回SUBSCRIPTION_IDENTIFIERS_NOT_SUPPORTED原因码行为分两种见 emqx_bridge_mqtt_ingress.erl普通主题订阅自动降级去掉 Subscription-Identifier 属性后重新订阅保证普通场景的兼容性$queue/队列订阅不降级直接返回错误subscription_identifier_required_for_queue_subscription。因为队列消息的路由完全依赖该属性降级后必然导致消息丢失显式报错比静默丢消息更安全。maybe_retry_without_subscription_identifier( #{subscription_id : _SubscriptionId, remote : #{topic : $queue/, _/binary}}, {ok, _Props, ReasonCodes} Result, _Pid ) - case lists:member(?RC_SUBSCRIPTION_IDENTIFIERS_NOT_SUPPORTED, ReasonCodes) of true - {error, subscription_identifier_required_for_queue_subscription}; false - Result end;启动期强制校验connector 在创建 ingress 配置时调用ensure_queue_subscription_supported/2若发现配置了$queue/主题但协议不是 MQTT v5不支持 Subscription-Identifier直接抛出不可恢复错误见 emqx_bridge_mqtt_connector.erlensure_queue_subscription_supported(#{topic : Topic}, SubscriptionIdToHandlerIndex) - case is_queue_topic(Topic) andalso not supports_queue_subscription(SubscriptionIdToHandlerIndex) of false - ok; true - error({unrecoverable_error, subscription_identifier_required_for_queue_subscription}) end.其中 ETS 订阅标识符索引仅在proto_ver : v5时创建见 emqx_bridge_mqtt_connector.erl与上面的强制校验形成闭环。消息到达后的路由ingress 收到 PUBLISH 包后解析属性中的Subscription-Identifier支持单个整数或多个标识符的列表见 emqx_bridge_mqtt_ingress.erl到 ETS 索引中查找对应的通道配置find_channel_configs_by_subscription_identifier(Props, SubscriptionIdToHandlerIndex) - lists:flatmap( fun(SubscriptionId) - case ets:lookup(SubscriptionIdToHandlerIndex, SubscriptionId) of [{SubscriptionId, ChannelConfig}] - [ChannelConfig]; [] - [] end end, subscription_ids(Props) ).配置要求与限制综合上述代码逻辑使用 MQTT 桥接 Source 订阅远端 broker 的$queue/主题时需要满足以下条件远端与本地连接均须使用 MQTT v5connector 配置中proto_ver v5Subscription-Identifier 是 MQTT v5 特性emqx_bridge_mqtt_connector.erl仅在 v5 下创建标识符索引且$queue/订阅在不支持该特性时会启动失败并报subscription_identifier_required_for_queue_subscription远端 broker 须正确实现 MQTT v5 订阅标识符语义订阅端携带的 Subscription-Identifier 必须原样出现在投递的 PUBLISH 包属性中。EMQX 自身通过 emqx_session.erl 的注入逻辑满足该要求这也是本次修复的核心订阅标识符数量有上限每条 ingress 订阅占用一个标识符最大编号为 268435455实际可分配数量受?MAX_SUBSCRIPTION_ID约束。测试验证本次修复的各个环节均有对应测试用例覆盖可作为回归验证的入口。MQ 侧保留验证emqx_mq_SUITE.erl 中的t_subscription_identifier/1用例完整模拟了问题场景创建名为subid的非 lastvalue 队列主题过滤器为t/#用SubId 42、属性#{Subscription-Identifier 42}订阅$queue/subid/t/#向t/0发布消息断言收到的 PUBLISH 包中properties含Subscription-Identifier : 42否则ct:fail(t/0 message from MQ not received or missing Subscription-Identifier)。该用例直接验证了修复目标MQ 队列消息投递时 PUBLISH 包必须携带订阅标识符。桥接侧路由与协商验证emqx_bridge_mqtt_source_SUITE.erl 中针对该场景有以下用例t_mqtt_conn_bridge_ingress_subid_dispatch第 486 行起文档注释为 MQTT ingress dispatches queue deliveries by Subscription Identifier.覆盖远端$queue/orders/t/#场景下按订阅标识符分发队列投递t_mqtt_conn_bridge_rejects_queue_source_without_subscription_identifier第 523 行起文档注释为 MQTT ingress rejects queue subscriptions when Subscription Identifiers are unavailable.断言报错原因字符串为queue subscriptions require connector proto_ver v5验证了启动期强制校验t_mqtt_conn_bridge_ingress_retries_without_subid_on_a1第 543 行起验证远端不支持订阅标识符时普通订阅先发带 subid 的 SUBSCRIBE、收到不支持响应后自动降级重发不带 subid 的 SUBSCRIBE第 597-598 行分别断言subscribe_with_subid与subscribe_without_subid消息。结语本次修复changes/ee/fix-16999.en.md从根因上补齐了 MQ 投递链路中 PUBLISH 包缺失 Subscription-Identifier 的问题使 MQTT 桥接 Source 在远端 broker 开启 MQ 特性时能够正常从$queue/订阅接收消息。修复同时完善了三层防护MQ 投递侧正确注入订阅标识符、桥接 ingress 侧按标识符优先路由、以及队列订阅对不支持 MQTT v5 的远端显式报错而非静默丢消息。对于运维与集成人员核心要点是只要桥接远端涉及$queue/队列订阅连接协议必须配置为 MQTT v5且远端 broker 需完整支持 Subscription-Identifier 语义。赞分享后端物联网消息队列通信【免费下载链接】emqxThe most scalable and reliable MQTT broker for AI, IoT, IIoT and connected vehicles项目地址https://gitcode.com/gh_mirrors/em/emqx点击查看免费下载相关推荐EMQX 消息流Streams投递修复解析Subscription Identifier 等订阅选项的完整传递链路EMQX 消息流Streams投递修复解析Subscription Identifier 等订阅选项的完整传递链路 本指南围绕变更记录 changes/e后端物联网消息队列通信EMQX 修复解析Message Queue 启用时 $queue/ 订阅无法向 MQTT 桥接投递消息的问题EMQX 修复解析Message Queue 启用时 $queue/ 订阅无法向 MQTT 桥接投递消息的问题 本文基于仓库变更记录 changes/ee/f后端物联网消息队列通信EMQX MQTT Ingress 桥接队列订阅$queue与 MQTT 5 Subscription Identifiers 深度解析EMQX MQTT Ingress 桥接队列订阅 $queue 与 MQTT 5 Subscription Identifiers 深度解析 MQTT in后端物联网消息队列通信创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表