Dart中gRPC与REST集成问题:REST控制器无法等待gRPC处理结果
问题背景
我有一套通过gRPC通信的服务器与多个客户端(单元),当前测试单服务器单客户端场景。服务器同时接收用户应用的REST调用,整体流程为:用户应用(REST)<->服务器(gRPC)<->单元(执行额外REST调用后返回结果)。我尽量让示例代码结构与真实项目保持一致,但在REST API控制器与gRPC的连接环节遇到了问题。
当服务器收到REST调用时,我通过StreamController向客户端单元发送gRPC请求,单元会执行处理逻辑并发起额外REST调用后返回结果,但存在两个核心问题:
- 如何在REST控制器的
_setPreset函数中,读取单独文件的unitBus循环里处理的gRPC返回结果? - 如何让REST控制器等待整个gRPC调用完成后,再向用户应用返回REST响应?
日志信息
服务器端日志
unit bus processing await unitSubscribeReq: { unitId: IUZTR configId: asdf config: somethingWithConfiguration } return rest unit bus processing await
客户端日志
Got message unitSubscribeRes: { data: ok } Got message init: { config: message } VideoCore is offline else
从日志可见,服务器端的gRPC循环打印会延迟到下一次REST调用时才出现,这是异步处理时序错误导致的,大概率和async/await的使用不当有关。之前尝试过用额外的StreamController解决,但无法正确关联整个流程。我对REST技术较为熟悉,但gRPC是全新领域,希望得到解决方向。
解决方向与实现建议
核心思路:绑定REST请求与gRPC调用的生命周期
当前用StreamController传递结果的方式不适合单次请求-响应的场景,会导致结果被后续请求消费,且无法让REST请求等待gRPC调用完成。推荐用Completer+Future的方式为每个请求绑定独立的结果通道。
1. 用请求ID关联REST请求与gRPC响应
- 在服务器端维护一个全局Map,存储请求ID与对应的
Completer(用于生成等待结果的Future)。 - 每个REST请求触发时,生成唯一请求ID,创建
Completer并存入Map,然后发送带请求ID的gRPC请求。 - gRPC响应返回后,根据请求ID找到对应的
Completer,调用complete()传递结果,再从Map中移除该条目。 - 在REST控制器的
_setPreset函数中,await该Completer的Future,等待gRPC结果返回后再生成REST响应。
2. 修正gRPC循环的异步逻辑
检查unitBus中的循环处理代码,确保所有异步操作(比如gRPC调用、IO操作)都正确使用await,避免阻塞事件循环导致结果延迟返回。
3. 选择合适的gRPC调用模式
如果当前场景是单次请求-单次响应,优先使用gRPC的一元调用(UnaryCall),而非流式调用。一元调用本身就是基于Future的异步模型,能直接和REST请求的异步流程适配,无需自行维护StreamController。
关键代码调整示例
- 服务器端维护请求映射:
import 'package:uuid/uuid.dart'; final Map<String, Completer<UnitSubscribeRes>> _pendingRequests = {}; final _uuid = const Uuid();
- REST控制器的
_setPreset函数:
Future<Response> _setPreset(Request request) async { // 生成唯一请求ID final requestId = _uuid.v4(); final completer = Completer<UnitSubscribeRes>(); _pendingRequests[requestId] = completer; // 发送带请求ID的gRPC请求到单元 _unitBus.add(UnitSubscribeReq( unitId: 'IUZTR', configId: 'asdf', config: 'somethingWithConfiguration', requestId: requestId, // 新增请求ID字段,需要在proto中定义 )); // 等待gRPC响应结果 final grpcResponse = await completer.future; // 返回REST响应 return Response.ok(grpcResponse.data); }
- unitBus的处理循环:
void runUnitBus() async { await for (final req in _unitBus.stream) { print('unit bus processing await'); print('unitSubscribeReq: $req'); // 发起gRPC调用到单元客户端 final grpcRes = await _grpcClient.subscribe(req); // 根据请求ID找到对应的Completer并完成Future final completer = _pendingRequests.remove(req.requestId); if (completer != null) { completer.complete(grpcRes); print('gRPC result returned for request ${req.requestId}'); } } }
注意事项
- 需要在gRPC的proto文件中为
UnitSubscribeReq添加requestId字段,确保请求与响应能正确关联。 - 要处理超时场景:可以为
completer.future添加超时逻辑,避免REST请求无限等待。 - 如果是多客户端场景,需要确保
_pendingRequests的线程安全(Dart单线程模型下无需额外锁,但如果是Isolate场景需要注意)。
内容的提问来源于stack exchange,提问作者Liniik

