如何基于抽象Stream类异步读取NetworkStream/SslStream且不阻塞?
问题背景
我希望通过抽象父类Stream,从NetworkStream或SslStream异步读取数据。流的异步读取方式有三种:
- 异步编程模型(APM):使用
BeginRead和EndRead操作。 - 任务并行库(TPL):使用
Task并创建任务延续。 - 基于任务的异步模式(TAP):操作以Async为后缀,可使用
async和await关键字。
我主要关注用TAP模式实现异步读取。以下代码可异步读取至流末尾,返回字节数组形式的数据:
internal async Task<byte[]> ReadToEndAsync(Stream stream) { byte[] buffer = new byte[4096]; using (MemoryStream memoryStream = new MemoryStream()) { int bytesRead = await stream.ReadAsync(buffer, 0, buffer.Length); while (bytesRead != 0) { // Received datas were aggregated to a memory stream. await memoryStream.WriteAsync(buffer, 0, bytesRead); bytesRead = await stream.ReadAsync(buffer, 0, buffer.Length); } return memoryStream.ToArray(); } }
缓冲区大小为4096字节,数据超过缓冲区时会持续读取至返回0(流结束)。该方法在FileStream上正常,但在NetworkStream或SslStream的ReadAsync操作处无限挂起。
原因是网络流的ReadAsync仅在Socket通信关闭时返回0,而我不希望每次传输数据都关闭通信。
问题
如何在不关闭Socket通信的情况下,避免ReadAsync的阻塞调用?
解决方案
核心思路是让读取操作明确知道当前数据块的结束标记,而非等待流关闭。以下是几种可靠的实现方式:
1. 先传输数据长度(推荐方案)
这是网络通信中最常用的可靠方案,通过"长度前缀+实际数据"的协议格式,让接收方明确知道要读取多少字节:
- 发送方先写入一个固定长度的整数(比如4字节的
int),表示后续数据的总字节数。 - 接收方先读取这个长度值,再循环读取直到获取到指定字节数的内容。
示例代码:
internal async Task<byte[]> ReadExactLengthAsync(Stream stream) { // 读取数据长度(4字节int) byte[] lengthBuffer = new byte[4]; int bytesRead = await stream.ReadAsync(lengthBuffer, 0, lengthBuffer.Length); if (bytesRead != lengthBuffer.Length) throw new InvalidOperationException("未能读取完整的数据长度"); int dataLength = BitConverter.ToInt32(lengthBuffer, 0); if (dataLength <= 0) return Array.Empty<byte>(); // 读取指定长度的数据 using (MemoryStream memoryStream = new MemoryStream(dataLength)) { int remainingBytes = dataLength; byte[] buffer = new byte[Math.Min(4096, remainingBytes)]; while (remainingBytes > 0) { bytesRead = await stream.ReadAsync(buffer, 0, Math.Min(buffer.Length, remainingBytes)); if (bytesRead == 0) throw new IOException("连接意外中断"); await memoryStream.WriteAsync(buffer, 0, bytesRead); remainingBytes -= bytesRead; } return memoryStream.ToArray(); } }
发送方对应逻辑:
internal async Task SendDataAsync(Stream stream, byte[] data) { byte[] lengthBytes = BitConverter.GetBytes(data.Length); await stream.WriteAsync(lengthBytes, 0, lengthBytes.Length); await stream.WriteAsync(data, 0, data.Length); await stream.FlushAsync(); }
2. 使用特殊分隔符标记数据结束
如果数据内容中不会出现特定的字节序列(比如0xFF 0xFF 0xFF),可以用该序列作为数据块的结束标记:
- 接收方持续读取数据,直到检测到分隔符为止。
- 需处理缓冲区中部分包含分隔符的情况,避免截断有效数据。
示例代码(简化版):
internal async Task<byte[]> ReadUntilDelimiterAsync(Stream stream, byte[] delimiter) { using (MemoryStream memoryStream = new MemoryStream()) { byte[] buffer = new byte[4096]; int bytesRead; while (true) { bytesRead = await stream.ReadAsync(buffer, 0, buffer.Length); if (bytesRead == 0) throw new IOException("连接意外中断"); // 检查缓冲区中是否包含分隔符 int delimiterIndex = FindDelimiter(buffer, bytesRead, delimiter); if (delimiterIndex != -1) { // 写入分隔符之前的内容 await memoryStream.WriteAsync(buffer, 0, delimiterIndex); break; } // 写入全部缓冲区内容 await memoryStream.WriteAsync(buffer, 0, bytesRead); } return memoryStream.ToArray(); } } // 辅助方法:查找分隔符在缓冲区中的位置 private int FindDelimiter(byte[] buffer, int bufferLength, byte[] delimiter) { for (int i = 0; i <= bufferLength - delimiter.Length; i++) { bool match = true; for (int j = 0; j < delimiter.Length; j++) { if (buffer[i + j] != delimiter[j]) { match = false; break; } } if (match) return i; } return -1; }
3. 设置读取超时(临时替代方案)
若无法修改通信协议,可给ReadAsync设置超时时间,超时后停止读取。但这种方式可能截断正常数据,仅适合非关键场景:
internal async Task<byte[]> ReadWithTimeoutAsync(Stream stream, int timeoutMs) { using (MemoryStream memoryStream = new MemoryStream()) using (CancellationTokenSource cts = new CancellationTokenSource(timeoutMs)) { byte[] buffer = new byte[4096]; try { while (true) { int bytesRead = await stream.ReadAsync(buffer, 0, buffer.Length, cts.Token); if (bytesRead == 0) break; await memoryStream.WriteAsync(buffer, 0, bytesRead); } } catch (OperationCanceledException) { // 超时,返回已读取的数据 } return memoryStream.ToArray(); } }
总结
优先选择先传输数据长度的方案,它稳定可靠,适配绝大多数网络通信场景。分隔符方案适合文本类或有明确边界的数据,超时方案仅作为临时替代,不推荐用于生产环境。
内容的提问来源于stack exchange,提问作者Péter Szilvási
相关产品推荐
相关产品推荐

