.NET上位机踩坑:用Pipelines替代环形缓冲区(番外篇)

发布时间:2026/9/30 21:24:50
.NET上位机踩坑:用Pipelines替代环形缓冲区(番外篇) 目录前言Pipelines管道实例后记前言大家好我是 wacky。书接上回我们上回探讨了如何解决TCP粘包和半包的问题并在最后引入了环形缓冲区的概念。虽然环形缓冲区主要是通过固定数组读写索引模运算来实现循环但是实际开发过程中还存在手动管理 offset、count数组扩容、数据移动容易越界、内存拷贝多等一系列的问题。而实际上在.NET中依然有更简洁高效的API来代替手写环形缓冲区它就是Pipelines管道。如果还想回顾一下手写环形缓冲区的概念可以从传送门出发.NET上位机踩坑为什么有时读取数据需要SleepPipelines管道命名空间为System.IO.Pipelines它是.NET内置的高性能内存缓冲组件。之前在和一些技术大佬聊环形缓冲区的过程中他们很多人已经把这个组件用于实际生产环境中了那我们今天就来讲一讲在解决Modbus协议TCP粘包和半包的问题中这个组件要怎么用。它主要包含以下几部分内容Pipe内置一对PipeWriter 和PipeReaderWriter 负责往管道塞网络收到的数据Reader 负责从管道读取、解析报文。Pipe 内部自带环形缓冲。Socket 异步接收用Socket.ReceiveAsync 持续接收 PLC 下发的字节流写入PipeWriter。TCP 是字节流没有边界这一步只管收字节不关心是不是完整帧。PipeReader 循环解析ReadAsync()从管道拿一段可用内存无需拷贝直接内存切片尝试在这段内存里查找完整 ModbusTCP 帧MBAP 头固定 7 字节事务 ID (2) 协议 ID (2) 长度 (2) 单元 ID (1)后面跟 N 个功能码数据如果找到完整帧AdvanceTo 标记已经消费掉的字节交给业务处理剩下半包留在管道缓冲区下次继续解析如果不够一帧停止解析等待后续 Socket 继续收到数据写入管道4. 断开 / 异常Complete PipeReader/PipeWriter释放资源。实例现在我们来结合C#实例继续讲解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报文回调上位机在这里解析数据 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)); // 两个独立任务1. Socket收数据写入PipeWriter2. PipeReader解析报文 _writeTask FillPipeFromSocketAsync(_socket, _pipe.Writer, _cts.Token); _readTask ParsePipeFramesAsync(_pipe.Reader, _cts.Token); } // 从Socket接收字节写入PipeWriter相当于填充缓冲区 private async Task FillPipeFromSocketAsync(Socket socket, PipeWriter writer, CancellationToken ct) { try { while (!ct.IsCancellationRequested) { // 获取一块内存最小分配512字节可根据PLC调大小 Memorybyte memory writer.GetMemory(512); int readLen await socket.ReceiveAsync(memory, SocketFlags.None, ct); if (readLen 0) break; // 远端关闭连接 writer.Advance(readLen); // 告诉writer实际收到多少字节 FlushResult flushResult await writer.FlushAsync(ct); if (flushResult.IsCompleted) break; } } catch (Exception ex) { Console.WriteLine($接收异常:{ex.Message}); } finally { writer.Complete(); } } // PipeReader 循环解析Modbus TCP报文核心替代环形缓冲区解析 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; // 循环解析缓冲区里所有完整Modbus帧处理粘包一次多个报文 while (TryParseModbusFrame(buffer, out var frame, out var consumed)) { OnModbusFrameReceived?.Invoke(frame.ToArray()); consumedPos consumed; buffer buffer.Slice(consumed); // 切掉已经解析完的数据 } // AdvanceTo标记消费位置、查看位置Pipelines自动回收内存 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字节 [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头半包等待更多数据 // 读取MBAP第5、6字节PDU长度 var headerReader new SequenceReaderbyte(seq); headerReader.Advance(4); headerReader.TryReadBigEndian(out short pduLen); int totalFrameLen 7 pduLen; if (seq.Length totalFrameLen) return false; // 收到的数据不够完整报文等待 frame seq.Slice(0, totalFrameLen); consumed seq.GetPosition(totalFrameLen); return true; } // 发送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个异步方法分开来执行用于职责分离。socket在异步拿到数据后会告诉PipeWriter实际收到了多少字节然后把数据提交给PipeReader用于读取。在这一步socket的职责只有把TCP字节流输入管道中不做任何的协议解析。而在解析的方法中只负责拿出管道中当前可用的所有数据进行解析这里引入了ReadOnlySequence的概念这个类型是代表跨多个内存块的连续逻辑字节流不需要合并数组。其中又包含一个我们自定义的TryParseModbusFrame方法用于解析ModbusTCP帧。在TryParseModbusFrame方法内部我们处理了半包或者粘包的场景由于一个完整的MBAP(Modbus TCP)报文头的长度是7字节因此我们在不足7字节的时候直接return false判断此为半包不会继续进行解析。然后我们会计算整个帧的总长度totalFrameLen 7 pduLen这里我们会继续判断收到的数据是否为完整的报文长度如果不足那么依然return false继续等待数据完整后将buffer切片输出给调用方。最后我们通过AdvanceTo来标记消费位置、查看位置用于下次继续从指定的位置继续读取这样我们就可以有效解决TCP存在粘包和半包的问题。后记现代编程语言API的进化总是让人欣喜直接引入Pipelines就可以避免手写RingBuffer存在的很多问题可谓是大大提升了效率。但是我们依然还是需要知其然知其所以然明白了存在什么问题再去根据问题的本质进行解决懂得怎么解决之后再做改进各位不知今天看懂了吗引入地址