拓十年匠心定制 · 商业建站与技术教学双线并行 咨询热线:400-886-1026 service@lmnt.cn
ARTICLE DETAIL

资讯详情

深耕网站建设与运营推广的一线实战洞察。

工业上位机把PLC数据推上MQTT_三_离线队列与QoS1可靠投递

工业上位机把PLC数据推上MQTT_三_离线队列与QoS1可靠投递

工业上位机把 PLC 数据推上 MQTT(三):离线队列 + QoS1,断网不丢、连上补发

系列说明:这是《工业上位机把 PLC 数据推上 MQTT》三部曲第 3 篇,也是收尾篇。
第 1 篇讲了整体架构和发布链路;第 2 篇讲了死区 + 最小间隔(流量治理);本篇讲"断网了数据怎么办"。
前两篇:(一)架构与完整链路 · (二)死区与节流


前两篇把"发什么、发多频繁"定了。最后这篇讲最要命的:网络抖一下,那段时间的数据不能丢。

车间网络什么德行干过现场的都懂——交换机重启、光纤被挖断、4G 基站抽风,短则几秒长则几十分钟。要是没有兜底,这段时间 PLC 的数据就真空了,MES 那边对账对不平,工艺组第一个找你。

我们用了两层保险:QoS1 保"发出去的被确认",离线队列保"没发出去的先存着"。

一、QoS1:至少一次,靠确认

MQTT 的 QoS 三级里,设备上报一般用 QoS1(至少一次)。QoS0 是"发了就不管",工厂数据不敢用;QoS2"恰好一次"握手太重,工业场景性价比低。QoS1 的语义是:Broker 收到必须回 PUBACK,没收到发布端就重发。

关键在"重发"怎么跟踪。项目里有个容易混的三个概念,我在qos1state.h里专门分开了:

// qos1state.h// msgId :应用层持久标识(单调递增),用于日志/去重/对账// token :Paho MQTTAsync_token(int),发送时库返回,用于投递完成回调关联// packetId :线上 Packet Identifier(Paho 内部,不暴露)

混了到底会怎样(真踩过的坑):早期版本我图省事,把 Paho 回调里拿到的token直接当应用层msgId去查m_byToken表。结果 Paho 的token是"发送批次"维度的,同一批次多条消息共用一个 token,而m_byToken是按单条消息建的——onAck(token)命中后只清了表里一条,剩下几条永远停在Inflight。表象就是状态面板inflight只增不减、重连后这些"幽灵消息"被全量补发翻倍,消费端 ts 去重都救不回来(因为 ts 虽同但根本没发成功过,被当成新消息又发一遍)。教训:token 只用来关联"这次投递完成回调",msgId 才是业务去重/对账的主键,两者必须分开存。

发出去时建记录,收到 ACK 时清记录:

发出去时建记录,收到 ACK 时清记录:

voidQos1Tracker::markInflight(qint64 msgId,MQTTAsync_token token){Qos1Record r;r.msgId=msgId;r.token=token;r.status=Qos1Status::Inflight;m_byToken.insert(token,r);// 按 token 关联}voidQos1Tracker::onAck(MQTTAsync_token token,Qos1Status st){autoit=m_byToken.find(token);// 投递完成回调按 token 命中if(it!=m_byToken.end()){it->status=st;m_byToken.erase(it);}}

m_qos1.inflightCount()直接喂给状态面板,运维一眼能看到"还有几条没被 Broker 确认"。

一个必须说清楚的设计点:断线重连后,之前 QoS1 没被确认的消息会被全量补发,这必然产生重复。这是 QoS1 的本职,不是 bug。解决办法在消费端——按(tagKey, ts)做幂等去重。ts是发布时就带上的采样时刻,重复的消息ts一样,落库时去重即可。所以我们的 change 载荷里永远带着ts,就是这个用处:

{"tag":"DefaultPLC/Main_Plc/回水温度","value":52.3,"quality":"GOOD","ts":"2026-09-28T09:12:00.123Z"}

二、离线队列:连不上就先存盘

QoS1 只管"发出去→确认"。但要是压根没连上 Broker,sendRaw直接返回 false,消息得有个地方先待着。这就是enqueue→ 离线队列。

队列分两级,内存不够再落盘:

// mqttpublisher.cpp::enqueueconstintcap=m_cfg?m_cfg->maxQueueMem:2000;// 内存上限 2000 条if(m_memQueue.size()<cap){m_memQueue.append(m);}elseif(m_dbOk){diskInsert(channel,topic,payload,qos);// 超限转 SQLite 落盘}else{m_memQueue.append(m);// 无磁盘兜底,best-effort}

落盘用的是 SQLite,独立连接,一张outbound表:

CREATETABLEIFNOTEXISTSoutbound(idINTEGERPRIMARYKEYAUTOINCREMENT,channelTEXT,topicTEXT,payloadBLOB,qosINT,created_atINTEGER);

两个细节是现场能救命的:

① 有界,不撑爆磁盘。落盘也有限额(默认 2MB),超了删最旧一行:

if(m_cfg&&(m_diskBytes+payload.size()>m_cfg->maxQueueDiskBytes)){del.exec("DELETE FROM outbound WHERE id = (SELECT MIN(id) FROM outbound)");}

网络中断半小时,队列把"最早的历史"让给"最新的实时",这符合监控语义——宁可丢最老的,也要保最新的。

② 并发保护。离线队列的写(落盘)和读(冲刷)可能跨线程抢 SQLite,我们设了QSQLITE_BUSY_TIMEOUT=5000,避免直接SQLITE_BUSY报错丢数据:

m_db.setConnectOptions("QSQLITE_BUSY_TIMEOUT=5000");// 关键:避免 SQLITE_BUSY

三、连上即补发:flushQueue

重连成功后第一件事,是把攒着的消息吐出去:

// mqttpublisher.cpp::handleConnected → flushQueuevoidMqttPublisher::flushQueue(){while(!m_memQueue.isEmpty()){// 先内存队列(先进先出)QueuedMessage m=m_memQueue.takeFirst();if(!sendRaw(m.topic,m.payload,m.qos)){m_memQueue.prepend(m);break;}m_totalPublished++;}if(m_dbOk){// 再磁盘,按 id ASC 顺序for(auto&m:diskLoad()){if(!sendRaw(m.topic,m.payload,m.qos))break;diskDelete(m.id);// 发成功才删,失败留着下轮}}}

关于顺序(先内存后磁盘会不会乱序?):入队逻辑是"内存没满就append,满了(size()>=cap)才落盘",所以内存里是更早的消息、磁盘里是更晚的消息;冲刷时内存 FIFO(内升序)+ 磁盘id ASC(内升序),单次连续断网、一次冲刷到底时,消费方收到的是全局时间升序,没问题。唯一的边角场景:上轮冲刷到一半又断、且新消息把内存填满溢出到磁盘,磁盘里残留的更早消息会排在内存新消息之后,出现局部倒挂。生产硬化做法见下——按createdAt合并排序再发,所有场景都稳。

注意"发成功才删":冲刷中途又断了,sendRaw失败就break,没发的留着,下一轮连上接着补。这条链路接在第 1 篇的publishFormatted末尾——连不上走enqueue,连上了走flushQueue,闭环。

3.1 生产硬化①:按创建时间合并排序(彻底杜绝倒挂)

把内存和磁盘的消息一起按createdAt升序排好再发。需要diskLoad()顺带返回created_at列(表里本来就有):

// 合并内存 + 磁盘 → 按 createdAt 升序 → 逐个发structItem{qint64 seq;QueuedMessage m;boolfromDisk;};QVector<Item>all;for(constauto&m:m_memQueue)all.append({m.createdAt,m,false});if(m_dbOk)for(constauto&m:diskLoad())all.append({m.createdAt,m,true});std::sort(all.begin(),all.end(),[](constItem&a,constItem&b){returna.seq<b.seq;});for(constauto&it:all){if(!sendRaw(it.m.topic,it.m.payload,it.m.qos))break;if(it.fromDisk){diskDelete(it.m.id);m_diskBytes-=it.m.payload.size();}m_totalPublished++;}m_memQueue.clear();

3.2 生产硬化②:限速分批冲刷(防消息风暴)

MQTTAsync_send是异步非阻塞,紧循环一次能把内存 2000 条 + 磁盘数万条在毫秒级全甩出去,打满 Broker 的max_inflight_messages、或占满 4G 带宽。用定时器分批:

// m_flushPerTick 默认 50、间隔 50ms(≈1000 条/秒),现场按 Broker 能力调voidMqttPublisher::drainTick(){intsent=0;while(!m_drain.isEmpty()&&sent<m_flushPerTick){constautom=m_drain.takeFirst();if(!sendRaw(m.topic,m.payload,m.qos)){m_drain.prepend(m);break;}if(m.fromDisk){diskDelete(m.id);m_diskBytes-=m.payload.size();}m_totalPublished++;sent++;}if(!m_drain.isEmpty())m_flushTimer->singleShot(50,this,&MqttPublisher::drainTick);}

落盘队列也别一次性SELECT全表进内存(默认 2MB 有界,但数万行仍占一块),diskLoad改成LIMIT分批游标、边取边发,和上面的限速冲刷合在一起最稳。

3.3 生产硬化③:过期丢弃(TTL)

断网 2 小时,重连后把 2 小时前的温度补发给实时看板,消费方可能误判"当前值"。配置加maxAgeMs(0=不过期),冲刷/落盘前丢超期消息:

// 落盘/冲刷前判断:now - createdAt > maxAgeMs 则丢弃if(m_cfg->maxAgeMs>0&&(now-m.createdAt>m_cfg->maxAgeMs)){if(m.fromDisk)diskDelete(m.id);// 内存的直接跳过continue;}

更彻底的办法是升级MQTT 5.0 的Message Expiry Interval,让 Broker 自动丢弃超期补发(见第七节)。

四、一个真踩过的线程坑

Paho 的 C 回调(onConnectionLost/onDeliveryComplete等)是在库自己的网络线程里触发的。我们一开始在回调里直接写 SQLite、等 ACK、加锁——结果偶发死锁,UI 卡死。

根因:回调线程不能碰主线程的资源(SQLite 连接、QMutex 持有的业务状态)。方案是回调里只 marshal 回主线程,绝不阻塞:

// mqttpublisher.cpp::onDeliveryComplete(网络线程触发)voidMqttPublisher::onDeliveryComplete(void*context,MQTTAsync_token token){auto*self=static_cast<MqttPublisher*>(context);QMetaObject::invokeMethod(self,[self,token](){self->m_qos1.onAck(token,Qos1Status::Acked);// 回主线程再改状态},Qt::QueuedConnection);}

Qt::QueuedConnection把活儿排队到主线程事件循环,网络线程立刻返回。业务状态(inflight 计数、离线队列读写)永远只在主线程动,死锁消失。这条规矩写在头文件注释里:回调禁止阻塞(写 SQLite/等 ACK/加锁都会死锁),统一marshal回主线程。

五、报警通道:绕过节流、走 QoS1、靠 ts 去重

第 1 篇提过报警是独立通道,这里把它的"特殊待遇"说清。publishAlarm()直接formatAlarm → publishFormatted,根本不调shouldPublish——死区和最小间隔都拦不到它。这恰恰是对的:报警事件绝不能因为"值没变够多"或"离上次太近"被吞掉,该报就报。

报警 QoS 走m_cfg->qos(默认 1)。有人问"报警要不要 QoS2 防重复?“我的建议是不用:报警最大的风险是"漏报"而不是"重复报”,QoS1 + 消费端按(tagKey, ts)幂等去重已经够用;QoS2 握手更重,而且重复问题同样靠 ts 去重解决,上 QoS2 是亏本买卖。

六、Broker 端配套配置清单(上位机再稳也得 Broker 配合)

断线期间 Broker 没配好,消息照样丢。Mosquitto 最低配置参考:

# mosquitto.conf max_queued_messages 0 # 0=不限制 Broker 侧队列(或按内存设大,配合上位机限速) max_inflight_messages 100 # 单客户端在途上限,避免一个客户端占满 persistence true # 开启持久化,Broker 重启不丢 retained / 会话 # 若上位机 cleanSession=false(我们默认就是 false),务必开持久会话, # 否则断线期间订阅关系与会话丢失,重连后收不到补发

EMQX 对应:mqtt.max_inflight、开启retainer、配置session_expiry。给甲方交付时这份清单直接附上,别光说自己客户端可靠。

七、可观测性与运维监控

status()已经把关键指标吐出来了,运维面板至少该挂这几项:

  • 离线队列深度=queuedMem+queuedDiskBytes,超阈值(比如内存 > 80% 上限 / 磁盘 > 1MB)就告警——意味着网络长期不稳或 Broker 收不动;
  • inflight长时间不归零= Broker 不回 PUBACK(卡死 / 网络半通),告警;
  • 发布失败率= 失败计数 /totalPublished,暴露;
  • lastError进日志,断连原因可追溯(是onConnectionLost报的 cause,还是MQTTAsync_send失败)。

这几项不展示,断网了你都不知道数据在丢。

八、MQTT 5.0 与版本差异

全文基于 PahoMQTTAsync(MQTT 3.1.1)。工业场景直接能用上 5.0 的三个特性:

  • Message Expiry Interval:消息级过期,Broker 自动丢弃超期补发,正好解决第三节 3.3 的 TTL 问题,比应用层maxAgeMs更彻底;
  • Session Expiry Interval:替代cleanSession那个别扭的布尔,精细控制会话保留时长;
  • Shared Subscription:多消费者负载均衡,适合一个 topic 被多个后端抢着处理的场景。

升级路径建议:先用 3.3 的应用层maxAgeMs兜底,等现场 Broker 支持 5.0 再切原生过期,平滑过渡。

九、三篇串起来

回到第 1 篇那张链路图,现在三道闸都齐了:

采集 → onTagValueChanged → shouldPublish (第2篇:死区+最小间隔,砍掉没用的) → formatChange (JSON 格式化、带 ts) → publishFormatted ├─ sendRaw 成功 → m_totalPublished++ └─ sendRaw 失败 → enqueue(本篇:内存→落盘,有界) 重连成功 → flushQueue(本篇:连上即补发,发成功才删) QoS1 → markInflight / onAck(本篇:至少一次 + 幂等去重)
  • 第 1 篇管架构:三通道、主题、JSON、异步客户端;
  • 第 2 篇管流量:死区 + 最小间隔,解决"发太多";
  • 第 3 篇管可靠:QoS1 + 离线队列,解决"断网丢"。

现场配的时候我的习惯:先开 change + 默认死区(0.5%)看流量落不落得下来;要历史全貌再开 snapshot;最后确认离线队列上限按现场断网时长估(2MB 大概够撑一阵,真长断网调大maxQueueDiskBytes)。QoS 默认 1 别动,消费端记得按(tagKey, ts)去重。


本篇配置速查卡(可靠投递)

配置默认说明
qos1默认 QoS1(至少一次),消费端按(tagKey, ts)去重
cleanSessionfalse持久会话,断线重连可续订(Broker 须开持久化)
keepAliveSec/connectTimeoutSec60/30心跳 / 连接超时
reconnectBackoffSec5断线重连退避
maxQueueMem2000内存队列上限(条)
maxQueueDiskBytes2MB落盘上限,超了丢最旧保最新
maxAgeMs0(硬化新增)补发消息最大龄期,0=不过期;升级 MQTT5 用Message Expiry更彻底
m_flushPerTick/ 间隔50/50ms(硬化新增)限速分批冲刷,≈1000 条/秒,按 Broker 调

完整性清单(Broker 配置 / 可观测性 / MQTT 5.0)见本篇第六~八节。

本文及 PLCMonitor 系列文章均为免费分享。
本文免费分享,如需转载,请联系作者获取授权。
如果觉得这篇文章对你有帮助,欢迎点赞收藏。
源码获取地址:https://github.com/freddiezhang1990/plcmonitor

返回列表