文章

Dart Stream 流式编程深度解析:数据管道、背压与响应式开发的正确姿势

Stream 是 Dart 响应式编程的管道:讲清单订阅/广播流、listen、await for、背压(Backpressure)与 StreamBuilder。 掌握后能答出「Stream 和 Future 的区别及背压怎么处理」。

Dart Stream 流式编程深度解析:数据管道、背压与响应式开发的正确姿势

一句话概括

Stream 是 Dart 的异步事件序列——如果说 Future 是”一次性外卖”,Stream 就是”持续出餐的传送带”。它解决了 WebSocket 推送、文件流读取、状态变化广播等场景中”数据分批次到达”的问题。

核心知识点

1. Stream 的两种模式:单播 vs 广播

1
2
3
4
5
6
7
8
9
10
// 单订阅(single-subscription)—— 只能有一个 listener
final single = Stream.fromIterable([1, 2, 3]);
single.listen(print); // ✅
// single.listen(print); // ❌ 运行时错误

// 广播(broadcast)—— 多对多,类似 EventEmitter
final broadcast = StreamController<int>.broadcast();
broadcast.stream.listen(print); // listener A
broadcast.stream.listen(print); // listener B ✅
broadcast.add(42); // A 和 B 都收到

面试要点:默认是单播,需要广播用 .asBroadcastStream() 或 StreamController.broadcast()。单播的 Stream 是按需生产(懒加载),广播是推模型。

2. StreamController:手动控制的事件源

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
final controller = StreamController<String>();

// 从外部往里塞
controller.add('事件1');
controller.add('事件2');

// 错误处理
controller.addError(Exception('出错了'));

// 监听
controller.stream.listen(
  (data) => print('收到: $data'),
  onError: (e) => print('错误: $e'),
  onDone: () => print('流关闭'),
);

// 完成后关闭
await controller.close();

StreamController 是 Stream 的”写端”,StreamController.stream 是”读端”。这和 Completer + Future 的设计模式一致。

3. 用 await for 消费 Stream

1
2
3
4
5
6
7
Future<void> processAll(Stream<int> numbers) async {
  await for (final n in numbers) {
    print('处理: $n');
    if (n > 100) break; // 提前终止
  }
  print('流结束或手动终止');
}

await for 等价于 .listen() 后的异步回调,但写法更像同步。退出循环后 Stream 的 subscription 会自动 cancel。

4. 转换操作符:map、where、transform

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
Stream<int> ticker() => Stream.periodic(
  Duration(seconds: 1), (i) => i + 1
);

// 链式转换
ticker()
  .where((n) => n.isEven)           // 只取偶数
  .map((n) => '第 $n 秒')            // 类型转换 Stream<int> → Stream<String>
  .take(5)                           // 只取前5个
  .listen(print);                    // 第2秒, 第4秒, ...

// StreamTransformer:复杂转换
final parser = StreamTransformer<int, String>.fromHandlers(
  handleData: (data, sink) {
    if (data > 0) sink.add('有效: $data');
    else sink.addError('无效数据');
  }
);

常见操作符:map、where、take、skip、distinct、handleError、timeout。这些都是惰性的,只有被订阅时才实际执行。

5. 背压控制:pausable 订阅

1
2
3
4
5
6
7
8
9
// 消费者处理不过来时暂停生产
final sub = fastProducer.listen((data) async {
  await heavyProcess(data); // 可能很慢
});
sub.pause();  // 暂停接收
sub.resume(); // 恢复接收

// StreamTransformer 自带背压——sink 会等到 add 完成后才继续
final pipeline = slowProducer.transform(Processor());

Stream 默认是推模型——生产者不管消费者能不能处理。背压通过 pause/resume 或 StreamTransformer 的内部缓冲机制来控制流量。


其实你每天都在用

  • StreamBuilder:Flutter 里监听 Stream 更新 UI 的官方 Widget,内部自动管理订阅
  • TextField.onChanged:返回的是一个 Stream 吗?不是,是 callback。但 RxDart 把它转成了 Observable
  • BLoC 模式:核心就是 StreamController —— Events in, States out
  • Firebase onSnapshot():返回 Stream,每次数据变更自动推送
  • File.openRead():返回 Stream<List<int>>,逐块读取大文件,不会一次性撑爆内存

常见误解(FAQ)

❌ 误区:「Stream 和 Future 只是『一个值』和『多个值』的区别」

不只是数量。Future 是即时启动的(一旦创建就开始执行),单播 Stream 是懒启动的(只有被订阅时才生产数据)。这个执行时机的差异很关键。

❌ 误区:「await for 会阻塞 UI 线程」

跟 await 一样,不会阻塞。await for 在每次迭代之间让出控制权给事件循环。UI 照样 60fps。

❌ 误区:「Stream 多订阅要小心,默认是广播」

反过来——默认是单播。忘记这点的结果是运行时 “Stream already listened” 错误。需要广播时显式声明。

❌ 误区:「RxDart 可以完全替代 Stream」

RxDart 是 Stream API 的扩展(Observable 实现 Stream 接口),不是替代。几乎所有 RxDart 操作符都能回退到纯 Stream + Transformer 实现。


一句话总结

Stream 是数据管道,不是数据容器——它描述”数据将如何流动”,而不是”数据现在在哪里”。掌握 StreamController(写端)+ Stream(读端)+ Transformer(处理)= 就能驾驭任何流式场景。

本文由作者按照 CC BY 4.0 进行授权