上位機(jī)把PLC數(shù)據(jù)推上MQTT_三_離線隊(duì)列與QoS1可靠投遞)
工業(yè)上位機(jī)把 PLC 數(shù)據(jù)推上 MQTT三離線隊(duì)列 QoS1斷網(wǎng)不丟、連上補(bǔ)發(fā)系列說明這是《工業(yè)上位機(jī)把 PLC 數(shù)據(jù)推上 MQTT》三部曲第 3 篇也是收尾篇。第 1 篇講了整體架構(gòu)和發(fā)布鏈路第 2 篇講了死區(qū) 最小間隔流量治理本篇講斷網(wǎng)了數(shù)據(jù)怎么辦。前兩篇(一)架構(gòu)與完整鏈路 · (二)死區(qū)與節(jié)流前兩篇把發(fā)什么、發(fā)多頻繁定了。最后這篇講最要命的網(wǎng)絡(luò)抖一下那段時(shí)間的數(shù)據(jù)不能丟。車間網(wǎng)絡(luò)什么德行干過現(xiàn)場(chǎng)的都懂——交換機(jī)重啟、光纖被挖斷、4G 基站抽風(fēng)短則幾秒長則幾十分鐘。要是沒有兜底這段時(shí)間 PLC 的數(shù)據(jù)就真空了MES 那邊對(duì)賬對(duì)不平工藝組第一個(gè)找你。我們用了兩層保險(xiǎn)QoS1 保發(fā)出去的被確認(rèn)離線隊(duì)列保沒發(fā)出去的先存著。一、QoS1至少一次靠確認(rèn)MQTT 的 QoS 三級(jí)里設(shè)備上報(bào)一般用 QoS1至少一次。QoS0 是發(fā)了就不管工廠數(shù)據(jù)不敢用QoS2恰好一次握手太重工業(yè)場(chǎng)景性價(jià)比低。QoS1 的語義是Broker 收到必須回 PUBACK沒收到發(fā)布端就重發(fā)。關(guān)鍵在重發(fā)怎么跟蹤。項(xiàng)目里有個(gè)容易混的三個(gè)概念我在qos1state.h里專門分開了// qos1state.h// msgId 應(yīng)用層持久標(biāo)識(shí)單調(diào)遞增用于日志/去重/對(duì)賬// token Paho MQTTAsync_tokenint發(fā)送時(shí)庫返回用于投遞完成回調(diào)關(guān)聯(lián)// packetId 線上 Packet IdentifierPaho 內(nèi)部不暴露混了到底會(huì)怎樣真踩過的坑早期版本我圖省事把 Paho 回調(diào)里拿到的token直接當(dāng)應(yīng)用層msgId去查m_byToken表。結(jié)果 Paho 的token是發(fā)送批次維度的同一批次多條消息共用一個(gè) token而m_byToken是按單條消息建的——onAck(token)命中后只清了表里一條剩下幾條永遠(yuǎn)停在Inflight。表象就是狀態(tài)面板inflight只增不減、重連后這些幽靈消息被全量補(bǔ)發(fā)翻倍消費(fèi)端 ts 去重都救不回來因?yàn)?ts 雖同但根本沒發(fā)成功過被當(dāng)成新消息又發(fā)一遍。教訓(xùn)token 只用來關(guān)聯(lián)這次投遞完成回調(diào)msgId 才是業(yè)務(wù)去重/對(duì)賬的主鍵兩者必須分開存。發(fā)出去時(shí)建記錄收到 ACK 時(shí)清記錄發(fā)出去時(shí)建記錄收到 ACK 時(shí)清記錄voidQos1Tracker::markInflight(qint64 msgId,MQTTAsync_token token){Qos1Record r;r.msgIdmsgId;r.tokentoken;r.statusQos1Status::Inflight;m_byToken.insert(token,r);// 按 token 關(guān)聯(lián)}voidQos1Tracker::onAck(MQTTAsync_token token,Qos1Status st){autoitm_byToken.find(token);// 投遞完成回調(diào)按 token 命中if(it!m_byToken.end()){it-statusst;m_byToken.erase(it);}}m_qos1.inflightCount()直接喂給狀態(tài)面板運(yùn)維一眼能看到還有幾條沒被 Broker 確認(rèn)。一個(gè)必須說清楚的設(shè)計(jì)點(diǎn)斷線重連后之前 QoS1 沒被確認(rèn)的消息會(huì)被全量補(bǔ)發(fā)這必然產(chǎn)生重復(fù)。這是 QoS1 的本職不是 bug。解決辦法在消費(fèi)端——按(tagKey, ts)做冪等去重。ts是發(fā)布時(shí)就帶上的采樣時(shí)刻重復(fù)的消息ts一樣落庫時(shí)去重即可。所以我們的 change 載荷里永遠(yuǎn)帶著ts就是這個(gè)用處{tag:DefaultPLC/Main_Plc/回水溫度,value:52.3,quality:GOOD,ts:2026-09-28T09:12:00.123Z}二、離線隊(duì)列連不上就先存盤QoS1 只管發(fā)出去→確認(rèn)。但要是壓根沒連上 BrokersendRaw直接返回 false消息得有個(gè)地方先待著。這就是enqueue→ 離線隊(duì)列。隊(duì)列分兩級(jí)內(nèi)存不夠再落盤// mqttpublisher.cpp::enqueueconstintcapm_cfg?m_cfg-maxQueueMem:2000;// 內(nèi)存上限 2000 條if(m_memQueue.size()cap){m_memQueue.append(m);}elseif(m_dbOk){diskInsert(channel,topic,payload,qos);// 超限轉(zhuǎn) SQLite 落盤}else{m_memQueue.append(m);// 無磁盤兜底best-effort}落盤用的是 SQLite獨(dú)立連接一張outbound表CREATETABLEIFNOTEXISTSoutbound(idINTEGERPRIMARYKEYAUTOINCREMENT,channelTEXT,topicTEXT,payloadBLOB,qosINT,created_atINTEGER);兩個(gè)細(xì)節(jié)是現(xiàn)場(chǎng)能救命的① 有界不撐爆磁盤。落盤也有限額默認(rèn) 2MB超了刪最舊一行if(m_cfg(m_diskBytespayload.size()m_cfg-maxQueueDiskBytes)){del.exec(DELETE FROM outbound WHERE id (SELECT MIN(id) FROM outbound));}網(wǎng)絡(luò)中斷半小時(shí)隊(duì)列把最早的歷史讓給最新的實(shí)時(shí)這符合監(jiān)控語義——寧可丟最老的也要保最新的。② 并發(fā)保護(hù)。離線隊(duì)列的寫落盤和讀沖刷可能跨線程搶 SQLite我們?cè)O(shè)了QSQLITE_BUSY_TIMEOUT5000避免直接SQLITE_BUSY報(bào)錯(cuò)丟數(shù)據(jù)m_db.setConnectOptions(QSQLITE_BUSY_TIMEOUT5000);// 關(guān)鍵避免 SQLITE_BUSY三、連上即補(bǔ)發(fā)flushQueue重連成功后第一件事是把攢著的消息吐出去// mqttpublisher.cpp::handleConnected → flushQueuevoidMqttPublisher::flushQueue(){while(!m_memQueue.isEmpty()){// 先內(nèi)存隊(duì)列先進(jìn)先出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);// 發(fā)成功才刪失敗留著下輪}}}關(guān)于順序先內(nèi)存后磁盤會(huì)不會(huì)亂序入隊(duì)邏輯是內(nèi)存沒滿就append滿了size()cap才落盤所以內(nèi)存里是更早的消息、磁盤里是更晚的消息沖刷時(shí)內(nèi)存 FIFO內(nèi)升序 磁盤id ASC內(nèi)升序單次連續(xù)斷網(wǎng)、一次沖刷到底時(shí)消費(fèi)方收到的是全局時(shí)間升序沒問題。唯一的邊角場(chǎng)景上輪沖刷到一半又?jǐn)唷⑶倚孪褍?nèi)存填滿溢出到磁盤磁盤里殘留的更早消息會(huì)排在內(nèi)存新消息之后出現(xiàn)局部倒掛。生產(chǎn)硬化做法見下——按createdAt合并排序再發(fā)所有場(chǎng)景都穩(wěn)。注意發(fā)成功才刪沖刷中途又?jǐn)嗔藄endRaw失敗就break沒發(fā)的留著下一輪連上接著補(bǔ)。這條鏈路接在第 1 篇的publishFormatted末尾——連不上走enqueue連上了走flushQueue閉環(huán)。3.1 生產(chǎn)硬化①按創(chuàng)建時(shí)間合并排序徹底杜絕倒掛把內(nèi)存和磁盤的消息一起按createdAt升序排好再發(fā)。需要diskLoad()順帶返回created_at列表里本來就有// 合并內(nèi)存 磁盤 → 按 createdAt 升序 → 逐個(gè)發(fā)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 生產(chǎn)硬化②限速分批沖刷防消息風(fēng)暴MQTTAsync_send是異步非阻塞緊循環(huán)一次能把內(nèi)存 2000 條 磁盤數(shù)萬條在毫秒級(jí)全甩出去打滿 Broker 的max_inflight_messages、或占滿 4G 帶寬。用定時(shí)器分批// m_flushPerTick 默認(rèn) 50、間隔 50ms≈1000 條/秒現(xiàn)場(chǎng)按 Broker 能力調(diào)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);}落盤隊(duì)列也別一次性SELECT全表進(jìn)內(nèi)存默認(rèn) 2MB 有界但數(shù)萬行仍占一塊diskLoad改成LIMIT分批游標(biāo)、邊取邊發(fā)和上面的限速?zèng)_刷合在一起最穩(wěn)。3.3 生產(chǎn)硬化③過期丟棄TTL斷網(wǎng) 2 小時(shí)重連后把 2 小時(shí)前的溫度補(bǔ)發(fā)給實(shí)時(shí)看板消費(fèi)方可能誤判當(dāng)前值。配置加maxAgeMs0不過期沖刷/落盤前丟超期消息// 落盤/沖刷前判斷now - createdAt maxAgeMs 則丟棄if(m_cfg-maxAgeMs0(now-m.createdAtm_cfg-maxAgeMs)){if(m.fromDisk)diskDelete(m.id);// 內(nèi)存的直接跳過continue;}更徹底的辦法是升級(jí)MQTT 5.0 的Message Expiry Interval讓 Broker 自動(dòng)丟棄超期補(bǔ)發(fā)見第七節(jié)。四、一個(gè)真踩過的線程坑Paho 的 C 回調(diào)onConnectionLost/onDeliveryComplete等是在庫自己的網(wǎng)絡(luò)線程里觸發(fā)的。我們一開始在回調(diào)里直接寫 SQLite、等 ACK、加鎖——結(jié)果偶發(fā)死鎖UI 卡死。根因回調(diào)線程不能碰主線程的資源SQLite 連接、QMutex 持有的業(yè)務(wù)狀態(tài)。方案是回調(diào)里只 marshal 回主線程絕不阻塞// mqttpublisher.cpp::onDeliveryComplete網(wǎng)絡(luò)線程觸發(fā)voidMqttPublisher::onDeliveryComplete(void*context,MQTTAsync_token token){auto*selfstatic_castMqttPublisher*(context);QMetaObject::invokeMethod(self,[self,token](){self-m_qos1.onAck(token,Qos1Status::Acked);// 回主線程再改狀態(tài)},Qt::QueuedConnection);}Qt::QueuedConnection把活兒排隊(duì)到主線程事件循環(huán)網(wǎng)絡(luò)線程立刻返回。業(yè)務(wù)狀態(tài)inflight 計(jì)數(shù)、離線隊(duì)列讀寫永遠(yuǎn)只在主線程動(dòng)死鎖消失。這條規(guī)矩寫在頭文件注釋里回調(diào)禁止阻塞寫 SQLite/等 ACK/加鎖都會(huì)死鎖統(tǒng)一marshal回主線程。五、報(bào)警通道繞過節(jié)流、走 QoS1、靠 ts 去重第 1 篇提過報(bào)警是獨(dú)立通道這里把它的特殊待遇說清。publishAlarm()直接formatAlarm → publishFormatted根本不調(diào)shouldPublish——死區(qū)和最小間隔都攔不到它。這恰恰是對(duì)的報(bào)警事件絕不能因?yàn)橹禌]變夠多或離上次太近被吞掉該報(bào)就報(bào)。報(bào)警 QoS 走m_cfg-qos默認(rèn) 1。有人問報(bào)警要不要 QoS2 防重復(fù)“我的建議是不用報(bào)警最大的風(fēng)險(xiǎn)是漏報(bào)而不是重復(fù)報(bào)”QoS1 消費(fèi)端按(tagKey, ts)冪等去重已經(jīng)夠用QoS2 握手更重而且重復(fù)問題同樣靠 ts 去重解決上 QoS2 是虧本買賣。六、Broker 端配套配置清單上位機(jī)再穩(wěn)也得 Broker 配合斷線期間 Broker 沒配好消息照樣丟。Mosquitto 最低配置參考# mosquitto.conf max_queued_messages 0 # 0不限制 Broker 側(cè)隊(duì)列或按內(nèi)存設(shè)大配合上位機(jī)限速 max_inflight_messages 100 # 單客戶端在途上限避免一個(gè)客戶端占滿 persistence true # 開啟持久化Broker 重啟不丟 retained / 會(huì)話 # 若上位機(jī) cleanSessionfalse我們默認(rèn)就是 false務(wù)必開持久會(huì)話 # 否則斷線期間訂閱關(guān)系與會(huì)話丟失重連后收不到補(bǔ)發(fā)EMQX 對(duì)應(yīng)mqtt.max_inflight、開啟retainer、配置session_expiry。給甲方交付時(shí)這份清單直接附上別光說自己客戶端可靠。七、可觀測(cè)性與運(yùn)維監(jiān)控status()已經(jīng)把關(guān)鍵指標(biāo)吐出來了運(yùn)維面板至少該掛這幾項(xiàng)離線隊(duì)列深度queuedMemqueuedDiskBytes超閾值比如內(nèi)存 80% 上限 / 磁盤 1MB就告警——意味著網(wǎng)絡(luò)長期不穩(wěn)或 Broker 收不動(dòng)inflight長時(shí)間不歸零 Broker 不回 PUBACK卡死 / 網(wǎng)絡(luò)半通告警發(fā)布失敗率 失敗計(jì)數(shù) /totalPublished暴露lastError進(jìn)日志斷連原因可追溯是onConnectionLost報(bào)的 cause還是MQTTAsync_send失敗。這幾項(xiàng)不展示斷網(wǎng)了你都不知道數(shù)據(jù)在丟。八、MQTT 5.0 與版本差異全文基于 PahoMQTTAsyncMQTT 3.1.1。工業(yè)場(chǎng)景直接能用上 5.0 的三個(gè)特性Message Expiry Interval消息級(jí)過期Broker 自動(dòng)丟棄超期補(bǔ)發(fā)正好解決第三節(jié) 3.3 的 TTL 問題比應(yīng)用層maxAgeMs更徹底Session Expiry Interval替代cleanSession那個(gè)別扭的布爾精細(xì)控制會(huì)話保留時(shí)長Shared Subscription多消費(fèi)者負(fù)載均衡適合一個(gè) topic 被多個(gè)后端搶著處理的場(chǎng)景。升級(jí)路徑建議先用 3.3 的應(yīng)用層maxAgeMs兜底等現(xiàn)場(chǎng) Broker 支持 5.0 再切原生過期平滑過渡。九、三篇串起來回到第 1 篇那張鏈路圖現(xiàn)在三道閘都齊了采集 → onTagValueChanged → shouldPublish 第2篇死區(qū)最小間隔砍掉沒用的 → formatChange JSON 格式化、帶 ts → publishFormatted ├─ sendRaw 成功 → m_totalPublished └─ sendRaw 失敗 → enqueue本篇內(nèi)存→落盤有界 重連成功 → flushQueue本篇連上即補(bǔ)發(fā)發(fā)成功才刪 QoS1 → markInflight / onAck本篇至少一次 冪等去重第 1 篇管架構(gòu)三通道、主題、JSON、異步客戶端第 2 篇管流量死區(qū) 最小間隔解決發(fā)太多第 3 篇管可靠QoS1 離線隊(duì)列解決斷網(wǎng)丟?,F(xiàn)場(chǎng)配的時(shí)候我的習(xí)慣先開 change 默認(rèn)死區(qū)0.5%看流量落不落得下來要?dú)v史全貌再開 snapshot最后確認(rèn)離線隊(duì)列上限按現(xiàn)場(chǎng)斷網(wǎng)時(shí)長估2MB 大概夠撐一陣真長斷網(wǎng)調(diào)大maxQueueDiskBytes。QoS 默認(rèn) 1 別動(dòng)消費(fèi)端記得按(tagKey, ts)去重。本篇配置速查卡可靠投遞配置默認(rèn)說明qos1默認(rèn) QoS1至少一次消費(fèi)端按(tagKey, ts)去重cleanSessionfalse持久會(huì)話斷線重連可續(xù)訂Broker 須開持久化keepAliveSec/connectTimeoutSec60/30心跳 / 連接超時(shí)reconnectBackoffSec5斷線重連退避maxQueueMem2000內(nèi)存隊(duì)列上限條maxQueueDiskBytes2MB落盤上限超了丟最舊保最新maxAgeMs0硬化新增補(bǔ)發(fā)消息最大齡期0不過期升級(jí) MQTT5 用Message Expiry更徹底m_flushPerTick/ 間隔50/50ms硬化新增限速分批沖刷≈1000 條/秒按 Broker 調(diào)完整性清單Broker 配置 / 可觀測(cè)性 / MQTT 5.0見本篇第六~八節(jié)。本文及 PLCMonitor 系列文章均為免費(fèi)分享。本文免費(fèi)分享如需轉(zhuǎn)載請(qǐng)聯(lián)系作者獲取授權(quán)。如果覺得這篇文章對(duì)你有幫助歡迎點(diǎn)贊收藏。源碼獲取地址https://github.com/freddiezhang1990/plcmonitor