拓十年匠心定制 · 商业建站与技术教学双线并行 咨询热线:400-886-1026 service@lmnt.cn
ARTICLE DETAIL

资讯详情

深耕网站建设与运营推广的一线实战洞察。

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

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

目录

前言

Pipelines管道

实例

后记


前言

大家好,我是 wacky。

书接上回,我们上回探讨了如何解决TCP粘包和半包的问题,并在最后引入了环形缓冲区的概念。虽然环形缓冲区主要是通过固定数组+读写索引+模运算来实现循环,但是实际开发过程中,还存在手动管理 offset、count,数组扩容、数据移动,容易越界、内存拷贝多等一系列的问题。

而实际上在.NET中,依然有更简洁高效的API来代替手写环形缓冲区,它就是Pipelines管道。如果还想回顾一下手写环形缓冲区的概念,可以从传送门出发:.NET上位机踩坑:为什么有时读取数据需要Sleep?

Pipelines管道

命名空间为System.IO.Pipelines,它是.NET内置的高性能内存缓冲组件。之前在和一些技术大佬聊环形缓冲区的过程中,他们很多人已经把这个组件用于实际生产环境中了,那我们今天就来讲一讲,在解决Modbus协议TCP粘包和半包的问题中这个组件要怎么用。

它主要包含以下几部分内容:

  1. Pipe:内置一对PipeWriter 和PipeReader,Writer 负责往管道塞网络收到的数据;Reader 负责从管道读取、解析报文。Pipe 内部自带环形缓冲。

  2. Socket 异步接收:用Socket.ReceiveAsync 持续接收 PLC 下发的字节流,写入PipeWriter。TCP 是字节流,没有边界,这一步只管收字节,不关心是不是完整帧。

  3. 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 Action<byte[]> 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收数据写入PipeWriter;2. 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调大小 Memory<byte> 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); ReadOnlySequence<byte> 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 ReadOnlySequence<byte> seq, out ReadOnlySequence<byte> frame, out SequencePosition consumed) { frame = default; consumed = seq.Start; if (seq.Length < 7) return false; // 不足MBAP头,半包,等待更多数据 // 读取MBAP第5、6字节:PDU长度 var headerReader = new SequenceReader<byte>(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存在的很多问题,可谓是大大提升了效率。但是我们依然还是需要"知其然,知其所以然",明白了存在什么问题,再去根据问题的本质进行解决,懂得怎么解决之后,再做改进,各位不知今天看懂了吗?

引入地址

返回列表