Day 60 · MQTT 与工业物联网协议 — 发布/订阅的世界观,QoS 的保证书
这周你的服务器一直在"点对点"推数据——它必须知道每个客户端要什么、连没连、断了怎么补。而工业现场的世界观完全不同:设备只管往 Topic 上发布,谁在订阅、有多少人订阅,设备根本不知道也不关心。这就是 MQTT 的发布/订阅模型,工业物联网的事实标准(PLC、网关、阿里云 IoT、AWS IoT 全都原生支持)。今天理解它的模型与 QoS 机制,并用 mqtt.js 把浏览器接上真 Broker——BOSS 战的数据链路从此有了"工业正统性"。
目录
- 一、为什么工厂设备侧都用 MQTT
- 二、核心概念:Broker / Topic / 发布订阅
- 三、通配符:订阅的威力
- 四、QoS 0/1/2:三档送达保证
- 五、retain 与遗嘱消息
- 六、浏览器接入:mqtt.js 实战
- 七、MQTT vs 自建 WS 推送:选型决策
- 八、常见坑点
- 九、自测挑战
- 十、总结
一、为什么工厂设备侧都用 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_temp、temp_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 给你打的预防针,明天全部派上用场。