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

如何用.NET的Utf8JsonReader从网络流逐读完整JSON对象?

在.NET中流式读取网络流中的JSON对象

可以实现流式读取,不需要完全手动拆分JSON消息,但需要结合Utf8JsonReader和自定义缓冲区管理——Utf8JsonReader本身仅处理字节序列,不直接对接Stream,但通过维护缓冲区可以实现暂停-续读的逻辑。

核心实现思路

Utf8JsonReader支持逐token解析字节数据,当遇到不完整的JSON时会停止读取,此时只需从流中获取更多数据追加到缓冲区,就能继续解析。关键是跟踪解析状态,识别完整的根级JSON对象(即顶级{}),并保留未处理的字节供下次解析使用。

具体代码实现

以下是一个异步方法,可从网络流中逐个返回完整的JSON对象:

using System.Buffers;
using System.Text.Json;

public static async IAsyncEnumerable<JsonElement> ReadJsonObjectsFromStreamAsync(Stream stream, CancellationToken cancellationToken = default)
{
    // 复用字节数组减少内存分配,初始大小可根据业务调整
    var buffer = ArrayPool<byte>.Shared.Rent(4096);
    var remainingBuffer = new Memory<byte>();
    
    try
    {
        int bytesRead;
        while ((bytesRead = await stream.ReadAsync(buffer, cancellationToken).ConfigureAwait(false)) > 0)
        {
            // 合并剩余未处理数据与新读取的字节
            var combinedLength = remainingBuffer.Length + bytesRead;
            var combinedBuffer = combinedLength > buffer.Length 
                ? ArrayPool<byte>.Shared.Rent(combinedLength) 
                : buffer;
            
            remainingBuffer.CopyTo(combinedBuffer);
            new Memory<byte>(buffer, 0, bytesRead).CopyTo(combinedBuffer.AsMemory(remainingBuffer.Length));
            
            var reader = new Utf8JsonReader(combinedBuffer.AsMemory(0, combinedLength), isFinalBlock: false, state: default);
            int objectStartIndex = 0;
            
            while (reader.Read())
            {
                if (reader.TokenType == JsonTokenType.StartObject)
                {
                    // 记录根对象的起始位置
                    objectStartIndex = (int)reader.BytesConsumed;
                }
                else if (reader.TokenType == JsonTokenType.EndObject && reader.CurrentDepth == 0)
                {
                    // 提取完整的根JSON对象字节
                    var objectBytes = combinedBuffer.AsMemory(objectStartIndex, (int)(reader.BytesConsumed - objectStartIndex));
                    using var doc = JsonDocument.Parse(objectBytes);
                    yield return doc.RootElement.Clone();
                    
                    // 更新剩余缓冲区为未处理的字节
                    remainingBuffer = combinedBuffer.AsMemory((int)reader.BytesConsumed);
                    // 重置阅读器,从剩余数据开始解析
                    reader = new Utf8JsonReader(remainingBuffer, isFinalBlock: false, state: default);
                    objectStartIndex = 0;
                }
            }
            
            // 释放临时租用的缓冲区(如果不是原buffer的话)
            if (combinedBuffer != buffer)
            {
                ArrayPool<byte>.Shared.Return(combinedBuffer);
            }
            
            // 保留未处理的剩余数据
            if (reader.BytesConsumed < remainingBuffer.Length)
            {
                remainingBuffer = remainingBuffer.Slice((int)reader.BytesConsumed);
            }
            else
            {
                remainingBuffer = Memory<byte>.Empty;
            }
        }
        
        // 处理流结束时剩余的最后一个完整JSON对象
        if (!remainingBuffer.IsEmpty)
        {
            using var doc = JsonDocument.Parse(remainingBuffer);
            yield return doc.RootElement.Clone();
        }
    }
    finally
    {
        ArrayPool<byte>.Shared.Return(buffer);
        if (!remainingBuffer.IsEmpty && remainingBuffer.IsOwned)
        {
            ArrayPool<byte>.Shared.Return(remainingBuffer.ToArray());
        }
    }
}

关键细节说明

  • 缓冲区复用:使用ArrayPool<byte>管理字节数组,避免频繁创建销毁数组带来的性能损耗。
  • 根对象识别:通过CurrentDepth判断是否回到根级别,防止将嵌套对象误判为完整的根对象。
  • 续读支持:设置isFinalBlock: false,告知阅读器后续还有数据,避免因不完整数据抛出异常;未处理的字节会保留到下一轮解析。
  • 元素克隆:JsonElement.Clone()确保返回的元素在JsonDocument释放后仍可正常使用,因为JsonDocument是一次性对象。

为什么不能用JsonDocument.ParseAsync(Stream)

该方法会强制读取整个流直到末尾才开始解析,完全不适合长期保持开放、间歇性接收数据的网络流,无法满足“收到一个完整JSON就立即返回”的需求。

内容的提问来源于stack exchange,提问作者ygoe

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.07 15:25:10