尧图精选

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

🕒 发布时间:2026/10/1 20:01:19 📁 来源:尧图网络
工业上位机把 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_tokenint发送时库返回用于投递完成回调关联// packetId 线上 Packet IdentifierPaho 内部不暴露混了到底会怎样真踩过的坑早期版本我图省事把 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.msgIdmsgId;r.tokentoken;r.statusQos1Status::Inflight;m_byToken.insert(token,r);// 按 token 关联}voidQos1Tracker::onAck(MQTTAsync_token token,Qos1Status st){autoitm_byToken.find(token);// 投递完成回调按 token 命中if(it!m_byToken.end()){it-statusst;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 只管发出去→确认。但要是压根没连上 BrokersendRaw直接返回 false消息得有个地方先待着。这就是enqueue→ 离线队列。队列分两级内存不够再落盘// mqttpublisher.cpp::enqueueconstintcapm_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_diskBytespayload.size()m_cfg-maxQueueDiskBytes)){del.exec(DELETE FROM outbound WHERE id (SELECT MIN(id) FROM outbound));}网络中断半小时队列把最早的历史让给最新的实时这符合监控语义——宁可丢最老的也要保最新的。② 并发保护。离线队列的写落盘和读冲刷可能跨线程抢 SQLite我们设了QSQLITE_BUSY_TIMEOUT5000避免直接SQLITE_BUSY报错丢数据m_db.setConnectOptions(QSQLITE_BUSY_TIMEOUT5000);// 关键避免 SQLITE_BUSY三、连上即补发flushQueue重连成功后第一件事是把攒着的消息吐出去// mqttpublisher.cpp::handleConnected → flushQueuevoidMqttPublisher::flushQueue(){while(!m_memQueue.isEmpty()){// 先内存队列先进先出QueuedMessage mm_memQueue.takeFirst();if(!sendRaw(m.topic,m.payload,m.qos)){m_memQueue.prepend(m);break;}m_totalPublished;}if(m_dbOk){// 再磁盘按 id ASC 顺序for(autom: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;};QVectorItemall;for(constautom:m_memQueue)all.append({m.createdAt,m,false});if(m_dbOk)for(constautom:diskLoad())all.append({m.createdAt,m,true});std::sort(all.begin(),all.end(),[](constItema,constItemb){returna.seqb.seq;});for(constautoit: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(){intsent0;while(!m_drain.isEmpty()sentm_flushPerTick){constautomm_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 小时前的温度补发给实时看板消费方可能误判当前值。配置加maxAgeMs0不过期冲刷/落盘前丢超期消息// 落盘/冲刷前判断now - createdAt maxAgeMs 则丢弃if(m_cfg-maxAgeMs0(now-m.createdAtm_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*selfstatic_castMqttPublisher*(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 / 会话 # 若上位机 cleanSessionfalse我们默认就是 false务必开持久会话 # 否则断线期间订阅关系与会话丢失重连后收不到补发EMQX 对应mqtt.max_inflight、开启retainer、配置session_expiry。给甲方交付时这份清单直接附上别光说自己客户端可靠。七、可观测性与运维监控status()已经把关键指标吐出来了运维面板至少该挂这几项离线队列深度queuedMemqueuedDiskBytes超阈值比如内存 80% 上限 / 磁盘 1MB就告警——意味着网络长期不稳或 Broker 收不动inflight长时间不归零 Broker 不回 PUBACK卡死 / 网络半通告警发布失败率 失败计数 /totalPublished暴露lastError进日志断连原因可追溯是onConnectionLost报的 cause还是MQTTAsync_send失败。这几项不展示断网了你都不知道数据在丢。八、MQTT 5.0 与版本差异全文基于 PahoMQTTAsyncMQTT 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
上一篇/下一篇内容由系统自动关联 返回资讯列表 →