不少工厂的"上云"是这样做的:云端每 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 为主题保留最后一条,新订阅立即收到设备上线/下线状态、参数快照,别用于高频遥测
CleanSessionfalse 时 Broker 保留订阅与离线 QoS1 消息配合固定 ClientId,但保存窗口和条数有上限
KeepAlive心跳周期,1.5 倍周期无报文 Broker 判定离线30~60 秒;NAT 超时短的网络调到 20 秒
LWT 遗嘱连接时登记,异常断线时由 Broker 代发设备非正常掉线自动广播 offline

三、QoS 怎么选:不是越高越好

数据类型QoSRetain理由
高频过程量(温度/转速秒级)0丢一两帧无所谓,最新值才有意义,追求低开销
班产量、工单计数、条码过站1一条都不能丢,允许重复,云端按消息ID去重
报警/停机事件1必须送达且可审计,事件ID幂等
在线状态1LWT + 上线消息都 Retain,订阅方一打开就看到当前状态
下行参数/指令1上位机必须回执行结果到独立响应主题,不能只靠 QoS
QoS1 的"至少一次"意味着云端会收到重复消息。每条消息带唯一 ID(设备号+本地自增序号),云端入库唯一约束去重,这是断网补传时消息重发的必然要求——不能假设 QoS1 只来一次。

四、主题分层:上线第一天就定好,后面改不动

主题是设备和云端之间的契约,一旦有数据沉淀和报表依赖,改主题等于动接口。建议四层定位设备、一层区分消息类型:

ying/v1/f/SH01/l/A01/d/PLC-260914-03/telemetry // 上行·周期遥测(产量/温度/节拍) ┆ ┆ ┆ └event // 上行·报警/停机/换型等事件 ┆ ┆ └status // 上行·online/offline(Retain,LWT) ┆ └cmd // 下行·云端指令(参数下发/远程启停)cmdresp/{msgId} // 上行·指令执行结果,按消息ID关联 // 云端一条订阅收全厂:ying/v1/f/+/l/+/d/+/telemetry // 单台设备调试只订: ying/v1/f/SH01/l/A01/d/PLC-260914-03/#
  • 开头带 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才算成功
}
重连不要写 while(true) 立即重试。断网时疯狂重连会让工控机网卡队列堆满、Broker 日志爆掉。指数退避 1→2→4…封顶 60 秒,再加 ±20% 随机抖动,避免几百台设备同一秒集体重连造成惊群。

六、断网补传:持久会话靠不住,必须有应用层 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
恢复按id顺序补发
每批50限速
云端按msg_id去重
按collect_ts对齐
补传要有兜底策略。队列超过阈值(如断网超过本地可存时长)要现场声光+界面报警,不能让几 GB 历史数据在恢复瞬间把产线带宽和云端打满;更不能为了清队列把实时数据通道堵住——必要时给补传单独限速、实时遥测优先。

七、消息体约定:让数据十年后还能读

{
  "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 没有意义;质量码照实标,坏值比缺值更有诊断价值。

八、上线前必查的几个安全与运维点

  1. 生产环境一律 8883 + TLS + 一机一密,1883 明文只在隔离内网调试用;凭证随装机流程预烧录,不写死在代码里;
  2. ClientId 全局唯一:开发机、测试机别复用设备 ClientId,否则两边互相踢线,现场表现为"每隔几十秒掉线重连",极难排查;
  3. Broker 侧按主题做 ACL:设备证书只能发布/订阅自己 SN 路径下的主题,一台设备失陷不能伪造全厂数据;
  4. 监控四个指标:在线率、消息丢失/重传率、outbox 积压深度、指令响应超时率——出问题时这四个数能直接定位是网络、设备还是云端;
  5. 先小批灰度:拿一条产线跑一周,人为拔网线、关 Broker、重启设备各演练几次,确认补传不重不丢、状态正确翻转,再全厂推广。

总结一下这套架构:MQTT 负责实时通道,QoS1+消息ID负责可靠送达,LWT 负责在线感知,SQLite outbox 负责断网不丢,云端幂等负责补传不重。把这五件事做齐,产线数据上云才算真正可靠——剩下的报表、看板、大屏,都是在可靠数据之上的应用层工作。