Python客户端向C# gRPC服务请求流上传报错排查求助
我正在尝试用gRPC实现Python客户端到C#服务端的请求流上传功能,目前第一个参数消息已成功发送,但服务端执行requestStream.MoveNext()时立即触发错误,已尝试多种方案未解决,恳请帮忙排查。
Proto定义
syntax = "proto3"; service PerceiveAPIDataService { rpc UploadResource (stream UploadResourceRequest) returns (UploadResourceResponse); } message UploadResourceRequest { oneof request_data { ResourceChunk resource_chunk = 1; UploadResourceParameters parameters = 2; } } message ResourceChunk { bytes content = 1; } message UploadResourceParameters { string path = 1; }
C#服务端实现代码
public override async Task<UploadResourceResponse> UploadResource(IAsyncStreamReader<UploadResourceRequest> requestStream, ServerCallContext context) { if (!await requestStream.MoveNext()) { throw new RpcException(new Status(StatusCode.FailedPrecondition, "No upload parameters found.")); } var initialMessage = requestStream.Current; if (initialMessage.RequestDataCase != UploadResourceRequest.RequestDataOneofCase.Parameters) { throw new RpcException(new Status(StatusCode.FailedPrecondition, "First message must contain upload parameters.")); } var path = initialMessage.Parameters.Path; if (string.IsNullOrWhiteSpace(path)) { throw new RpcException(new Status(StatusCode.InvalidArgument, "Upload path is required.")); } using (var ms = new MemoryStream()) { while (await requestStream.MoveNext()) { var chunk = requestStream.Current.ResourceChunk; if (chunk == null) { continue; // Skip any messages that are not resource chunks } await ms.WriteAsync(chunk.Content.ToByteArray().AsMemory(0, chunk.Content.Length)); } ms.Seek(0, SeekOrigin.Begin); // Reset memory stream position to the beginning for reading during upload var uploadResult = await _dataService.UploadResourceAsync(path, ms); return new UploadResourceResponse { Succeeded = uploadResult.IsSuccessful }; } }
Python客户端代码
def generate_request(self, data: bytearray, next_cloud_path: str) -> Generator: first_req = perceive_api_data_service_pb2.UploadResourceRequest( parameters=perceive_api_data_service_pb2.UploadResourceParameters(path=next_cloud_path) ) yield first_req print("Sent initial request with path:", next_cloud_path) chunk_size = 2048 total_chunks = (len(data) + chunk_size - 1) // chunk_size # Ceiling division to get total number of chunks print(f"Data size: {len(data)} bytes, chunk size: {chunk_size} bytes, total chunks: {total_chunks}") for i in range(0, len(data), chunk_size): chunk = data[i:i+chunk_size] yield perceive_api_data_service_pb2.UploadResourceRequest( resource_chunk=perceive_api_data_service_pb2.ResourceChunk(content=chunk) ) print(f"Sent chunk {((i // chunk_size) + 1)} of {total_chunks}") async def upload_file(self, data: bytearray, next_cloud_path: str) -> bool: async with grpc.aio.insecure_channel("localhost:5228") as channel: stub = perceive_api_data_service_pb2_grpc.PerceiveAPIDataServiceStub(channel) request_iterator = self.generate_request(data, next_cloud_path) response = await stub.UploadResource(request_iterator) return response.succeeded
服务端报错信息
Grpc.AspNetCore.Server.ServerCallHandler: Error: Error when executing service method 'UploadResource'.
System.IO.IOException: The client reset the request stream.
at System.IO.Pipelines.Pipe.GetReadResult(ReadResult& result)
at System.IO.Pipelines.Pipe.GetReadAsyncResult()
客户端报错信息
File "C:\Users\user_name\AppData\Local\Programs\Python\Python311\Lib\site-packages\grpc\aio_call.py", line 690, in _conduct_rpc
serialized_response = await self._cython_call.stream_unary(
^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^ File "src\python\grpcio\grpc_cython_cygrpc/aio/call.pyx.pxi", line 458, in stream_unary File "src\python\grpcio\grpc_cython_cygrpc/aio/callback_common.pyx.pxi", line 166, in _receive_initial_metadata File "src\python\grpcio\grpc_cython_cygrpc/aio/callback_common.pyx.pxi", line 99, in execute_batch asyncio.exceptions.CancelledError
排查方向及解决方案
同步生成器不兼容异步gRPC
Python的grpc.aio要求请求迭代器为异步生成器,当前同步Generator会导致流处理异常。修改生成器为异步类型:async def generate_request(self, data: bytearray, next_cloud_path: str) -> AsyncGenerator: first_req = perceive_api_data_service_pb2.UploadResourceRequest( parameters=perceive_api_data_service_pb2.UploadResourceParameters(path=next_cloud_path) ) yield first_req print("Sent initial request with path:", next_cloud_path) chunk_size = 2048 total_chunks = (len(data) + chunk_size - 1) // chunk_size print(f"Data size: {len(data)} bytes, chunk size: {chunk_size} bytes, total chunks: {total_chunks}") for i in range(0, len(data), chunk_size): chunk = data[i:i+chunk_size] yield perceive_api_data_service_pb2.UploadResourceRequest( resource_chunk=perceive_api_data_service_pb2.ResourceChunk(content=chunk) ) print(f"Sent chunk {((i // chunk_size) + 1)} of {total_chunks}")服务端未处理取消信号
在读取流时传入取消令牌,避免模糊IO异常:while (await requestStream.MoveNext(context.CancellationToken)) { context.CancellationToken.ThrowIfCancellationRequested(); var chunk = requestStream.Current.ResourceChunk; if (chunk == null) { continue; } await ms.WriteAsync(chunk.Content.ToByteArray().AsMemory(0, chunk.Content.Length)); }版本兼容性问题
检查Python的grpcio/grpcio-tools与C#的Grpc.AspNetCore版本是否匹配,建议两边升级至最新稳定版,避免跨语言版本差异导致的流逻辑不兼容。客户端通道配置优化
调整通道参数,避免因超时或消息大小限制导致流中断:async with grpc.aio.insecure_channel("localhost:5228", options=[ ("grpc.max_send_message_length", -1), ("grpc.max_receive_message_length", -1), ("grpc.keepalive_time_ms", 10000) ]) as channel: # 原有调用逻辑
内容的提问来源于stack exchange,提问作者Luuk Wuijster

