形緩沖區(qū)(番外篇))
目錄前言Pipelines管道實例后記前言大家好我是 wacky。書接上回我們上回探討了如何解決TCP粘包和半包的問題并在最后引入了環(huán)形緩沖區(qū)的概念。雖然環(huán)形緩沖區(qū)主要是通過固定數(shù)組讀寫索引模運算來實現(xiàn)循環(huán)但是實際開發(fā)過程中還存在手動管理 offset、count數(shù)組擴容、數(shù)據(jù)移動容易越界、內(nèi)存拷貝多等一系列的問題。而實際上在.NET中依然有更簡潔高效的API來代替手寫環(huán)形緩沖區(qū)它就是Pipelines管道。如果還想回顧一下手寫環(huán)形緩沖區(qū)的概念可以從傳送門出發(fā).NET上位機踩坑為什么有時讀取數(shù)據(jù)需要SleepPipelines管道命名空間為System.IO.Pipelines它是.NET內(nèi)置的高性能內(nèi)存緩沖組件。之前在和一些技術(shù)大佬聊環(huán)形緩沖區(qū)的過程中他們很多人已經(jīng)把這個組件用于實際生產(chǎn)環(huán)境中了那我們今天就來講一講在解決Modbus協(xié)議TCP粘包和半包的問題中這個組件要怎么用。它主要包含以下幾部分內(nèi)容Pipe內(nèi)置一對PipeWriter 和PipeReaderWriter 負責(zé)往管道塞網(wǎng)絡(luò)收到的數(shù)據(jù)Reader 負責(zé)從管道讀取、解析報文。Pipe 內(nèi)部自帶環(huán)形緩沖。Socket 異步接收用Socket.ReceiveAsync 持續(xù)接收 PLC 下發(fā)的字節(jié)流寫入PipeWriter。TCP 是字節(jié)流沒有邊界這一步只管收字節(jié)不關(guān)心是不是完整幀。PipeReader 循環(huán)解析ReadAsync()從管道拿一段可用內(nèi)存無需拷貝直接內(nèi)存切片嘗試在這段內(nèi)存里查找完整 ModbusTCP 幀MBAP 頭固定 7 字節(jié)事務(wù) ID (2) 協(xié)議 ID (2) 長度 (2) 單元 ID (1)后面跟 N 個功能碼數(shù)據(jù)如果找到完整幀AdvanceTo 標(biāo)記已經(jīng)消費掉的字節(jié)交給業(yè)務(wù)處理剩下半包留在管道緩沖區(qū)下次繼續(xù)解析如果不夠一幀停止解析等待后續(xù) Socket 繼續(xù)收到數(shù)據(jù)寫入管道4. 斷開 / 異常Complete PipeReader/PipeWriter釋放資源。實例現(xiàn)在我們來結(jié)合C#實例繼續(xù)講解internal class ModbusTcpPipelinesClient { private readonly string _ip; private readonly int _port; private Socket _socket; private Pipe _pipe; private Task _readTask; private Task _writeTask; private CancellationTokenSource _cts; // 收到完整Modbus報文回調(diào)上位機在這里解析數(shù)據(jù) public Actionbyte[] OnModbusFrameReceived { get; set; } public ModbusTcpPipelinesClient(string ip, int port 502) { _ip ip; _port port; } public async Task ConnectAsync() { _cts new CancellationTokenSource(); _socket new Socket(AddressFamily.InterNetwork, SocketType.Stream, ProtocolType.Tcp); await _socket.ConnectAsync(_ip, _port, _cts.Token); _pipe new Pipe(new PipeOptions(pauseWriterThreshold: 4096, resumeWriterThreshold: 2048)); // 兩個獨立任務(wù)1. Socket收數(shù)據(jù)寫入PipeWriter2. PipeReader解析報文 _writeTask FillPipeFromSocketAsync(_socket, _pipe.Writer, _cts.Token); _readTask ParsePipeFramesAsync(_pipe.Reader, _cts.Token); } // 從Socket接收字節(jié)寫入PipeWriter相當(dāng)于填充緩沖區(qū) private async Task FillPipeFromSocketAsync(Socket socket, PipeWriter writer, CancellationToken ct) { try { while (!ct.IsCancellationRequested) { // 獲取一塊內(nèi)存最小分配512字節(jié)可根據(jù)PLC調(diào)大小 Memorybyte memory writer.GetMemory(512); int readLen await socket.ReceiveAsync(memory, SocketFlags.None, ct); if (readLen 0) break; // 遠端關(guān)閉連接 writer.Advance(readLen); // 告訴writer實際收到多少字節(jié) FlushResult flushResult await writer.FlushAsync(ct); if (flushResult.IsCompleted) break; } } catch (Exception ex) { Console.WriteLine($接收異常:{ex.Message}); } finally { writer.Complete(); } } // PipeReader 循環(huán)解析Modbus TCP報文核心替代環(huán)形緩沖區(qū)解析 private async Task ParsePipeFramesAsync(PipeReader reader, CancellationToken ct) { try { while (!ct.IsCancellationRequested) { ReadResult result await reader.ReadAsync(ct); ReadOnlySequencebyte buffer result.Buffer; SequencePosition? consumedPos null; // 循環(huán)解析緩沖區(qū)里所有完整Modbus幀處理粘包一次多個報文 while (TryParseModbusFrame(buffer, out var frame, out var consumed)) { OnModbusFrameReceived?.Invoke(frame.ToArray()); consumedPos consumed; buffer buffer.Slice(consumed); // 切掉已經(jīng)解析完的數(shù)據(jù) } // AdvanceTo標(biāo)記消費位置、查看位置Pipelines自動回收內(nèi)存 reader.AdvanceTo(consumedPos ?? buffer.Start, buffer.End); if (result.IsCompleted) break; } } catch (Exception ex) { Console.WriteLine($解析異常:{ex.Message}); } finally { reader.Complete(); } } /// summary /// 嘗試從ReadOnlySequence解析ModbusTCP幀 /// MBAP: 7字節(jié) [TransId(2)ProtoId(2)Len(2)UnitId(1)] PDU /// /summary private bool TryParseModbusFrame(in ReadOnlySequencebyte seq, out ReadOnlySequencebyte frame, out SequencePosition consumed) { frame default; consumed seq.Start; if (seq.Length 7) return false; // 不足MBAP頭半包等待更多數(shù)據(jù) // 讀取MBAP第5、6字節(jié)PDU長度 var headerReader new SequenceReaderbyte(seq); headerReader.Advance(4); headerReader.TryReadBigEndian(out short pduLen); int totalFrameLen 7 pduLen; if (seq.Length totalFrameLen) return false; // 收到的數(shù)據(jù)不夠完整報文等待 frame seq.Slice(0, totalFrameLen); consumed seq.GetPosition(totalFrameLen); return true; } // 發(fā)送Modbus請求 public async Task SendAsync(byte[] data) { if (_socket null || !_socket.Connected) throw new InvalidOperationException(未連接PLC); int sendTotal 0; while (sendTotal data.Length) { int sent await _socket.SendAsync(data.AsMemory(sendTotal), SocketFlags.None); sendTotal sent; } } public async Task CloseAsync() { _cts?.Cancel(); try { await Task.WhenAll(_readTask, _writeTask); } catch { } _socket?.Close(); _pipe null; } }在上述代碼段中我們把寫入和解析作為2個異步方法分開來執(zhí)行用于職責(zé)分離。socket在異步拿到數(shù)據(jù)后會告訴PipeWriter實際收到了多少字節(jié)然后把數(shù)據(jù)提交給PipeReader用于讀取。在這一步socket的職責(zé)只有把TCP字節(jié)流輸入管道中不做任何的協(xié)議解析。而在解析的方法中只負責(zé)拿出管道中當(dāng)前可用的所有數(shù)據(jù)進行解析這里引入了ReadOnlySequence的概念這個類型是代表跨多個內(nèi)存塊的連續(xù)邏輯字節(jié)流不需要合并數(shù)組。其中又包含一個我們自定義的TryParseModbusFrame方法用于解析ModbusTCP幀。在TryParseModbusFrame方法內(nèi)部我們處理了半包或者粘包的場景由于一個完整的MBAP(Modbus TCP)報文頭的長度是7字節(jié)因此我們在不足7字節(jié)的時候直接return false判斷此為半包不會繼續(xù)進行解析。然后我們會計算整個幀的總長度totalFrameLen 7 pduLen這里我們會繼續(xù)判斷收到的數(shù)據(jù)是否為完整的報文長度如果不足那么依然return false繼續(xù)等待數(shù)據(jù)完整后將buffer切片輸出給調(diào)用方。最后我們通過AdvanceTo來標(biāo)記消費位置、查看位置用于下次繼續(xù)從指定的位置繼續(xù)讀取這樣我們就可以有效解決TCP存在粘包和半包的問題。后記現(xiàn)代編程語言API的進化總是讓人欣喜直接引入Pipelines就可以避免手寫RingBuffer存在的很多問題可謂是大大提升了效率。但是我們依然還是需要知其然知其所以然明白了存在什么問題再去根據(jù)問題的本質(zhì)進行解決懂得怎么解決之后再做改進各位不知今天看懂了嗎引入地址