如何用.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
相关产品推荐
相关产品推荐

