【WebSocket+MQTT】day60-mqtt

作者:mario 发布时间: 2026-09-03 阅读量:5 评论数:0

Day 60 · MQTT 与工业物联网协议 — 发布/订阅的世界观,QoS 的保证书

这周你的服务器一直在"点对点"推数据——它必须知道每个客户端要什么、连没连、断了怎么补。而工业现场的世界观完全不同:设备只管往 Topic 上发布,谁在订阅、有多少人订阅,设备根本不知道也不关心。这就是 MQTT 的发布/订阅模型,工业物联网的事实标准(PLC、网关、阿里云 IoT、AWS IoT 全都原生支持)。今天理解它的模型与 QoS 机制,并用 mqtt.js 把浏览器接上真 Broker——BOSS 战的数据链路从此有了"工业正统性"。


目录


一、为什么工厂设备侧都用 MQTT

1.1 工业现场的网络现实

现实 MQTT 的应对
设备在矿井/厂房角落,网络弱且常断 QoS 1/2 协议级保证送达 + 会话保持
设备是嵌入式 MCU,内存以 KB 计 协议头最小 2 字节,实现极简
几百个设备、几十个应用,关系复杂 发布/订阅彻底解耦(互不认识)
断电/断网是常态不是异常 遗嘱消息(LWT)自动广播"设备掉线"

1.2 和裸 WebSocket 点对点的本质差异

自建 WS 推送(本周前几天):
  服务器 ──必须维护──▶ 每个客户端的订阅表、状态、补传逻辑(全要自己写)

MQTT:
  设备 ──publish──▶ Broker(订阅表/QoS/离线消息全由 Broker 负责)◀──subscribe──▶ 应用
  你的服务端代码量断崖式下降

二、核心概念:Broker / Topic / 发布订阅

2.1 三要素

┌────────┐  publish "factory/line1/furnace-1/temp"  ┌────────┐
│ 加热炉  │ ──────────────────────────────────────▶ │        │
└────────┘            85.3 ℃                        │ Broker │  (Mosquitto/EMQX/
                                                      │        │   云厂商 IoT 平台)
┌────────┐  subscribe "factory/+/furnace-+/temp"    │        │
│ 监控页  │ ──────────────────────────────────────▶ │        │
└────────┘◀────── Broker 把匹配的消息推给订阅者 ──────┘
└────────┘
  • Broker:消息中枢(不生产数据,只路由)。独立进程,有成熟开源(Mosquitto、EMQX)和云服务
  • Topic:层级化的主题路径,用 / 分隔——它是 MQTT 的"地址系统"
  • 发布/订阅:发布者和订阅者互不知道对方存在,只认识 Topic

2.2 Topic 设计规范(工业项目第一步)

推荐结构:{站点}/{产线}/{设备}/{通道}

factory/line1/furnace-1/temperature     ✅ 温度数据
factory/line1/furnace-1/status          ✅ 设备状态
factory/line1/furnace-1/alarm           ✅ 告警事件
factory/line2/press-3/pressure

设计纪律:
1. 层级从粗到细(站点→产线→设备→通道),匹配通配符的威力
2. 全小写 + 连字符(别混下划线/驼峰)
3. 不要把数据放 Topic 里("factory/temp/85.3" ❌——Topic 是地址不是内容)

三、通配符:订阅的威力

// 单层通配符 +:匹配恰好一层
"factory/+/furnace-1/temp"
// 匹配 factory/line1/furnace-1/temp、factory/line2/furnace-1/temp
// 不匹配 factory/line1/furnace-2/temp

// 多层通配符 #:匹配任意层级(必须在末尾)
"factory/line1/#"
// 匹配 factory/line1 下所有消息(任何设备任何通道)

// 组合:大屏只关心温度告警
"factory/+/+/alarm"

一次订阅,挂全厂设备——这就是 Topic 层级设计的回报。反过来,设计糟糕的 Topic(furnace1_temptemp_furnace1 混用)会让通配符永远使不上劲。


四、QoS 0/1/2:三档送达保证

MQTT 最核心的机制,也是它征服工业现场的本钱:

QoS 名称 机制 代价 场景
0 至多一次 发出去就不管 可能丢 高频传感器流(丢一帧无所谓)
1 至少一次 PUBACK 确认,没收到就重发 可能重复 告警事件、状态变更(重复可去重)
2 恰好一次 四步握手(PUBREC/PUBREL/PUBCOMP) 最慢最重 计费、开环控制指令

4.1 QoS 1 的"至少"意味着什么

发送方                          接收方
  │ ── PUBLISH(msg, id=42) ──▶  │   收到,处理
  │                              │ ── PUBACK(42) ──▶  │   ACK 在路上丢了!
  │ ✗ 没收到 PUBACK(42)          │
  │ ── PUBLISH(msg, id=42) 重发 ▶ │   又收到 id=42
                                  │   → 重复投递(业务层要按 id 去重)

🎯 QoS 1 把"可靠"拆给了协议和业务:协议保证不丢,业务负责去重(通常按消息 id 或时间戳)。这正好和你 Day 61 要写的补传去重逻辑是同一个思想——先在 MQTT 里见过它,再自己写一遍,理解就立体了

4.2 工业选型速查

温度振动流(1kHz)        → QoS 0(重传一帧旧数据反而是污染)
设备启停/告警            → QoS 1(宁可重复告警,不能漏)
远程设定值下发/配方下发   → QoS 2(执行两次 = 事故)

五、retain 与遗嘱消息

5.1 retain:新订阅者立刻拿到最后一条

痛点:监控页 10:05 才打开,设备 10:00 发的"当前温度"已经过去了——
      新页面要傻等下一次推送(1s~更久)。

retain 消息:Broker 为每个 Topic 存最后一条 retain 消息,
新订阅者连上的瞬间立刻收到它。
// 发布端:标记 retain
client.publish("factory/line1/furnace-1/temp", "85.3", { retain: true });

// 任何新订阅者 connect 后立刻收到这条"温度 85.3",不用等下一次变化

工业语义:状态类数据(当前温度/当前模式)发布时永远带 retain,事件类数据(告警触发)不带。

5.2 遗嘱消息(LWT):设备掉线的自动广播

痛点:设备断电了,谁去告诉 Broker 和所有监控页"我下线了"?
      设备自己肯定来不及说。

LWT:设备连接时预先登记一条"遗嘱"——
     Broker 检测到设备异常掉线(keepalive 超时),替它发布这条遗嘱。
// 设备端:连接时登记遗嘱
const client = mqtt.connect("wss://broker.emqx.io:8084/mqtt", {
  clientId: "furnace-1",
  will: {                                // ⭐ 遗嘱
    topic: "factory/line1/furnace-1/status",
    payload: JSON.stringify({ online: false, ts: Date.now() }),
    qos: 1,
    retain: true,
  },
});
// 设备正常在线时自己发布 { online: true }(retain)
// 设备崩溃/断电 → Broker 代发遗嘱 { online: false } → 所有监控页立刻感知

💡 LWT 就是设备侧的心跳方案——和 Day 58 你给浏览器写的心跳异曲同工,只是它把判死逻辑下沉到了 Broker(keepalive 机制内建)。两个方向的"死亡判决"都见过了,容错设计的感觉应该开始成形


六、浏览器接入:mqtt.js 实战

6.1 安装与连接

npm install mqtt
import mqtt from "mqtt";

/**
 * 浏览器接入 MQTT(MQTT over WebSocket——浏览器不能裸连 TCP 1883)
 */
const client = mqtt.connect("wss://broker.emqx.io:8084/mqtt", {
  clientId: `panel-${Math.random().toString(16).slice(2, 8)}`,
  keepalive: 30,               // 心跳间隔(秒)——Broker 侧判死依据
  clean: false,                // false 保留会话:断线期间 QoS1 消息由 Broker 暂存
  reconnectPeriod: 2000,       // mqtt.js 内建重连(体会:库帮你做了 Day 58 的事)
});

client.on("connect", () => {
  console.log("已连接 Broker");
  // 通配符订阅:全厂所有设备的温度和告警
  client.subscribe(
    ["factory/+/+/temperature", "factory/+/+/alarm"],
    { qos: 1 },
    (err) => { if (!err) console.log("订阅成功"); }
  );
});

client.on("message", (topic: string, payload: Buffer) => {
  // 按 Topic 分发(MqttDataSource 的雏形)
  if (topic.endsWith("/temperature")) {
    const value = Number(payload.toString());
    console.log(`温度 ${topic} = ${value}`);
  } else if (topic.endsWith("/alarm")) {
    const alarm = JSON.parse(payload.toString());
    console.log("告警", alarm);
  }
});

client.on("offline", () => console.log("连接断开,自动重连中…"));

6.2 用 MQTTX 当"模拟设备"

不用写设备代码就能测订阅——图形化客户端 MQTTX:

1. MQTTX 连同一个 Broker(wss://broker.emqx.io:8084/mqtt)
2. 发布:Topic = factory/line1/furnace-1/temperature,内容 = 85.3,QoS 1,勾选 retain
3. 回浏览器:onmessage 触发 ✅
4. 刷新浏览器页面:retain 消息立刻送达 ✅

6.3 MqttDataSource:兑现 DataSource 接口

import type { DataSource, PanelSnapshot } from "./types";
import mqtt from "mqtt";

/**
 * MQTT 数据源:实现与 MockDataSource 相同的 DataSource 契约
 * 面板代码零改动即可切换 —— 第 8 周架构承诺的又一次兑现
 */
export class MqttDataSource implements DataSource {
  private client = mqtt.connect("wss://broker.emqx.io:8084/mqtt", {
    keepalive: 30,
    reconnectPeriod: 2000,
  });
  private listeners: ((s: PanelSnapshot) => void)[] = [];

  constructor() {
    this.client.on("connect", () => {
      this.client.subscribe("factory/+/+/+", { qos: 1 });
    });
    this.client.on("message", (_topic, payload) => {
      // 生产实现:按 topic 聚合各设备数据,组装 PanelSnapshot 后分发
      // (聚合逻辑与 BOSS 战的 WsDataSource 同构,留到 Day 63 统一实现)
      this.listeners.forEach((l) => l(JSON.parse(payload.toString())));
    });
  }

  subscribe(listener: (s: PanelSnapshot) => void): () => void {
    this.listeners.push(listener);
    return () => { this.listeners = this.listeners.filter((l) => l !== listener); };
  }
  start(): void { /* mqtt.js 自动连接,无需额外启动 */ }
  stop(): void { this.client.end(true); }
}

七、MQTT vs 自建 WS 推送:选型决策

维度 自建 WS(Node + ws) MQTT(现成 Broker)
服务端代码 全自己写(订阅表/重连/补传) 几乎为零(Broker 全包)
送达保证 自己实现(Day 61) QoS 1/2 内建
离线消息 自己实现 clean:false + 会话保持
设备掉线感知 自己实现 LWT 遗嘱内建
二进制高频帧 完全自由(Day 59) 支持(payload 是 Buffer)
浏览器接入 原生 WebSocket mqtt.js over WS
运维成本 自己维护服务进程 部署/购买 Broker

结论

  • 有设备侧接入需求的工业项目 → MQTT(你写的只是消费端)
  • 纯 Web 产品、数据源本来就是自家服务 → 自建 WS 更直接
  • BOSS 战两条都做:WS 版主链路(练全栈理解)+ MQTT 版替换实验(验证架构弹性)

八、常见坑点

坑 1:浏览器连 mqtt://(TCP)失败

浏览器只能走 ws:// / wss:// 的 MQTT。连公共 Broker 用 wss://broker.emqx.io:8084/mqtt(注意结尾 /mqtt 路径——EMQX 的 WS 挂载点)。

坑 2:clientId 重复互踢

两个页面用同一个 clientId 连接 → Broker 按"同名会话互踢"处理,来回掉线死循环。修复:clientId 加随机后缀(见 6.1)。

坑 3:retain 消息"删不掉"

发了 retain 消息后想清除:发布一条空 payload 的 retain 消息到同一 Topic(空 retain = 删除指令)。

坑 4:QoS 2 当万金油

高频数据也用 QoS 2 → 四步握手把吞吐拖垮。QoS 按消息语义选,不按"我要可靠"选(第四节速查表)。

坑 5:Topic 大小写不匹配

Factory/line1 订阅 factory/line1 收不到——Topic 区分大小写。规范先行:全小写。

坑 6:clean: false 但 brokerId 每次变

想用会话保持收离线消息,clientId 却每次随机 → Broker 认为是新客户端,离线消息作废。要离线消息:固定 clientId + clean: false 成对出现


九、自测挑战

T1 · 通配符实验(20 分钟)

用 MQTTX 发消息到 5 个不同 Topic,浏览器端分别用 +#、精确 Topic 三种方式订阅,记录每种收到哪些。今天核心作业

T2 · QoS 差异体验(25 分钟)

MQTTX 分别用 QoS 0 和 QoS 1 发 10 条消息到同一 Topic,浏览器订阅 QoS 1。在 DevTools Network→WS 里观察 PUBACK 帧的差异。

T3 · 遗嘱实验(30 分钟)

写一个"模拟设备"脚本(Node + mqtt):连上后定期发布温度,10 秒后 process.exit(1) 模拟崩溃。浏览器订阅其 status Topic,验证 LWT 自动到达。

T4 · MqttDataSource 接面板(进阶,40 分钟)

把 6.3 的 MqttDataSource 补全(聚合各设备 Topic 组装 PanelSnapshot),替换上周面板的 MockDataSource 跑起来——用 MQTTX 手动发数据驱动面板。这是 BOSS 战 MQTT 版的预演


十、总结

概念 工业价值
发布/订阅 + Topic 设备与应用彻底解耦,通道层可扩展
通配符 + / # 一次订阅挂全厂
QoS 0/1/2 按消息语义选送达保证(流 0 / 事件 1 / 指令 2)
retain 新订阅者秒拿最新状态
LWT 遗嘱 设备掉线自动广播(Broker 侧心跳)
mqtt.js over WS 浏览器接入工业数据的正解

明天回到自建 WS 链路,把 Day 58 留下的最后一块硬骨头啃掉:断线期间丢的数据怎么补回来——消息缓冲、时间戳去重、无缝衔接实时流。MQTT 的 QoS 1 给你打的预防针,明天全部派上用场。

评论