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

GRPC双向流问题:请求流意外关闭异常排查求助

GRPC双向流异常问题解决方案

问题根源

  1. 服务器使用Task.WhenAny(ProcessRequests(), PublishResponses()),只要其中一个任务完成(比如客户端关闭请求流后ProcessRequests完成),HandleUpdate方法就会立即返回,GRPC框架会终结整个双向流调用,导致后续读取请求流时抛出Can't read messages after the request is complete.异常。
  2. 服务器未正确维护流的生命周期,未等待读写任务同步执行,导致流被提前关闭。
  3. 客户端未添加流状态检查,可能在服务器关闭流后仍尝试写入,加剧异常触发概率。

修复方案

服务器端代码修改

将Task.WhenAny替换为Task.WhenAll,确保读写任务同步运行,同时细化异常处理,记录关键状态:

private async Task HandleUpdate(IAsyncStreamReader<UpdateRequest> requestStream,
    IServerStreamWriter<UpdateResponse> responseStream, CancellationToken cancellationToken)
{
    async Task ProcessRequests()
    {
        try
        {
            await foreach (var request in requestStream.ReadAllAsync(cancellationToken))
            {
                _logger.LogInformation($"Update: {request}");
            }
            _logger.LogInformation("客户端已关闭请求流");
        }
        catch (RpcException ex) when (ex.StatusCode == StatusCode.Cancelled)
        {
            _logger.LogInformation("请求流被取消");
        }
        catch (Exception e)
        {
            _logger.LogError(e, "处理请求流时发生异常");
        }
    }

    async Task PublishResponses()
    {
        try
        {
            while (!cancellationToken.IsCancellationRequested)
            {
                // 替换为实际的响应发送逻辑
                // await responseStream.WriteAsync(new UpdateResponse());
                await Task.Delay(2000, cancellationToken);
            }
        }
        catch (TaskCanceledException)
        {
            _logger.LogInformation("响应发布任务被取消");
        }
        catch (Exception e)
        {
            _logger.LogError(e, "发布响应时发生异常");
        }
    }

    try
    {
        // 等待读写任务全部完成或被取消
        await Task.WhenAll(ProcessRequests(), PublishResponses());
    }
    catch (Exception ex)
    {
        _logger.LogError(ex, "HandleUpdate发生未捕获异常");
    }
}

客户端代码修改

添加流状态检查,细化异常捕获逻辑,确保写入任务完成后再释放资源:

using (var call = client.StreamTableUpdate())
{
    try
    {
        var writeTask = Task.Run(async () =>
        {
            for (int i = 0; i < 1000; i++)
            {
                if (call.RequestStream.IsCompleted)
                {
                    _logger.LogWarning("请求流已关闭,停止写入");
                    break;
                }
                try
                {
                    var request = new UpdateRequest();
                    await call.RequestStream.WriteAsync(request);
                    // 可选:添加短暂延迟避免过快发送
                    // await Task.Delay(10);
                }
                catch (RpcException ex) when (ex.StatusCode == StatusCode.Cancelled)
                {
                    _logger.LogInformation("写入请求流时被取消");
                    break;
                }
                catch (Exception e)
                {
                    _logger.LogError(e, $"写入第{i}个请求时发生异常");
                    break;
                }
            }

            try
            {
                await call.RequestStream.CompleteAsync();
                _logger.LogInformation("客户端已完成请求流写入");
            }
            catch (Exception e)
            {
                _logger.LogError(e, "关闭请求流时发生异常");
            }
        });

        try
        {
            await foreach (var response in call.ResponseStream.ReadAllAsync())
            {
                // 替换为实际的响应处理逻辑
                _logger.LogInformation($"收到响应: {response}");
            }
            _logger.LogInformation("响应流已关闭");
        }
        catch (RpcException ex) when (ex.StatusCode == StatusCode.Cancelled)
        {
            _logger.LogInformation("读取响应流时被取消");
        }
        catch (Exception e)
        {
            _logger.LogError(e, "读取响应流时发生异常");
        }

        // 等待写入任务完全完成
        await writeTask;
    }
    catch (Exception e)
    {
        _logger.LogError(e, "处理双向流时发生异常");
    }
}

关键说明

  • 服务器通过Task.WhenAll确保读写任务同步执行,直到客户端关闭请求流或触发取消信号,避免流被提前终结。
  • 客户端添加call.RequestStream.IsCompleted检查,避免在流关闭后继续写入,同时针对GRPC特定的RpcException(如取消状态码)做针对性处理。
  • 移除无意义的异常吞入,改为日志记录,便于排查问题。

内容的提问来源于stack exchange,提问作者Wonkyung Wayne Kim

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.24 06:34:57