资讯动态

EMQX订阅上下线监控实战:从MQTT事件到规则引擎与Webhook

发布时间:2026/9/15 2:55:21 来源:尧图企业网站定制
做物联网平台开发的朋友大概率都被同一个问题折磨过设备到底在不在线它有没有订阅到正确的主题它什么时候悄悄退订了某个关键主题办公室里几十台设备一起跑哪台掉线、哪个通道异常全靠用户打电话来报障才后知后觉。我第一次正经接触 EMQX 的订阅上下线监控就是被这种场景逼出来的——设备侧明明显示“已连接”服务端却一直等不到数据最后排查出来是设备发送 SUBSCRIBE 时用的 topic 拼错了一个字符。从那以后我就明白订阅事件和上下线事件必须作为平台的基础数据自己接一份不能指望“看起来连上了”。这篇文章我会沿着“为什么需要 → 事件从哪来 → 怎么把 Docker 里的 EMQX 跑起来 → 怎么用规则引擎把事件推到自己的服务 → 数据落到业务后有哪些坑”这条线把 EMQX 订阅上下线这件事完整讲透。不管你是正在搞设备接入平台还是只想在本地快速验证一下 EMQX 的事件能力这篇文章应该都能让你少走不少弯路。1. 先想清楚订阅上下线监控到底解决什么问题1.1 一个典型物联网平台的在线状态困境大多数物联网平台的第一步需求永远是随时知道设备是否在线。设备在线状态影响的是后续一切业务——指令下发要不要发、告警要不要触发、计费要不要暂停、OTA 任务要不要推送。很多平台的第一个版本用的是“应用层心跳”设备每隔 30 秒上报一个心跳主题服务端一旦超过 90 秒没收到就判定离线。这个方案本身没问题但它有一个致命缺陷心跳只能证明“设备还在上报”证明不了“设备已经成功订阅了某个业务主题”。举个我实际遇到的例子。某设备厂商的固件里写死了一系列订阅主题比如device/{sn}/command。后来他们升级了一次服务端把命令主题改成了device/{sn}/cmd但老设备固件没升级。设备心跳正常、连接正常、在线状态绿灯可平台下发的所有指令都石沉大海。这种问题单靠心跳完全看不出来只有把订阅事件也纳入监控才能第一时间发现“这个设备订阅了哪些主题、订阅是否和预期一致”。1.2 订阅上下线事件能带来哪些场景价值把 EMQX 的订阅事件和上下线事件沉淀成一条数据流之后很多从前要靠猜的问题就变成了可以精确回答的问题在线状态实时感知设备 CONNECT 成功、DISCONNECT 断开按 clientid 维护一张在线状态表比心跳轮询更即时而且不消耗设备额外流量。订阅关系审计每次 SUBSCRIBE / UNSUBSCRIBE 都记录下来出现“设备没收到命令”时可以直接回溯它到底订阅了哪个主题、用的什么 QoS、什么时候退订的。动态订阅管理有些设备不会一次性把主题订齐而是根据遥控指令、物料属性、告警状态动态增加或撤销订阅。这时候只有订阅事件流能帮你还原设备完整的行为轨迹。异常行为识别如果某个 clientid 在短时间内高频订阅/退订大量主题或者订阅了明显不属于它的系统主题往往意味着固件 bug甚至可能存在非法设备恶意图订阅攻击。1.3 为什么事件驱动比“轮询查询”更适合做这件事有人可能会问EMQX 有 Dashboard里面不也能看到客户端列表和订阅列表吗定时去 Dashboard 查一下不也行行但仅限于人工排障。平台系统想要自动处理就必须走事件驱动。轮询存在两个天然问题一是实时性差两个轮询周期之间状态可能已经变化了好几次二是拿不到“变化过程”只有“某一瞬的快照”。而事件流天然带有时间戳、clientid、主题、QoS、原因码这些明细能让你完整还原“设备何时上线、何时订阅、何时断开、断开原因是什么”。说白了订阅上下线事件就是 MQTT broker 的“审计日志”它的价值不在于单个事件本身而在于把这些事件串成线索去回答业务上的“为什么”。2. 搞清楚事件的源头MQTT 协议与 EMQX 内部机制2.1 从 CONNECT 到 DISCONNECT协议层的完整生命周期在进入 EMQX 配置之前先花点时间把底层协议讲清楚。MQTT 连接的生命周期并不复杂CONNECT客户端发起 TCP 连接后第一条报文必须是 CONNECT携带 clientid、username、password、keepalive、clean_session3.1.1或 clean_start5.0等参数。EMQX 校验通过后回复 CONNACK此时 broker 内部就产生了一次客户端上线事件。SUBSCRIBE客户端可以在连接成功后的任意时刻发送 SUBSCRIBE 报文订阅一个或多个主题过滤器并声明每个订阅的 QoS。broker 完成权限校验、与现有订阅合并后回复 SUBACK此时产生订阅事件。PUBLISH订阅关系建立后消息的路由就在 broker 内完成按主题和订阅关系派发给对应客户端。UNSUBSCRIBE客户端可以随时退订broker 清理订阅关系后回复 UNSUBACK此时产生取消订阅事件。DISCONNECT / 异常断开客户端主动发送 DISCONNECT 是“优雅下线”网络中断、心跳超时、被服务端踢除、broker 重启则属于“非优雅断开”。无论哪种EMQX 都会产生断开事件只是事件里的原因码和连接持续时间不同。需要特别注意连接建立和订阅建立是两件独立的事中间可能隔了很久也可能根本没有发生订阅。一个客户端可以只建连不订阅任何主题这在 EMQX 看来依然是一次正常上线。2.2 clean_session 与持久会话上下线事件最容易误判的地方MQTT 的 clean_session/clean_start 参数直接决定了断线重连后订阅关系是否保留。这一点在监控订阅上下线时特别容易栽跟头clean_session trueclean start连接断开时broker 清除该客户端的所有会话状态订阅关系也随之消失。下次重连需要重新发送 SUBSCRIBE所以你会看到“断开事件 → 上线事件 → 订阅事件”的完整链路。clean_session false持久会话broker 会保留会话状态包括订阅关系和离线消息。客户端重连后只要不主动重新 SUBSCRIBEbroker 也会恢复之前的订阅关系。此时你只会看到“断开事件 → 上线事件”不会看到新的订阅事件。如果你拿“每次上线后一定会收到订阅事件”来开发在线状态逻辑遇到持久会话的设备时就会漏掉数据。正确做法是订阅事件只当作订阅关系变化的痕迹不要当作订阅关系的全量快照。如果业务需要全量快照可以在设备上线后主动向设备下发一次“订阅上报”请求或者通过 EMQX 的 API 拉取当前订阅列表。2.3 遗嘱消息在上下线判定里的作用遗嘱消息Will Message是 MQTT 的另一个关键机制。设备在 CONNECT 时带上遗嘱主题和遗嘱 payload之后如果 broker 发现连接异常断开非主动 DISCONNECT就会代设备发布这条遗嘱消息。很多方案会把“在线/离线”状态也做成一份遗嘱比如遗嘱 topic 为status/{clientid}payload 为offline。但要注意遗嘱消息不能替代 disconnected 事件。它是发布到一个普通主题的普通消息要经过订阅关系匹配才能送达如果没有任何客户端订阅这个遗嘱主题遗嘱消息就“消失”了。而 EMQX 的client.disconnected事件是 broker 内部事件无论有没有订阅者都会产生。所以做订阅上下线监控正确姿势是用client.connected/client.disconnected事件维护在线状态表用client.subscribed/client.unsubscribed事件维护订阅关系表遗嘱消息可以当作一条“离线广播”推给业务侧或者推给其他设备做联动但不应该作为唯一的数据源。2.4 EMQX 暴露事件的几条路径EMQX 把各类事件暴露给使用者的方式主要有三种理解它们各自的适用场景才能选对方案事件路径原理适合场景$SYS系统主题EMQX 在内部主题树上发布连接、断开、订阅、退订等事件普通 MQTT 客户端可以订阅这些主题本地调试、快速验证、人工巡检规则引擎通过预设的 SQL 匹配事件然后触发动作Webhook、消息持久化、消息重发布等生产环境对接业务系统最推荐统计 API / Dashboard通过 HTTP API 查询当前客户端、订阅列表低频管理操作、页面展示不依赖实时事件流其中$SYS主题只适合“人眼观察”不适合做业务数据源因为它的实时性、稳定性都没有保证而且不同 EMQX 版本的路径和格式有差异。生产环境我强烈建议用规则引擎 Webhook或规则引擎 消息队列。3. Docker 快速起一个 EMQX先把系统主题订阅出来3.1 一行命令启动 EMQX 5.x如果你只是为了验证订阅上下线本地用 Docker 跑一个 EMQX 是最省事的。我用的是 5.x 版本相对 4.x 的配置方式和事件格式都有变化所以这里以 5.x 为例。docker run -d --name emqx \ -p 1883:1883 \ -p 8083:8083 \ -p 8084:8084 \ -p 18083:18083 \ emqx/emqx:5.7.1端口说明1883MQTT 标准 TCP 端口8083MQTT over WebSocket 端口8084MQTT over TLS/SSL 端口18083Dashboard HTTP 端口启动完成后访问http://localhost:18083默认用户名admin默认密码public登录后建议马上改掉。如果是给本地长期开发环境用第一次启动还应该挂载数据卷否则容器重建后所有配置、证书、数据都丢了。我习惯用docker run -d --name emqx \ -p 1883:1883 -p 18083:18083 \ -v emqx-data:/opt/emqx/data \ -v emqx-etc:/opt/emqx/etc \ -v emqx-log:/opt/emqx/log \ emqx/emqx:5.7.13.2 订阅 $SYS 事件主题直接在 MQTTX 里观察事件启动好 EMQX 后我常用 MQTTX 这个客户端工具来观察事件流。先新建一个连接clientid 随便填比如monitor-001然后订阅以下主题过滤器$SYS/brokers//clients//connected $SYS/brokers//clients//disconnected $SYS/brokers//clients//subscribed $SYS/brokers//clients//unsubscribed再另起一个客户端命名为device001连接上来后订阅一个业务主题比如sensor/001/data。这时候回到第一个“监控客户端”你会看到类似这样的事件消息{ clientid: device001, username: device_user, ts: 1723024800000, sockname: 172.17.0.2:1883, peername: 172.17.0.1:54321, proto_ver: 4, keepalive: 60, connected_at: 2024-08-07T10:00:00.00008:00 }订阅事件大概是{ clientid: device001, username: device_user, ts: 1723024801000, topic: sensor/001/data, qos: 1 }当你断开device001时监控客户端会收到 disconnected 事件。你立刻就能直观理解每一次设备行为在 EMQX 里都有对应的事件痕迹。3.3 系统主题 payload 里的关键字段怎么读以 EMQX 5.x 实际输出为例我整理过一份常用字段对照字段出现的事件含义实际使用建议clientid全部客户端唯一标识永远用这个关联设备username全部连接时用户名可为空不能当唯一标识可伪造ts全部事件发生时间戳毫秒用于排序和审计peernameconnected/disconnected客户端 IP:端口判断来源网络socknameconnected/disconnected服务端监听地址多网卡时有用proto_verconnectedMQTT 协议版本 3/4/5兼容性排查keepaliveconnected心跳保活间隔秒判断连接质量connected_atconnected连接建立时间在线时长计算disconnected_atdisconnected断开时间在线时长计算topicsubscribed/unsubscribed订阅/退订的主题订阅关系维护qossubscribed订阅的 QoS 等级判定订阅是否生效reasondisconnected5.x断开原因码异常分析关键reason字段在线状态监控里价值很大。EMQX 5.x 的断开事件里reason会区分normal主动断开、keepalive_timeout心跳超时、tcp_closedTCP 被关闭、kicked被服务端踢除等。通过它你能快速判断设备是“正常下线”还是“异常掉线”。3.4 为什么不建议直接用 $SYS 主题做生产数据源看到这里你应该已经理解了$SYS主题的方便但我必须泼一盆冷水别在生产环境用它。原因有三$SYS主题在不同 EMQX 版本之间变化很大。4.x、5.0、5.7 的事件路径和 payload 格式都不完全一致升级可能导致解析逻辑崩掉。$SYS消息没有严格的投递保障它本身也是走 broker 消息路由的流量大时存在被丢弃或延迟的可能。它把 broker 内部事件暴露给了所有能订阅的客户端。如果权限控制不严格普通设备也能订阅到别人的上下线事件属于越权风险。所以$SYS只适合刚上手时的功能验证生产环境的订阅上下线监控请转向下一章要讲的规则引擎。4. 用规则引擎 Webhook 把订阅上下线事件推送到自己的服务4.1 规则引擎的基本逻辑事件 → SQL → 动作EMQX 的规则引擎可以理解成一条流水线输入是内部事件或消息中间经过一段 SQL 做字段提取和过滤输出是执行动作。处理客户端上下线、订阅退订这类事件就是典型的规则引擎用法。EMQX 5.x 内置了以下与订阅上下线相关的事件类型client.connected客户端连接成功client.disconnected客户端断开连接client.subscribed客户端订阅成功client.unsubscribed客户端取消订阅session.subscribed/session.unsubscribed会话级别订阅变化持久会话恢复时触发client.authorization授权检查结果可用于非法订阅审计在 Dashboard 的“规则”页面里选择对应事件类型编写 SQL再关联一个“动作”就能把事件转发到外部系统。最常见的动作就是 Webhook把事件用 HTTP POST 推到自己的后端。4.2 创建四类事件规则的分步配置我在生产项目中通常一次性创建四条规则分别处理连接、断开、订阅、退订。以“客户端连接事件”为例配置过程如下登录 EMQX Dashboard进入“管理 → 规则”。点击“创建”。规则名称填event-client-connectedSQL 处输入SELECT clientid, username, peername, proto_ver, keepalive, timestamp() AS event_time FROM $events/client_connected这里timestamp()是规则引擎内置函数会生成一个毫秒时间戳方便后面落库排序。如果你还需要原始 payload 的其他字段可以直接SELECT *但在生产环境我建议只挑需要的字段减少无效数据传输。在“添加动作”里选择“Webhook”配置请求 URLhttp://你的后端地址/api/emqx/eventHTTP 方法POSTHeadersContent-Type: application/jsonBody保持默认模板即可规则引擎会把我 SQL 选出的字段以 JSON 格式 POST 出去。点击“创建”后再手动确认规则已启用。同样的方式再创建client.disconnected、client.subscribed、client.unsubscribed三条规则。client.subscribed的 SQL 类似SELECT clientid, username, topic, qos, timestamp() AS event_time FROM $events/client_subscribedclient.unsubscribed把事件名换成$events/client_unsubscribed即可。4.3 写一个本地 HTTP 接收服务Python Flask 示例现在需要一个能接收 Webhook 的服务。这里给一个最简单的 Python Flask 示例方便联调用。你也可以用 Node.js、Spring Boot、Go 实现同样的接口逻辑都一样。先装依赖pip install flask flask-cors创建event_receiver.pyfrom flask import Flask, request, jsonify app Flask(__name__) app.route(/api/emqx/event, methods[POST]) def handle_event(): event request.get_json(forceTrue) # 这里先打印出来后续可以替换成写数据库、发告警等业务操作 print( EMQX Event ) print(json.dumps(event, ensure_asciiFalse, indent2)) return jsonify({code: 0, message: ok}) if __name__ __main__: app.run(host0.0.0.0, port5000)运行后你的本地服务就监听在5000端口。需要注意一个非常经典的坑如果你的 EMQX 跑在 Docker 容器里Webhook 地址里的“本机”不是127.0.0.1而是宿主机地址。在 Linux 上可以查宿主机局域网 IP在 Mac/Windows 的 Docker Desktop 里可以用host.docker.internal代替。否则EMQX 容器收到的请求目标127.0.0.1:5000指向的是它自己永远打不到你的 Flask 服务。如果你不想处理容器网络也可以把 EMQX 直接以--network host模式启动docker run -d --name emqx --network host emqx/emqx:5.7.1这样 EMQX 直接使用宿主机网络Webhook 地址写http://127.0.0.1:5000/api/emqx/event就能通。4.4 联调从两个客户端看完整事件链路服务端和规则引擎都就绪后我们来完整走一遍链路。我用 MQTTX 新建两个连接monitor-device连上后订阅sensor/001/datadevice001连上后订阅sensor/001/data然后退出当device001连接时Flask 后端会收到一条client.connected事件当它发出订阅时收到client.subscribed事件当它主动断开时收到client.disconnected事件。整个过程如下[事件1] client.connected clientiddevice001 peername192.168.1.20:51234 [事件2] client.subscribed clientiddevice001 topicsensor/001/data qos1 [事件3] client.disconnected clientiddevice001 reasonnormal这里有一个细节即使客户端订阅的是一个重名主题EMQX 仍然会产生订阅事件。如果device001已经订阅了sensor/001/data再次发送相同 SUBSCRIBE 报文EMQX 会更新订阅比如改 QoS事件同样会触发。这说明订阅事件反映的是“订阅操作发生”而不是“订阅关系新增”。5. 真实项目里处理订阅上下线数据的五个关键细节5.1 clientid 才是设备的唯一身份username 只是装饰把事件接到业务系统后第一个要确定的关联字段就是clientid。在 MQTT 协议里clientid 是区分客户端的唯一标识同一时刻同一 clientid 只能有一个连接存在。username 虽然业务语义更强但它不具备唯一性——多个设备可以共用同一个 username也可以完全留空。所以所有设备状态表、订阅关系表的主键都应该围绕 clientid 设计。如果业务侧有自己的设备编号比如 SN建议在设备接入时把 SN 直接放进 clientid比如sn:ABC123这样事件进来后不需要再做映射。5.2 断开原因码区分正常断开、踢下线、网络异常EMQX 5.x 的client.disconnected事件里reason字段能告诉你断开的具体原因。常见取值reason含义建议处理normal客户端主动发送 DISCONNECT正常离线不做告警keepalive_timeout心跳超时连接被 broker 关闭记录异常按需告警tcp_closedTCP 连接被关闭未收到 DISCONNECT往往伴随网络闪断关注频次kicked被管理员或同 clientid 新连接踢下线注意是否发生 session takeoverbroker_shutdownbroker 正在关闭导致连接断开平台维护时会成片出现protocol_error客户端发送了非法报文大概率是设备固件 bug我在实际项目中会写一个简单的分类逻辑normal不计告警keepalive_timeout和tcp_closed连续出现才告警kicked则直接关联 session takeover 逻辑。如果不看 reason只看到大量的 disconnected 就发告警那么一次网络抖动就能把告警渠道打爆。5.3 重复上线与 session takeover同一 clientid 互踢MQTT 协议规定如果客户端 A 已经连接此时具有相同 clientid 的客户端 B 再连接broker 会踢掉 A让 B 接管会话。这个过程叫 session takeover。在事件流里你会看到[事件] client.disconnected clientiddevice001 reasonkicked [事件] client.connected clientiddevice001 peername另一个IP如果平台同时有多个设备误配了相同 clientid就会导致设备反复上下线每次上线都会踢掉前一个连接。线上排查这类问题时reasonkicked就是最重要的线索。另一个排查点是看 disconnected 事件和 connected 事件的时间间隔。正常重连通常有几十毫秒到几秒的间隔如果间隔几乎为 0基本可以断定是同 clientid 互踢。处理方式也很直接在设备注册时强制生成唯一 clientid或在前置接入层校验 clientid 是否已经被占用并拒绝新连接。5.4 事件风暴防护批量断网时的告警轰炸与缓存策略做订阅上下线监控最大的考验不是单个事件而是瞬时事件风暴。比如某个区域停电、弱电井里交换机跳闸往往几十台设备一起掉线。如果你把每一条 disconnected 事件都直接转成告警、直接写入在线状态字段后端很可能会被打爆微信/短信告警也会变成轰炸。我的经验是把事件处理分成两段接收端做缓冲和聚合Webhook 接口收到事件后先写入本地队列或消息队列比如 Redis List、RabbitMQ、Kafka不做同步的慢操作。消费者做状态合并按 clientid 去重仅维护最新状态告警规则做成“连续 N 条异常才告警”而不是单条触发对 disconnected 风暴增加时间窗口统计完成后再批量推送。前端展示的状态也不需要每一条事件都刷新可以做成 5 秒或 10 秒聚合一次减少数据库写入压力。5.5 动态订阅场景设备反复订阅/退订如何聚合统计物联网设备里有一种很常见的“动态订阅”玩法设备初始只订阅一个控制主题收到平台下发的某个指令后再临时订阅另一个数据主题完成后退订。这种场景下订阅事件流是观察设备行为轨迹的唯一手段。如果你要统计“每个主题当前有多少活跃订阅者”直接在订阅事件上累计加一、在退订事件上减一是不可靠的——因为存在持久会话恢复的情况。更稳妥的做法是从 EMQX 的 HTTP API定时拉取一份订阅列表作为基准快照再用订阅/退订事件流做增量修正。这样两边的误差都可以校准。另外如果设备出现“反复订阅同一个主题、间隔非常短”的行为大概率是代码里循环调用了订阅接口。订阅事件里的ts字段能帮你精确还原这种行为模式从而定位到固件代码的哪一段循环逻辑出了问题。6. 踩坑记录与实用经验汇总6.1 用 Docker 部署 EMQX 时最容易踩的坑第一个坑是端口映射不全。有些老教程只映射了1883和18083导致你用 WebSocket 连接时一直不通查了半天发现8083没映射出来。如果你要用 MQTTX 的 WebSocket 模式记得把8083一起映射出来。第二个坑是 Dashboard 的初始密码。EMQX 5.x 默认用户admin/ 密码public如果部署在公网服务器上并且忘了改等于把管理后台裸奔在公网。建议在初始化脚本里就强制覆盖默认密码。第三个坑是配置文件挂载。如果你用-v挂载了/opt/emqx/etc在升级 EMQX 镜像时旧配置文件可能和新版本不兼容导致启动失败。我的建议是开发环境随便挂生产环境尽量用官方推荐的方式管理配置避免整个 etc 目录直接挂载。6.2 系统主题和事件字段的版本差异我见过很多人在网上抄了一段$SYS主题订阅的代码结果在 EMQX 5.0 上根本收不到消息。原因就是 4.x 的主题路径和 5.x 不一样4.x 常见的是$SYS/brokers/${node}/clients/${clientid}/connected5.x 把它移到了$SYS/brokers/${node}/clients/${clientid}/connected并调整了字段格式。最稳妥的方式是先到 Dashboard 的“问题排查 → 系统主题”里查看当前版本实际暴露的事件主题路径然后以实际为准。规则引擎的事件名也是一样。client.connected、client.disconnected、client.subscribed、client.unsubscribed在 4.x 和 5.x 都有但session.subscribed这类事件可能只在特定版本可用。我建议写规则前先建一条最基础的消息转发规则把事件输入打出来确认字段名后再写正式 SQL。6.3 我自己实际跑下来的几个建议最后分享几条我在这类项目里沉淀下来的经验不一定写在哪本手册里但非常实用不要把在线状态表建在关系型数据库里做高频更新。订阅上下线事件每分钟可能有几百上千条每次都 UPDATE 数据库行压力很大。我习惯用 Redis 维护在线状态TTL 设置成心跳周期的两倍数据库只保留审计日志。Webhook 接收接口一定要做幂等。EMQX 规则引擎的 Webhook 动作在网络异常时可能重试同一事件有可能被推送两次。后端处理时建议按clientid 事件类型 ts生成唯一键重复消息直接丢弃。规则引擎 SQL 里不要用SELECT *走生产。事件 payload 里有很多调试字段全部透传到业务端既浪费带宽也容易在日志里泄露内网地址。选取必要字段并在 Flink/数据库里做字段白名单校验。上线和订阅事件不是严格有序的。在分布式环境下client.subscribed和client.connected到达业务端的时间顺序可能颠倒处理时不要假设“先收到上线再收到订阅”。正确做法是收到订阅事件时如果在线状态表里还没有该 clientid先触发出一次“已存在但未上线”的初始化而不是丢弃事件。如果你只需要在线状态变化通知也可以考虑用 MQTT 遗嘱消息做广播但不要把业务核心逻辑完全寄托在上面。EMQX 事件流的数据完整性和调试便利性是遗嘱消息无法替代的。这套“订阅上下线”监控方案我后来复制到了不止一个项目里。从一开始的裸奔状态到后来每次设备在 Dashboard 上绿灯、但业务数据不出现时直接查订阅事件流就能锁定问题排障时间从“小时级”降到了“分钟级”。EMQX 订阅上下线事件真正给了我这种底气设备的一举一动只要发生必有痕迹。

读完文章,也想定制专属网站?

尧图设计师 24 小时内与您沟通定制方案

免费获取报价