You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何基于抽象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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.22 13:36:31