跳转至

从 MQTT 到 WebSocket:机器人状态如何实时抵达浏览器

设备用 MQTT 上报,管理后台用 HTTP 操作,浏览器靠 WebSocket 接收实时变化——这是机器人云平台常见的协议组合。真正容易出错的地方不是建立三条连接,而是消息跨越线程和协议以后,状态是否仍然一致。

本文以我的 Robot Cloud 演示项目为例,拆解一条消息从机器人到浏览器经历的关键边界。

先确定每种协议的职责

协议 适合处理 不适合承担
MQTT 设备遥测、连接状态、指令和回执 浏览器业务接口
HTTP 登录、设备 CRUD、分页查询、命令提交 高频主动推送
WebSocket 将设备事件实时推给已打开的管理页 作为设备消息 Broker

数据链路可以画成:

sequenceDiagram
    participant D as Device
    participant M as EMQX
    participant B as FastAPI
    participant W as Vue Browser

    D->>M: publish device/SN-01/telemetry
    M->>B: MQTT on_message
    B->>B: 解析并转换为统一事件
    B-->>W: WebSocket broadcast

这里 FastAPI 是协议翻译和业务边界,不应该让浏览器直接订阅所有设备 Topic。否则设备 ACL、用户权限、数据脱敏和前端协议都会绑死在 Broker 上。

Topic 要能表达方向和语义

演示项目采用以下 Topic:

device/{sn}/telemetry   # 遥测上行,QoS 0
device/{sn}/status      # 连接状态,QoS 1 + Retained
device/{sn}/cmd         # 指令下行,QoS 1
device/{sn}/cmd/ack     # 执行回执,QoS 1

status 使用 Retained 后,新订阅者可以立即得到设备最后已知状态。设备连接时还注册 Will 消息;如果进程被强杀或网络突然中断,Broker 在 keepalive 超时后代发 offline

需要注意:Will 表示连接异常结束,不是毫秒级在线检测。业务页面应显示“最后见报时间”,并根据场景设置合理的超时阈值。

同步 MQTT 回调与 asyncio 的桥接

paho-mqtt 的网络循环运行在后台线程,FastAPI 的 WebSocket 管理器运行在 asyncio 主循环。下面这种写法行不通:

def on_message(client, userdata, message):
    await manager.broadcast(...)  # 同步回调不能直接 await

也不应该在每条消息里随意 asyncio.run(),因为它会创建新的事件循环,而 WebSocket 对象属于 FastAPI 原来的循环。

项目启动时保存主循环:

main_loop = asyncio.get_running_loop()

MQTT 线程收到消息后,将协程安全投递回这个循环:

future = asyncio.run_coroutine_threadsafe(
    manager.broadcast(event),
    main_loop,
)

这一步的关键不是某个函数名,而是线程所有权:同步 SDK 可以保留自己的线程,但对异步资源的操作必须回到资源所属的事件循环。

数据库为什么也分同步与异步会话

HTTP 路由使用 SQLAlchemy AsyncSession,适合 FastAPI 请求生命周期;MQTT 回调处于同步线程,使用同步 Session 更直接:

HTTP / FastAPI loop  → aiomysql → AsyncSession
MQTT / paho thread   → pymysql  → SyncSession

两套会话共享同一数据库模型,但不共享 Session 对象。Session 不是线程安全的,也不应该跨请求长期保存。

另一种方案是 MQTT 回调只把消息写入线程安全队列,由主循环中的消费者统一异步落库和广播。这在消息量更大时更容易做批处理和背压,也是生产扩展时我更倾向的方向。

命令下发不能止于 publish

当浏览器调用:

POST /api/v1/devices/SN-01/commands

后端为命令生成唯一 ID,状态先记为 pending,然后发布 MQTT。设备 ACK 必须带回同一 ID:

{
  "command_id": "a1b2c3d4e5f6",
  "status": "success"
}

最终状态机是:

stateDiagram-v2
    [*] --> pending
    pending --> success: 收到成功 ACK
    pending --> failed: 收到失败 ACK
    pending --> timeout: 超过 15 秒

即使 MQTT 使用 QoS 1,仍然只能保证“至少一次”传递,业务上必须考虑重复消息。生产系统应让设备按 command_id 去重,后端 ACK 更新也要幂等。

浏览器重连需要退避和心跳

前端 WebSocket 封装采用 15 秒心跳和指数退避:

const delay = Math.min(1000 * 2 ** retry, 30_000)

立即无限重连会在服务恢复前制造连接风暴;固定 30 秒重连又会让短暂断网后的体验很迟钝。指数退避在两者之间取得平衡。生产环境还应加入随机抖动,避免大量客户端在同一时刻重连。

组件卸载时必须区分“用户主动关闭”和“网络异常关闭”,前者不应触发自动重连;重连成功后还应重新拉取一次 REST 快照,再继续消费增量事件,避免断线期间的数据缺口。

单进程广播的边界

演示项目的 ConnectionManager 保存当前进程内的 WebSocket 列表。这在单进程开发环境足够清楚,但部署多个 Uvicorn worker 后,不同连接可能分布在不同进程,某个 worker 收到 MQTT 消息时无法访问其他 worker 的连接。

可扩展结构应是:

MQTT consumer → Redis Stream / PubSub → each WebSocket worker → local clients

命令状态也应从内存字典迁移到 Redis 或数据库,才能在重启和多副本下保持一致。

小结

实时设备链路需要同时回答五个问题:

  1. 消息语义由哪个 Topic 和字段表达?
  2. 同步回调如何安全跨到异步事件循环?
  3. 断线、重连和状态快照如何收敛?
  4. 指令如何通过唯一 ID、ACK 和超时形成闭环?
  5. 单进程实现扩展到多副本时,状态放在哪里?

把这些问题解决以后,MQTT、WebSocket 和 FastAPI 才真正组成一个系统,而不是三段可以单独演示的代码。