Azure Function中自定义序列化器的CosmosClient查询CosmosDB异常问题
问题:Azure Function中自定义CosmosSerializer导致查询无结果并抛出流未关闭异常
场景说明
在Azure Function中使用带自定义序列化器的CosmosClient查询CosmosDB时,接口返回200 OK但无数据返回,通过Data Explorer验证查询语句本身正常。调试发现自定义序列化器执行后,触发了SDK内部CosmosJsonSerializerWrapper的异常:InvalidOperationException: Json Serializer left an open stream.
自定义序列化器代码
public class CosmosDbSerializerService : CosmosSerializer { private readonly JsonSerializer jsonSerializer; public CosmosDbSerializerService() { jsonSerializer = JsonSerializer .Create(JsonSerializerFactory.GetJsonSerializerSettings()); } public override T FromStream<T>(Stream stream) { using (StreamReader streamReader = new(stream)) using (JsonTextReader jsonTextReader = new(streamReader)) { return jsonSerializer .Deserialize<T>(jsonTextReader); } } public override Stream ToStream<T>(T input) { using (StringWriter stringWriter = new()) using (JsonTextWriter jsonTextWriter = new(stringWriter)) { jsonSerializer.Serialize(jsonTextWriter, input, input.GetType()); return new MemoryStream(Encoding.UTF8.GetBytes(stringWriter.ToString())); } } }
CosmosClient依赖注入配置
builder.Register(_ => { CosmosClientBuilder cosmosClientBuilder = new CosmosClientBuilder(productionDatabaseSettings.ConnectionString) .WithCustomSerializer(new CosmosDbSerializerService()) .WithHttpClientFactory(() => { HttpMessageHandler httpMessageHandler = new HttpClientHandler() { ServerCertificateCustomValidationCallback = HttpClientHandler.DangerousAcceptAnyServerCertificateValidator }; return new HttpClient(httpMessageHandler); }) .WithConnectionModeGateway(); return cosmosClientBuilder.Build(); });
问题分析
1. 流未正确处理的异常原因
自定义FromStream方法中,StreamReader默认会在using块结束时关闭底层流,但SDK内部的CosmosJsonSerializerWrapper会校验流是否处于不可读状态。由于StreamReader的关闭操作未被SDK正确识别,导致检测到流仍可读取,从而抛出异常。
2. 为什么会用到CosmosJsonSerializerWrapper
这是Azure Cosmos DB .NET SDK的内部机制:所有自定义CosmosSerializer都会被自动包装到CosmosJsonSerializerWrapper中,用于添加额外的校验逻辑(如流状态检查),自定义序列化器本身已经生效,只是被SDK做了一层包装。
解决方案
修复FromStream方法
修改FromStream,确保处理完流后将其置为不可读状态:
public override T FromStream<T>(Stream stream) { // 创建StreamReader时设置leaveOpen: true,避免自动关闭底层流 using (StreamReader streamReader = new StreamReader(stream, Encoding.UTF8, true, 1024, leaveOpen: true)) using (JsonTextReader jsonTextReader = new JsonTextReader(streamReader)) { var result = jsonSerializer.Deserialize<T>(jsonTextReader); // 将流位置移到末尾,让CanRead返回false stream.Position = stream.Length; return result; } }
优化ToStream方法(可选)
原方法使用StringWriter存在额外开销,可直接写入MemoryStream并避免意外关闭:
public override Stream ToStream<T>(T input) { MemoryStream memoryStream = new MemoryStream(); using (StreamWriter streamWriter = new StreamWriter(memoryStream, Encoding.UTF8, 1024, leaveOpen: true)) using (JsonTextWriter jsonTextWriter = new JsonTextWriter(streamWriter)) { jsonSerializer.Serialize(jsonTextWriter, input, input.GetType()); jsonTextWriter.Flush(); streamWriter.Flush(); memoryStream.Position = 0; return memoryStream; } }
验证
修改后重新运行函数,查询将正常返回结果,不会再触发流未关闭的异常。
内容的提问来源于stack exchange,提问作者Attilio Gelosa
相关产品推荐
相关产品推荐

