EMQX MQTT 桥接陈旧连接状态修复解析:从「假 Connected」到真实健康检查与自动重连
2026/9/24 3:53:22 网站建设 项目流程
  • 后端
  • 物联网
  • 消息队列
  • 通信

【免费下载链接】emqx

The most scalable and reliable MQTT broker for AI, IoT, IIoT and connected vehicles

项目地址:https://gitcode.com/gh_mirrors/em/emqx
点击查看免费下载

本文围绕 EMQX 开源仓库中 changes/ee/fix-15603.en.md 记录的缺陷修复展开:当 MQTT 桥接(MQTT Connector / Bridge)的底层连接已失效(stale connection)时,连接器状态却仍显示为Connected,且连接不会自动重新建立。文章将结合 emqx_bridge_mqtt_connector.erl 等源码与对应测试用例,剖析状态上报的判定逻辑、健康检查与自动重连机制,并给出可直接落地的配置与运维建议。读完本文,你将理解「连接状态显示 Connected 但实际已断」这一问题的根因与修复思路,掌握通过健康检查参数、重连回调与测试手段确保桥接连接真实可用的方法。

一、问题背景:桥接显示 Connected,数据却不再流动

MQTT 桥接是 EMQX 与另一台 MQTT Broker 之间打通消息通道的核心能力,既可作为数据源(ingress/source,订阅远端主题并导入本集群),也可作为数据出口(egress/action,将本集群消息发布到远端)。桥接依赖一条常驻的 TCP/MQTT 长连接,连接的健康状况直接决定数据链路是否可用。

fix-15603 修复的问题可以用一句话概括:当桥接的底层连接已经失效(例如远端 Broker 异常关闭、网络中断且未收到 FIN/RST,或连接进程被异常终止)时,连接器状态仍然显示为Connected,并且系统不会主动重新建立连接,导致消息持续静默丢失,而运维侧从 Dashboard / API 看到的状态却是一切正常,难以定位。

为什么会出现「假 Connected」

从当前源码结构看,连接器状态并非由单一事件驱动,而是依赖「资源健康检查」周期性轮询每个 MQTT 客户端进程后聚合得出:

  • emqx_bridge_mqtt_connector.erl 中的on_get_status/2通过emqx_utils:pmap/3对连接池内所有 worker 并行执行get_status/1,超时窗口为?HEALTH_CHECK_TIMEOUT = 1000毫秒;
  • 每个 worker 的状态由 emqx_bridge_mqtt_ingress.erl 的status/1判定:它调用emqtt:info(Pid)socket字段,socketundefined则返回connected,否则返回connecting,进程已不存在则捕获exit:{noproc, _}返回disconnected

也就是说,状态判定依赖emqtt客户端进程内的 socket 信息与进程存活情况。如果连接实际已断但emqtt进程内没有及时感知(例如连接被对端静默丢弃、处于半开状态),或者旧连接进程的状态未被正确清理,健康检查就可能看到「进程存活 + socket 未清空」的假象,从而继续上报Connected,同时连接又不会被重新触发建立。这正是该缺陷的本质:状态显示与连接实际可用性脱节,且缺少兜底的重连触发路径

二、修复内容解读:让状态回归真实,让重连真正发生

fix-15603 的修复目标从 changelog 描述看非常明确:不再把失效连接显示为Connected,并且要重新建立连接。结合当前仓库代码,修复后的行为由以下几条机制共同保障:

  1. 健康检查结果可信on_get_status/2会真实探测每个 worker 的emqtt客户端状态,disconnected状态拥有最高优先级(见下文状态聚合),不再可能被「旧连接遗留信息」掩盖;
  2. 断连即重连emqtt客户端本身具备自动重连能力,连接器在 start_mqtt_clients/3 中为连接池显式配置了{auto_reconnect, ?AUTO_RECONNECT_INTERVAL_S},其中?AUTO_RECONNECT_INTERVAL_S = 2,即底层客户端会以 2 秒为间隔自动尝试重连;
  3. 重连后恢复订阅:对于 ingress(数据源)方向,重连成功后通过 emqx_bridge_mqtt_ingress.erl 注册的on_reconnect/2回调重新执行远端主题订阅,保证断连期间丢失的订阅关系在恢复后自动补齐。

需要说明的是:本文档对应的具体代码变更 diff 不在当前仓库快照内,上述机制是基于当前源码实现对该修复目标如何达成的推断性解读;但「陈旧连接不得继续显示为 Connected」「必须自动重建连接」这两点是 changelog 明确记载的事实,也是下文将要展开的源码机制的最终目标。

三、连接状态是如何上报的:健康检查链路与状态聚合

3.1 连接器资源的状态回调

MQTT 连接器实现了emqx_resource行为,on_get_status/2是资源框架周期性调用(间隔由resource_opts.health_check_interval控制)的状态探针。其执行链路为:

on_get_status/2 └─> emqx_utils:pmap(get_status/1, Workers, 1000ms) └─> ecpool_worker:client(Worker) % 取出 emqtt 客户端进程 └─> emqx_bridge_mqtt_ingress:status(Client) └─> emqtt:info(Pid) 检查 socket 字段 └─> combine_status/3 聚合所有 worker 结果

关键实现见 emqx_bridge_mqtt_connector.erl:

  • 健康检查超时上限 1000ms,若pmap超时,整体直接返回connecting状态(避免检查动作本身阻塞资源框架);
  • 任一 worker 取不到客户端进程,该 worker 即视为disconnected
  • 如果连接器未分配任何可用 clientid,状态会被标记为{disconnected, {unhealthy_target, ...}},并在 on_get_channel_status/3 中使通道快速失败并触发告警。

3.2 状态聚合规则:disconnected 优先

combine_status/3定义了多 worker 场景下的状态合并规则,注释中明确给出了自然序:[connected, connecting, disconnected],即:

  • disconnected权重最高,任何 worker 断开都会让连接器整体显示为disconnected(或携带具体原因);
  • connecting高于connected,只要有 worker 处于重连中,状态就不会显示为已连接;
  • 只有全部 worker 均健康时才显示connected

这一规则与修复目标直接相关:只要有任何一条桥接连接处于断开或重连状态,对外呈现的状态就不再是Connected,从机制上杜绝了「部分连接已断、整体仍显示已连接」的陈旧状态。

3.3 状态原因透传

combine_status/3还会把底层错误原因透传出来:explain_error/1econnrefusedtcp_closedframe_parse_error等常见错误映射为人类可读的说明文案(见 emqx_bridge_mqtt_connector.erl),最终通过 API 返回status_reason字段。运维可以从状态原因中直接看出是「连接被拒」「监听器已达上限」还是「对端返回了非 MQTT 数据」,而不再面对一个无法解释的Connected

四、自动重连机制:断线后如何恢复

4.1 客户端级自动重连

连接池的每个 worker 对应一个emqtt客户端进程(connect/1),连接器在启动连接池时传入{auto_reconnect, 2},因此客户端断开后会自动以 2 秒间隔重建连接。测试用例 t_reconnect 验证了这一点:通过 HTTP 接口强制踢掉连接池中的部分客户端连接进程后,emqtt客户端会自动重新连接,连接池 worker 数量最终恢复为初始pool_size

4.2 ingress 重连回调:重连后恢复订阅

对于数据源通道,仅重建 TCP/MQTT 连接还不够——远端订阅关系同样需要在重连后恢复。连接器在添加 source 通道时调用emqx_bridge_mqtt_ingress:add_reconnect_callback/2(见 emqx_bridge_mqtt_connector.erl),为连接池内每个 worker 注册重连回调;断线重连后on_reconnect/2会基于保存的ingress_config重新执行subscribe_channel_helper/5,即重连成功即重新订阅远端主题(emqx_bridge_mqtt_ingress.erl)。

源码 SUITE 中 t_reconnect_with_session 与「重连后重新订阅」相关用例对此有专门验证;?tp(debug, "mqtt_source_reconnected", ...)事件也作为可观测埋点出现在日志中,可用于确认重连是否完成。

4.3 clean_start = false 时的会话消息防丢

一个值得注意的细节:当连接器配置clean_start = false(复用远端会话)时,如果 MQTT 客户端启动前 topic handler 索引尚未建立,重连后远端会话中积压的消息到达时会找不到对应 handler 而被丢弃。为此 maybe_add_sources_with_sessions_to_topic_handler/3 会在启动客户端之前把带会话的 source 主题预注册到 handler 索引中,且该操作为幂等操作。这保证了「恢复连接 → 会话消息重放 → 正确路由」的完整闭环。

4.4 发送侧的容错与重试

egress(action)方向的可靠性由资源框架的缓冲队列保障:断连期间的消息进入队列,恢复后继续投递。测试 t_mqtt_conn_bridge_egress_reconnect 完整演示了这一过程:停掉本地 1883 监听器模拟断连 → 发布消息使其入队(指标queuing + inflight == 2)→ 重启监听器 → 连接器状态恢复connected→ 队列中消息全部送达且failed保持为 0。异步模式 t_mqtt_conn_bridge_egress_async_reconnect 也有相同结论。

此外,classify_error/1 将disconnectedecpool_emptytcp_closedclosed等归类为recoverable_error(可重试),而frame_parse_error、未识别错误等归类为unrecoverable_error;t_publish_while_tcp_closed_concurrently 专门构造了「健康检查判定健康的同时连接被强制关闭」的竞态,断言系统会触发重试而非误判。

五、与修复相关的连接器配置详解

MQTT 连接器的完整配置 schema 定义在 emqx_bridge_mqtt_connector_schema.erl,其中与连接健康、重连、状态上报直接相关的参数如下:

配置项默认值说明
server必填远端 Broker 地址,支持host:portmqtt://mqtts://形式
pool_size1连接池大小,即并行的 MQTT 客户端数量;ingress/egress 可各自覆盖
proto_verv4MQTT 协议版本:v3/v4/v5
clean_starttrue是否使用干净会话;false可复用远端会话(配合消息防丢逻辑)
keepalive160s心跳间隔,emqtt客户端会以force_ping主动探测对端
connect_timeout10s单次连接建立超时
retry_interval15sQoS 1/2 消息重发间隔(emqtt客户端层)
max_inflight32未确认的最大在途 QoS 消息数
bridge_modefalse桥接模式标志(MQTT 3.1.1 时代的遗留项,v5 下无效并告警)
clientid_prefix客户端 ID 前缀,最长 19 字节以保证拼接后 ≤23 字节
static_clientids[]按节点静态指定 clientid(含用户名/密码),用于集群确定性分配
username/password连接远端 Broker 的认证凭据
resource_opts.health_check_interval框架默认健康检查周期,测试中常设为500ms以加快状态收敛
reconnect_interval已废弃自 5.0.16 起废弃,自动重连间隔由客户端内置(2 秒)控制

要点解读:

  • health_check_interval直接决定「假 Connected」的暴露速度。健康检查是周期性的,间隔越短,断连状态越早被聚合上报,也就越早触发后续的告警与处置。生产环境建议结合监控告警阈值合理设置(如 15s~60s),测试环境可设500ms加速验证。
  • clean_start = false与 egress action 存在已知注意点:连接器在 on_add_channel/4 中会对此组合发出mqtt_publisher_clean_start_false告警——如果该 clientid 在远端已有订阅,重连后可能收到未被本端 handler 处理的消息,需提前规避。
  • static_clientids用于集群多节点桥接:通过 find_my_static_clientid_info/1 按本节点分配 clientid,未分配 clientid 的节点会通过on_get_channel_status快速失败并告警,避免无连接可用的节点继续收消息。

六、问题定位与修复验证实践

6.1 从日志与 API 确认状态

  • API 状态字段GET /api/v5/connectors/{name}返回statusstatus_reason;修复后断连时会返回disconnectedconnecting,而不再误报connected(测试断言见 emqx_bridge_mqtt_action_SUITE.erl)。
  • 可观测事件emqtt客户端启动/失败会输出ingress_client_startingingress_client_connect_failed日志,并附带explain字段;断连恢复场景可检索mqtt_source_reconnectedtrace 事件。
  • 监控指标:连接器/通道指标中的connecteddisconnected状态切换以及failedqueuinginflightretried等计数(见测试中对get_action_metrics_api的断言)可用于衡量断连期间的真实影响面。

6.2 复现与验证步骤(参考测试用例)

参照 t_reconnect 与 t_mqtt_conn_bridge_egress_reconnect 的思路,可以自行搭建验证环境:

  1. 创建 MQTT 连接器并关联一个 action/source 通道,health_check_interval设为较小值(如500ms);
  2. 停掉远端 Broker(测试中为停掉本地tcp:default监听器),观察连接器 API 状态应在connecting/disconnected之间切换,且status_reason给出明确原因;
  3. 在断连期间向本地发布消息,确认消息进入队列(queuing + inflight > 0)而非直接失败;
  4. 恢复远端 Broker,确认状态自动回到connected、队列消息全部投递成功(success增长、failed == 0)、ingress 订阅自动恢复。

6.3 升级与兼容性提示

该修复随 EMQX 5.x 与 6.x 版本线发布:changes 目录中 changes/e5.10.1.en.md 与 changes/e6.0.0.en.md 均收录了同一条修复记录(分别对应 5.10.1 与 6.0.0 版本线)。如果你的集群中 MQTT 桥接曾出现「状态显示 Connected 但数据不通」的现象,升级到包含该修复的版本后,建议在升级窗口内重点观察桥接状态字段与重连日志,确认行为符合预期。

七、小结

fix-15603 修复的本质是让 MQTT 桥接的连接状态与底层连接的真实可用性保持一致,并为失效连接补上自动重建的通路。从当前仓库源码看,这一目标通过三条线共同实现:

  1. 健康检查聚合on_get_status/2emqtt:info/1探测 +combine_status/3的「disconnected 优先」规则,杜绝陈旧连接继续显示为Connected
  2. 客户端自动重连:连接池 worker 内置 2 秒间隔的auto_reconnect,断线后自行恢复;
  3. 重连后的状态恢复:ingress 通过on_reconnect回调恢复远端订阅,clean_start = false场景预注册 handler 索引防止会话消息丢失,egress 由资源队列缓存断连期消息并在恢复后补投。

对于自建 MQTT 桥接的用户,本文提供的配置参数表、状态聚合规则与测试验证步骤,可以直接迁移到自己的集群排障与升级验收流程中。相关实现细节可继续查阅 emqx_bridge_mqtt_connector.erl、emqx_bridge_mqtt_ingress.erl 及 emqx_bridge_mqtt_action_SUITE.erl 中的对应用例。

  • 后端
  • 物联网
  • 消息队列
  • 通信

【免费下载链接】emqx

The most scalable and reliable MQTT broker for AI, IoT, IIoT and connected vehicles

项目地址:https://gitcode.com/gh_mirrors/em/emqx
点击查看免费下载

相关推荐

创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询