如何让Dart StreamController并行处理多个请求?
问题:StreamController结合asyncMap无法并行处理请求,结果按输入顺序返回而非完成顺序
我希望通过StreamController并行处理多个请求,但实际运行结果和预期不符。
代码实现
main.dart
final controller = StreamController<Req>.broadcast(); final resStream = controller.stream.asyncMap<Res>((request) { switch (request.method) { case "get": return Future.delayed( const Duration(seconds: 10), () => const Res('response after 10s'), ); case "post": default: return Future.delayed( const Duration(seconds: 5), () => const Res('response after 5s'), ); } }); resStream.listen((event) { print('+++ response from stream: ${event.value}'); }); controller.add(const Req('get')); controller.add(const Req('post'));
req.dart
class Req { const Req(this.method); final String method; }
res.dart
class Res { const Res(this.value); final String value; }
实际输出
+++ response from stream: response after 10s +++ response from stream: response after 5s
预期输出
+++ response from stream: response after 5s +++ response from stream: response after 10s
解答
问题原因
asyncMap的核心特性是串行处理流事件:它会等待当前事件对应的Future执行完成后,才会去处理流中的下一个事件。
你先添加了get请求,它的Future需要10秒才能完成,此时post请求的事件会被暂存在流中,直到get的Future执行完毕,才会开始处理post的Future。因此最终输出顺序和请求添加的顺序完全一致,而非任务完成的顺序。
解决方案:使用asyncExpand实现并行处理
要实现并行处理,让先完成的请求先输出结果,可以用asyncExpand替代asyncMap。asyncExpand会把每个输入事件转换成一个独立的Stream,然后将所有这些Stream合并成一个输出流,哪个子Stream先产生结果就先输出。
修改resStream的实现:
final resStream = controller.stream.asyncExpand<Res>((request) { Future<Res> future; switch (request.method) { case "get": future = Future.delayed( const Duration(seconds: 10), () => const Res('response after 10s'), ); break; case "post": default: future = Future.delayed( const Duration(seconds: 5), () => const Res('response after 5s'), ); break; } // 将每个Future转换成单元素Stream return Stream.fromFuture(future); });
修改后的输出
运行修改后的代码,输出会符合预期:
+++ response from stream: response after 5s +++ response from stream: response after 10s
额外说明
如果需要限制并发请求的数量(比如避免同时发起过多请求),可以借助package:async库中的StreamQueue或自定义限流逻辑,但对于基础的并行需求,asyncExpand已经足够满足。
内容的提问来源于stack exchange,提问作者Alex
相关产品推荐
相关产品推荐

