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

Python客户端向C# gRPC服务请求流上传报错排查求助

gRPC Python客户端向C#服务端请求流上传失败问题

我正在尝试用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


排查方向及解决方案

  1. 同步生成器不兼容异步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}")
    
  2. 服务端未处理取消信号
    在读取流时传入取消令牌,避免模糊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));
    }
    
  3. 版本兼容性问题
    检查Python的grpcio/grpcio-tools与C#的Grpc.AspNetCore版本是否匹配,建议两边升级至最新稳定版,避免跨语言版本差异导致的流逻辑不兼容。

  4. 客户端通道配置优化
    调整通道参数,避免因超时或消息大小限制导致流中断:

    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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.25 19:28:17