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

如何让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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.09 02:05:25