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