不少工厂的"上云"是这样做的:云端每 30 秒连一次工厂数据库,捞增量数据。结果现场一断外网,云端轮询全部超时报错;数据什么时候产生的、哪台设备发的,全靠猜;云端想下发个参数,还得反过来开端口。MQTT 就是为这种弱网、海量设备、双向通信场景设计的长连接协议。这篇讲清楚怎么用 C#(MQTTnet)把产线数据可靠地送上去:协议选型、QoS、主题设计、遗嘱保活,以及最容易被忽略的断网补传。
一、为什么是 MQTT:和 HTTP 轮询的账要算清
| 对比项 | HTTP 定时轮询 | MQTT 长连接 |
|---|---|---|
| 通信方向 | 只能云端主动拉;下行要另开端口/反向轮询 | 发布/订阅天然双向,云端下发指令同样实时 |
| 实时性 | 受轮询间隔限制,间隔短了压垮接口 | 变化即推,毫秒级到达 |
| 弱网表现 | 断网期间请求全部失败,要自己实现重试队列 | 连接保活+重连,QoS 与持久会话覆盖部分离线场景 |
| 连接开销 | 每次请求 TCP+TLS 握手(短连接时) | 一条 TCP 长连接复用,报文头最小 2 字节 |
| 在线感知 | 无法区分"没数据"和"设备死了" | 遗嘱消息(LWT)+ Retain 上线/下线秒级感知 |
| 对接成本 | 写接口简单,运维熟悉 | 需要部署/租用 Broker(EMQX 等) |
▲ 结论:报表类、低频查询继续走 HTTP/数据库;设备遥测、状态感知、远程下发这类高频双向链路用 MQTT。两者不是替代关系。
二、核心概念先对齐,代码才看得懂
| 概念 | 含义 | 现场用法 |
|---|---|---|
| Broker | 消息代理服务器,所有端只和它通信,端与端不直连 | 自建 EMQX 或云物联网平台,1883 明文 / 8883 TLS |
| ClientId | 连接唯一标识 | 一机一码,重复 ClientId 会被 Broker 互踢(重要坑) |
| Topic 主题 | 斜杠分层的消息地址,支持 +/# 通配符订阅 | 按 工厂/产线/设备/消息类型 分层(见第四节) |
| QoS 0 | 至多一次,发完不管 | 高频、可丢弃的过程值(温度秒级采样) |
| QoS 1 | 至少一次,必须 PUBACK 确认,可能重复 | 产量、报警、状态等关键消息,消费端做幂等 |
| QoS 2 | 恰好一次,四次握手 | 几乎不用:开销大,端侧/Broker 支持参差,幂等设计后 QoS1 足够 |
| Retain 保留消息 | Broker 为主题保留最后一条,新订阅立即收到 | 设备上线/下线状态、参数快照,别用于高频遥测 |
| CleanSession | false 时 Broker 保留订阅与离线 QoS1 消息 | 配合固定 ClientId,但保存窗口和条数有上限 |
| KeepAlive | 心跳周期,1.5 倍周期无报文 Broker 判定离线 | 30~60 秒;NAT 超时短的网络调到 20 秒 |
| LWT 遗嘱 | 连接时登记,异常断线时由 Broker 代发 | 设备非正常掉线自动广播 offline |
三、QoS 怎么选:不是越高越好
| 数据类型 | QoS | Retain | 理由 |
|---|---|---|---|
| 高频过程量(温度/转速秒级) | 0 | 否 | 丢一两帧无所谓,最新值才有意义,追求低开销 |
| 班产量、工单计数、条码过站 | 1 | 否 | 一条都不能丢,允许重复,云端按消息ID去重 |
| 报警/停机事件 | 1 | 否 | 必须送达且可审计,事件ID幂等 |
| 在线状态 | 1 | 是 | LWT + 上线消息都 Retain,订阅方一打开就看到当前状态 |
| 下行参数/指令 | 1 | 否 | 上位机必须回执行结果到独立响应主题,不能只靠 QoS |
四、主题分层:上线第一天就定好,后面改不动
主题是设备和云端之间的契约,一旦有数据沉淀和报表依赖,改主题等于动接口。建议四层定位设备、一层区分消息类型:
- 开头带 v1 版本号,协议大改时开 v2 并行,不动老设备;
- 层级只放稳定的拓扑属性(工厂/产线/设备序列号),产品型号、工单等会变的东西放消息体 JSON,别放主题——否则订阅规则和权限表会爆炸;
- 遥测一个主题打天下(JSON 里用字段区分指标),不要每个测点一个主题,几百个主题订阅和授权都难维护;
- 下行指令和执行结果分两个主题,响应里带 msgId,云端超时可重发同一条指令而不会让设备执行两次(设备端按 msgId 幂等)。
五、MQTTnet 客户端:连接、保活、遗嘱、退避重连
下面代码基于 MQTTnet 4.3(NuGet 安装即可)。一个能长期跑在工控机上的客户端要具备:固定 ClientId 的持久会话、LWT、自动重连退避、发布确认检查。
public class MqttUploader : IDisposable { private readonly IMqttClient _client; private readonly MqttClientOptions _options; private int _backoffSec = 1; public MqttUploader(string broker, string deviceSn) { _client = new MqttFactory().CreateMqttClient(); _client.DisconnectedAsync += OnDisconnected; _client.ApplicationMessageReceivedAsync += OnCommand; // 下行指令 var statusTopic = $"ying/v1/f/SH01/l/A01/d/{deviceSn}/status"; _options = new MqttClientOptionsBuilder() .WithTcpServer(broker, 8883) .WithTlsOptions(o => o.UseTls = true) .WithCredentials(deviceSn, "一机一密-预烧录") .WithClientId(deviceSn) // 固定且唯一,重启后复用同一会话 .WithCleanSession(false) // Broker保留订阅与离线QoS1消息 .WithKeepAlivePeriod(TimeSpan.FromSeconds(30)) .WithWillTopic(statusTopic) .WithWillPayload("\"offline\""u8.ToArray()) .WithWillQualityOfServiceLevel(MqttQualityOfServiceLevel.AtLeastOnce) .WithWillRetain(true) // 异常掉线由Broker代发offline .Build(); } public async Task StartAsync(CancellationToken ct) { await _client.ConnectAsync(_options, ct); await _client.SubscribeAsync(new MqttTopicFilterBuilder() .WithTopic($"ying/v1/+/+/+/d/{_client.Options.ClientId}/cmd") .WithAtLeastOnceQoS().Build(), ct); // 主动上线(Retain),覆盖掉遗嘱的 offline await PublishAsync(TopicOf("status"), "\"online\""u8.ToArray(), retain: true); _backoffSec = 1; } private async Task OnDisconnected(MqttClientDisconnectedEventArgs e) { Log.Warning($"MQTT断开,{_backoffSec}s 后重连"); await Task.Delay(TimeSpan.FromSeconds(_backoffSec)); try { await _client.ConnectAsync(_options); _backoffSec = 1; // 成功后退避归零 } catch { _backoffSec = Math.Min(_backoffSec * 2, 60); } // 1→2→4…封顶60s } }
发布关键消息时检查 PUBACK 结果,没拿到确认就转入本地补传队列:
public async Task<bool> PublishAsync(string topic, byte[] payload, MqttQualityOfServiceLevel qos = MqttQualityOfServiceLevel.AtLeastOnce, bool retain = false) { if (!_client.IsConnected) return false; // 不阻塞调用方,交给outbox var msg = new MqttApplicationMessageBuilder() .WithTopic(topic).WithPayload(payload) .WithQualityOfServiceLevel(qos).WithRetain(retain).Build(); var r = await _client.PublishAsync(msg, CancellationToken.None); return r.ReasonCode == MqttClientPublishReasonCode.Success; // QoS1收到PUBACK才算成功 }
六、断网补传:持久会话靠不住,必须有应用层 outbox
最常见的误解:"开了 CleanSession=false,断网数据 Broker 帮我存着。"要弄清它的边界:Broker 持久会话只在设备订阅保持期间缓存离线消息,且有队列长度/会话超时上限(EMQX 默认队列满了会丢最旧消息);工厂光缆被挖断半天、换机换 ClientId,缓存一律没有。关键业务必须本地落盘。
做法是经典的发件箱(Outbox)模式:遥测事件产生时,先写本地 SQLite 队列表,再尝试发布;QoS1 拿到 PUBACK 才删记录;连接恢复后后台线程按时间顺序补发:
| 字段 | 作用 |
|---|---|
| id | 本地自增主键,补发严格按它排序 |
| msg_id | 全局唯一消息ID(设备号+自增序号),云端去重的依据 |
| topic / payload / qos | 完整消息快照 |
| collect_ts | 采集时刻(不是发送时刻),云端按它对齐数据,补传不乱序 |
| retry_count / last_error | 重试次数与失败原因,超限告警,防止死循环 |
// 产生消息:先落盘再发送,发送失败/程序崩溃记录都还在 public async Task ReportTelemetryAsync(string topic, byte[] payload, bool critical) { if (critical) _db.Execute( "INSERT INTO mqtt_outbox(msg_id,topic,payload,qos,collect_ts) VALUES(?,?,?,?,?)", NewMsgId(), topic, payload, 1, DateTimeOffset.Now); var ok = await PublishAsync(topic, payload); if (ok && critical) _db.Execute("DELETE FROM mqtt_outbox WHERE msg_id=?", NewMsgId()); // 失败什么都不用做,补发线程会捞;QoS0过程值不落盘,丢就丢 } // 后台补发:连接恢复后每2秒扫一批,先发旧的,限速避免补传风暴 async Task DrainOutboxAsync(CancellationToken ct) { while (!ct.IsCancellationRequested) { if (_client.IsConnected) { var rows = _db.Query<OutboxRow>( "SELECT * FROM mqtt_outbox ORDER BY id LIMIT 50").ToList(); foreach (var r in rows) { var ok = await PublishAsync(r.Topic, r.Payload, (MqttQualityOfServiceLevel)r.Qos); if (ok) _db.Execute("DELETE FROM mqtt_outbox WHERE id=?", r.Id); else { MarkRetry(r); break; } // 一条发不动就停,保序 } } await Task.Delay(2000, ct); } }
PUBACK即删
落SQLite
每批50限速
按collect_ts对齐
七、消息体约定:让数据十年后还能读
{
"msgId": "PLC-260914-03-00084521", // 全局唯一,云端幂等去重
"ts": "2026-09-14T10:23:45+08:00", // 采集时刻,带时区,补传不变
"sn": "PLC-260914-03",
"metrics": {
"output": 1284, // 产量整数
"temp_c": 76.4,
"weight_g": "1250.30" // 称重等decimal用字符串,浮点不丢精度
},
"q": "GOOD" // 质量码 GOOD/BAD/UNCERTAIN,坏值也照发但标明
}
- 字段名固定英文、单位写进字段名(temp_c、weight_g),别让云端猜单位;
- decimal/高精度值用字符串传输,JSON 数字走 double,称重、金额会丢精度;
- 采集端时钟必须准:配合 NTP/SNTP 定时对时,否则 ts 没有意义;质量码照实标,坏值比缺值更有诊断价值。
八、上线前必查的几个安全与运维点
- 生产环境一律 8883 + TLS + 一机一密,1883 明文只在隔离内网调试用;凭证随装机流程预烧录,不写死在代码里;
- ClientId 全局唯一:开发机、测试机别复用设备 ClientId,否则两边互相踢线,现场表现为"每隔几十秒掉线重连",极难排查;
- Broker 侧按主题做 ACL:设备证书只能发布/订阅自己 SN 路径下的主题,一台设备失陷不能伪造全厂数据;
- 监控四个指标:在线率、消息丢失/重传率、outbox 积压深度、指令响应超时率——出问题时这四个数能直接定位是网络、设备还是云端;
- 先小批灰度:拿一条产线跑一周,人为拔网线、关 Broker、重启设备各演练几次,确认补传不重不丢、状态正确翻转,再全厂推广。
总结一下这套架构:MQTT 负责实时通道,QoS1+消息ID负责可靠送达,LWT 负责在线感知,SQLite outbox 负责断网不丢,云端幂等负责补传不重。把这五件事做齐,产线数据上云才算真正可靠——剩下的报表、看板、大屏,都是在可靠数据之上的应用层工作。
