后端物联网消息队列通信【免费下载链接】emqxThe most scalable and reliable MQTT broker for AI, IoT, IIoT and connected vehicles项目地址https://gitcode.com/gh_mirrors/em/emqx点击查看免费下载导读MQTT Ingress入站桥接让 EMQX 可以从远端 MQTT Broker 消费消息再通过规则引擎转发到本地主题。本文围绕 EMQX 仓库中changes/ee/fix-16979.en.md所记录的能力展开当远端 Broker 支持 MQTT 5 Subscription Identifiers订阅标识符时EMQX 的 MQTT 入站桥接可以直接订阅以$queue/{name}/{bind-filter}形式暴露的远端消息队列当 Subscription Identifiers 不可用时队列订阅会被拒绝而普通主题订阅则会自动降级重试去掉订阅标识符后再订阅。文章结合apps/emqx_bridge_mqtt的源码与测试用例逐层拆解这一机制的配置方式、判定逻辑与消息分发路径帮助你理解并正确使用基于队列的 MQTT 数据接入方案。背景为什么队列订阅依赖 Subscription Identifiers在 MQTT 协议中$queue/前缀是 EMQX 等 Broker 内部约定的一种特殊主题形式用于把一批主题的消息聚合到一个队列里消费常写作$queue/{name}/{bind-filter}{name}是队列名称{bind-filter}是绑定到该队列的主题过滤器如orders/t/#队列会汇集所有匹配该过滤器发布的消息。这类队列语义与标准 MQTT 的共享订阅$share/{group}/{filter}不同队列订阅者收到的是从队列投递出来的消息而不是针对原始主题的实时转发。由于$queue订阅者的原始订阅过滤器与最终消费的主题之间并不一一对应Broker 在向订阅者投递消息时需要一种机制让订阅者知道“这条消息对应的是哪一条订阅”。MQTT 5 的 Subscription Identifier订阅标识符恰好提供了这个能力订阅者在 SUBSCRIBE 报文中携带Subscription-Identifier属性Broker 在 PUBLISH 报文的属性中回填该标识符从而把“消息”与“订阅”关联起来。EMQX 正是借助这一点来实现$queue订阅的消息路由入站桥接为每条队列订阅分配一个唯一的 Subscription Identifier收到远端 PUBLISH 后通过报文属性中的Subscription-Identifier精确匹配到对应的桥接通道再交给规则引擎处理。队列订阅的开启条件与拒绝逻辑前置条件连接器协议版本必须是 v5队列订阅依赖 MQTT 5 的 Subscription Identifier 特性因此连接器Connector的proto_ver必须配置为v5。从源码看连接器启动时会根据proto_ver决定是否创建订阅标识符索引表在 emqx_bridge_mqtt_connector.erl 中maybe_new_subscription_id_index/1只在proto_ver : v5时创建 ETS 表subscription_id_to_handler_index其他协议版本返回undefinedsupports_queue_subscription/1直接以该索引是否存在作为判定依据见 emqx_bridge_mqtt_connector.erl索引为undefined即视为不支持队列订阅添加通道时ensure_queue_subscription_supported/2会检查如果主题以$queue/开头且不支持队列订阅则直接报错subscription_identifier_required_for_queue_subscription见 emqx_bridge_mqtt_connector.erl。连接器的proto_ver配置项定义在 emqx_bridge_mqtt_connector_schema.erl可选值为v3、v4、v5默认是v4。也就是说要把远端主题配置为$queue/...必须显式把连接器协议版本设为v5否则创建 Source数据源通道时会直接失败。测试验证非 v5 连接器拒绝队列订阅emqx_bridge_mqtt_source_SUITE中的t_mqtt_conn_bridge_rejects_queue_source_without_subscription_identifier用例验证了这条规则见 emqx_bridge_mqtt_source_SUITE.erl创建proto_ver v3的连接器后再尝试把 Source 的主题配置为$queue/orders/t/#HTTP API 返回400错误信息为queue subscriptions require connector proto_ver v5这说明队列订阅的约束在通道配置校验阶段就被拦截而不是等到连接远端 Broker 时才暴露。订阅标识符的分配与索引维护分配算法从 1 开始递增、避开已占用 ID每个 Source 通道在安装时都会通过maybe_attach_subscription_identifier/2被分配一个订阅标识符见 emqx_bridge_mqtt_ingress.erl。分配逻辑如下从1开始查找只要该 ID 尚未出现在subscription_id_to_handler_index这个 ETS 表中就分配给它上限为?MAX_SUBSCRIPTION_ID即268435455对应 MQTT 5 规范中 Subscription Identifier 的单字节编码上限见 emqx_bridge_mqtt_ingress.erl如果 ID 耗尽则抛出{no_available_subscription_id, SubscriptionId}错误。分配到的 ID 会连同完整通道配置一起写入 ETS 表ets:insert(SubscriptionIdToHandlerIndex, {SubscriptionId, Conf})见 emqx_bridge_mqtt_ingress.erl。注意$queue订阅不会被插入主题索引TopicToHandlerIndex因为队列消息不是按原始主题路由的而普通主题订阅则会写入emqx_topic_index便于按主题匹配见 emqx_bridge_mqtt_ingress.erl。订阅时携带标识符实际发起订阅时subscribe_properties/1会把通道配置中的subscription_id转换为 MQTT 5 的Subscription-Identifier属性随 SUBSCRIBE 报文一起发出见 emqx_bridge_mqtt_ingress.erlsubscribe_properties(#{subscription_id : SubscriptionId}) - #{Subscription-Identifier SubscriptionId}; subscribe_properties(_Ingress) - #{}.收到消息后按标识符分发远端消息到达后handle_publish/4会先尝试从 PUBLISH 报文的属性中提取Subscription-Identifier用它在 ETS 表中精确查找对应的通道配置只有当属性缺失或找不到匹配时才回退到按主题匹配见 emqx_bridge_mqtt_ingress.erl。这保证了$queue投递的消息能够准确路由到订阅它的那一条 Source 通道即使多条队列订阅绑定了重叠的主题过滤器也不会串线。远端 Broker 不支持时的两种降级行为MQTT 5 规范中如果 Broker 不支持 Subscription Identifier会在 SUBACK 中返回原因码0x9BSubscription identifiers not supported对应代码中的?RC_SUBSCRIPTION_IDENTIFIERS_NOT_SUPPORTED。入站桥接在maybe_retry_without_subscription_identifier/3中分两种情况处理见 emqx_bridge_mqtt_ingress.erl队列订阅$queue/...收到0x9B后直接返回错误{error, subscription_identifier_required_for_queue_subscription}——队列消息的投递离不开订阅标识符降级没有意义因此订阅失败普通主题订阅收到0x9B后自动重试这次不带Subscription-Identifier属性重新 SUBSCRIBE从而与不支持该特性的老版本 Broker 保持兼容。对应源码片段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;测试验证普通主题订阅在 0x9B 后重试t_mqtt_conn_bridge_ingress_retries_without_subid_on_a1用例通过 meck 拦截emqtt:subscribe/4模拟远端 Broker 对带标识符的订阅返回0x9B断言随后会观察到一次不带Subscription-Identifier的重新订阅见 emqx_bridge_mqtt_source_SUITE.erl。这一用例直接对应 changelog 中“regular topic subscriptions automatically retry without Subscription Identifiers”的行为描述。配置示例如何搭建一条队列订阅的 MQTT Source下面结合 emqx_bridge_mqtt_pubsub_schema.erl 中的 Source 结构给出一个完整可用的配置思路。MQTT Subscriber Source 通过mqtt类型的桥接创建其核心参数来自连接器配置与ingress_parameters两部分。连接器配置要点连接器Connector负责与远端 Broker 建立 MQTT 连接关键配置如下字段定义见 emqx_bridge_mqtt_connector_schema.erl配置项说明备注server远端 Broker 地址支持mqtt://、mqtts://等 scheme使用mqtts://时需同步开启 SSL否则启动即报错proto_ver协议版本v3/v4/v5队列订阅必须为v5默认v4clientid_prefixClient ID 前缀可选username/password连接认证凭据密码按 secret 类型处理clean_start是否干净会话默认true设为false时重启后会先恢复会话中的积压消息keepalive心跳间隔默认160sconnect_timeout连接超时默认10spool_size连接池大小普通主题订阅时只有首个 worker 真正订阅日志会提示mqtt_pool_size_ignored提示server地址若写成mqtts://TLS scheme但未开启ssl.enable连接器会在启动阶段直接拒绝并提示“Inconsistent server address and SSL settings”相关校验见 emqx_bridge_mqtt_connector.erl。Source 参数配置在连接器之上创建的 Source数据源通道其参数结构在ingress_parameters中定义见 emqx_bridge_mqtt_pubsub_schema.erl配置项默认值说明topic—远端订阅主题队列订阅写作$queue/{name}/{bind-filter}如$queue/orders/t/#qos—订阅 QoS配置为2时会被降级为1见parse_remote/2的downgrade_ingress_qos/1emqx_bridge_mqtt_ingress.erlno_localfalse是否禁止接收本客户端自己发布的消息retain_as_publishedtrue转发时是否保留消息的 retain 标志local#{}本地主题映射把远端消息渲染后发布到 EMQX 本地主题一个典型的队列订阅 Source 配置HOCON 形式示意如下bridges.mqtt.my_queue_source { connector mqtt:my_connector # 引用 proto_ver v5 的连接器 parameters { topic $queue/orders/t/# qos 1 } local { topic ingress/orders/${topic} } }创建成功后可以在连接器实例状态中看到该通道的ingress_list内带有自动分配的subscription_id这一点由测试用例t_mqtt_conn_bridge_ingress_subid_dispatch断言验证配置$queue/orders/t/#后通道状态中的subscription_id为1见 emqx_bridge_mqtt_source_SUITE.erl。消息流转从远端队列到本地主题整条链路的处理流程可以概括为连接器按pool_size启动 ecpool 连接池每个 worker 持有一个 MQTT 客户端Source 通道安装时subscribe_channel/2遍历所有 worker 并逐一发起订阅见 emqx_bridge_mqtt_ingress.erl。should_subscribe/5决定哪些 worker 真正订阅$share共享订阅所有 worker 都订阅普通主题订阅只有第一个 worker 订阅emqx_bridge_mqtt_ingress.erl远端 PUBLISH 到达后handle_publish/4依据报文属性中的Subscription-Identifier或回退到主题匹配找到通道配置handle_channel_config/3先触发on_message_received钩子把消息交给规则引擎/数据集成见 emqx_bridge_mqtt_connector.erl再按local映射把消息发布到 EMQX 本地主题maybe_publish_local/3emqx_bridge_mqtt_ingress.erl。断线重连时每个 worker 都会通过on_reconnect/2回调重新执行订阅逻辑emqx_bridge_mqtt_ingress.erl因此重连后队列订阅依然能带上正确的 Subscription Identifier。常见问题与排查建议创建 Source 报错queue subscriptions require connector proto_ver v5确认连接器的proto_ver已改为v5且连接器与 Source 属于同一桥接实例。远端 Broker 是 MQTT 5 但不支持 Subscription Identifier队列订阅会失败日志出现ingress_client_subscribe_failed与subscription_identifier_required_for_queue_subscription这是预期行为说明该 Broker 无法承载$queue队列投递请改用支持该特性的 Broker 或退回到普通主题订阅。普通主题订阅在旧版 Broker 上首次订阅失败入站桥接会自动去掉Subscription-Identifier重试一次无需人工干预可通过测试用例中的 meck 模拟方式在本地复现该降级路径。mqtt_pool_size_ignored警告当远端主题不是$share共享订阅时连接池中只有第一个 worker 订阅池子大小对消费并发没有影响这是设计使然。队列订阅与主题索引$queue订阅不会写入TopicToHandlerIndex因此不能通过“按主题匹配”的方式命中务必确保远端 Broker 在 PUBLISH 中回填了订阅标识符否则消息将因找不到通道配置而被丢弃。小结fix-16979这一变更让 EMQX 的 MQTT 入站桥接真正具备“消费远端消息队列”的能力只要远端 Broker 支持 MQTT 5 Subscription Identifiers就可以用$queue/{name}/{bind-filter}主题接入队列数据不支持时队列订阅会被明确拒绝避免静默错误普通主题订阅则自动降级重试兼顾了兼容性与可诊断性。核心实现集中在 emqx_bridge_mqtt_ingress.erl订阅、标识符分配、重试与分发与 emqx_bridge_mqtt_connector.erl索引生命周期、通道管理配套测试覆盖了拒绝路径、降级重试与按标识符分发三类关键场景可作为理解和验证该功能的首选参考。赞分享后端物联网消息队列通信【免费下载链接】emqxThe most scalable and reliable MQTT broker for AI, IoT, IIoT and connected vehicles项目地址https://gitcode.com/gh_mirrors/em/emqx点击查看免费下载相关推荐EMQX MQTT Ingress 桥接队列订阅$queue与 MQTT 5 Subscription Identifiers 深度解析EMQX MQTT Ingress 桥接队列订阅 $queue 与 MQTT 5 Subscription Identifiers 深度解析 MQTT in后端物联网消息队列通信EMQX MQTT 桥接 $queue/ 订阅消息丢失修复Subscription-Identifier 在 MQ 投递链路中的完整解析EMQX MQTT 桥接 $queue/ 订阅消息丢失修复Subscription Identifier 在 MQ 投递链路中的完整解析 本文基于仓库变更记录后端物联网消息队列通信EMQX 修复解析Message Queue 启用时 $queue/ 订阅无法向 MQTT 桥接投递消息的问题EMQX 修复解析Message Queue 启用时 $queue/ 订阅无法向 MQTT 桥接投递消息的问题 本文基于仓库变更记录 changes/ee/f后端物联网消息队列通信创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
