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

如何在Dart中按有序方式合并List<Stream<T>>为Stream<T>?

按轮询批次合并流的实现方案

你的需求是把多个流按**“先取所有流的第1个元素,再取所有流的第2个元素”**的顺序合并,这种逻辑既不是Dart内置Stream的merge(无序输出)、concat(逐个流全量输出),也不是rxdart里的ZipStream(必须所有流都有对应位置元素才输出),所以确实没有现成的内置或rxdart操作符直接实现,需要自定义逻辑。

核心实现思路

维护每个流的迭代器,循环轮询所有未完成的流,依次取出它们的下一个元素:

  • 给每个流创建StreamIterator,用于逐个获取元素;
  • 一轮轮遍历所有迭代器,能取出元素就输出,流完成后就标记为不再参与后续轮询;
  • 直到所有流都处理完,结束结果流。

完整代码实现

import 'dart:async';

Stream<T> mergeInBatches<T>(List<Stream<T>> streams) async* {
  // 初始化所有流的迭代器,移动到第一个元素位置
  final iterators = await Future.wait(
    streams.map((stream) async {
      final iterator = StreamIterator(stream);
      await iterator.moveNext();
      return iterator;
    }),
  );

  bool hasActiveStreams;
  do {
    hasActiveStreams = false;
    // 轮询每个流的迭代器
    for (var i = 0; i < iterators.length; i++) {
      final iterator = iterators[i];
      if (iterator != null && iterator.current != null) {
        hasActiveStreams = true;
        // 输出当前元素
        yield iterator.current;
        // 尝试获取下一个元素
        final hasNext = await iterator.moveNext();
        if (!hasNext) {
          // 流已结束,标记为null不再处理
          iterators[i] = null;
        }
      }
    }
  } while (hasActiveStreams);

  // 清理所有迭代器
  for (final iterator in iterators) {
    iterator?.cancel();
  }
}

// 测试示例
void main() async {
  final a = Stream.fromIterable([1, 2]);
  final b = Stream.fromIterable([3, 4]);
  final c = Stream.fromIterable([5, 6]);
  
  final streams = [a, b, c];
  
  await for (final num in mergeInBatches(streams)) {
    print(num); // 输出顺序:1 → 3 → 5 → 2 → 4 → 6
  }
}

适配流大小不固定的场景

如果某个流的元素数量和其他流不一致,比如修改测试代码:

final a = Stream.fromIterable([1]); // 仅1个元素
final b = Stream.fromIterable([3, 4]);
final c = Stream.fromIterable([5, 6, 7]); // 3个元素

输出顺序会是1 → 3 → 5 → 4 → 6 → 7,完全符合预期——a流结束后自动退出轮询,剩下的流继续按批次输出后续元素。

关于rxdart的说明

rxdart提供的concat、merge、zip、combineLatest等操作符,都无法直接实现这种“按轮询批次取元素”的逻辑,所以自定义实现是当前最合适的方案。

内容的提问来源于stack exchange,提问作者Eray Erdin

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.12 05:12:30