GRPC双向流问题:请求流意外关闭异常排查求助
GRPC双向流异常问题解决方案
问题根源
- 服务器使用
Task.WhenAny(ProcessRequests(), PublishResponses()),只要其中一个任务完成(比如客户端关闭请求流后ProcessRequests完成),HandleUpdate方法就会立即返回,GRPC框架会终结整个双向流调用,导致后续读取请求流时抛出Can't read messages after the request is complete.异常。 - 服务器未正确维护流的生命周期,未等待读写任务同步执行,导致流被提前关闭。
- 客户端未添加流状态检查,可能在服务器关闭流后仍尝试写入,加剧异常触发概率。
修复方案
服务器端代码修改
将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
相关产品推荐
相关产品推荐

