Skip to content

Dart Stream:异步数据序列 / 生成器 / 背压 ​

回到总览:Dart 并发总览:Future / event loop / isolate
建议前置:先理解 Dart 事件循环与 Future(Dart EventLoop、Microtask Queue 与 Event Queue 机制、Dart Future、async/await 拆解与单线程并发)

一句话定义 ​

Dart 的异步是 Future + Stream 两条腿:Future 表达"将来会有一个值",Stream 表达"将来会有一串值(异步数据序列)"。Stream 支持单订阅 / 广播两种模式、listen 消费、pause/resume 手动背压,以及 async* / yield 生成器。

代码索引 ​

对应 Lab:dart-stream-async-generator

主题Lab 说明源码
Stream / StreamController / async* / 背压async* 生成器、手动推送、pause/resume 背压stream_async_generator_demo.dart

为什么需要 ​

很多工程师会用 Future / await,但遇到数据"来好几次"的场景就含糊:

  • Future 和 Stream 到底什么关系?什么时候必须用 Stream?
    • 一句话答:Future 表示单个异步值的终态回调,Stream 表示随时间产生的连续异步序列;面对 WebSocket、按钮点击、传感器等多事件场景必须用 Stream。(详见下文 §1)
  • 单订阅流和广播流差在哪?为什么广播流更适合多页面监听?
    • 一句话答:单订阅流只允许 listen 一次且默认缓存事件,广播流(Broadcast)允许多次订阅且事件随产随发,适合多 UI 组件共同监听广播。(详见下文 §2)
  • Stream.listen 的 pause / resume 是干什么的?
    • 一句话答:它们是 Dart Stream 底层自带的背压(Backpressure)流量控制开关,可以暂缓或恢复上游事件的派发。(详见下文 §3)
  • async* / yield 和普通 async / return 有什么本质区别?
    • 一句话答:async 函数返回单个 Future(用 return),而 async* 是异步生成器,使用 yield 逐步向外发射 Stream 数据元素。(详见下文 §4)
  • Stream 上的事件怎么排队?会卡住主 isolate 吗?
    • 一句话答:Stream 事件是以微任务(Microtask)或 Event 形式排入主 isolate 的事件循环执行;事件处理本身若有 CPU 重计算仍会卡住主 isolate。(详见下文 §1、§5)

如果说不清,工程上就很容易出现:

  • 键盘输入、传感器、WebSocket 消息这些"多次回调"硬用 Future 包,回调越包越乱
  • 多页面监听同一个流,用单订阅流导致 Bad state: Stream has already been listened to
  • 大流量事件流不做 pause/resume,UI 被事件淹没
  • 把 async* 当 async 写,return 和 yield 混用编译不过

底层机制 ​

1. Future = 单值,Stream = 多值 ​

对应 Lab:dart-stream-async-generator

dart
Future<int> one() async => 42;              // 将来返回一个值
Stream<int> many() async* {                 // 将来产出一串值
  yield 1;
  yield 2;
  yield 3;
}

预期现象:Future 的"完成"只有一次;Stream 的"事件"可以有很多次(onData 会被多次调用),直到 onDone 结束。选择依据是"这个数据是一次性的,还是持续产生的一串"。

2. 单订阅流 vs 广播流 ​

  • 单订阅流(默认) :只能有一个 listener,第二个 listen 抛 Bad state。适合"生产者 → 单个消费者"(如 File 逐行读、HttpClient 响应体)。
  • 广播流(broadcast) :多个 listener 同时监听,每个事件分发给所有订阅者。适合"多页面 / 多组件共享"(如全局通知)。注意:广播流不缓存事件,晚订阅的收不到之前的事件(除非自己缓存)。
dart
final controller = StreamController<int>.broadcast();
controller.stream.listen((v) => print('A: $v'));
controller.stream.listen((v) => print('B: $v'));   // 广播流允许第二个监听
controller.add(1);   // A 和 B 都收到

预期现象:单订阅流像"专用水管"(只能接一个水龙头);广播流像"广播电台"(谁调频谁听,但来晚了听不到之前的节目)。

3. StreamController:手动控制流的"水龙头" ​

dart
final ctrl = StreamController<int>();
ctrl.add(1);            // 发事件
ctrl.addError(Exception());  // 发错误
ctrl.close();           // 结束(触发 onDone)
ctrl.stream.listen(
  onData: (v) => print(v),
  onError: (e) => print('err'),
  onDone: () => print('done'),
);

预期现象:StreamController 让你在"没有 async* 生成器"的地方(回调、事件源、外部系统)手动产生流;add 之前有没有 listener 都行(单订阅流会缓冲直到有 listener?不——单订阅流在 listen 前的 add 会被缓冲,广播流则直接丢)。

细节:单订阅 StreamController 在无 listener 时 add 的数据会缓冲等待第一个 listener;广播流则丢弃无订阅者的事件。这决定了"先启动生产还是先监听"的工程顺序。

4. pause / resume:手动背压 ​

dart
final sub = stream.listen(expensive);
sub.pause();    // 暂停接收(生产者侧感知 onPause)
sub.resume();   // 恢复接收(生产者侧感知 onResume)

预期现象:pause 让流"暂停输送",适合下游处理不过来时(如 UI 繁忙)节流;配合 StreamController 的 onPause / onResume 回调,生产者可以据此调整产生速率——这是 Dart 侧手动背压的手段(相比 Kotlin Flow 的 suspend 自动背压更"手动")。

5. async* / yield:异步生成器 ​

dart
Stream<int> countdown(int n) async* {
  for (var i = n; i > 0; i--) {
    await Future.delayed(Duration(seconds: 1));  // 可 await 挂起
    yield i;                                      // 产出值并"让出"控制权
  }
}

对照点:

  • async 函数返回 Future,用 return 交出一个值;
  • async* 函数返回 Stream,用 yield 逐个产出值(yield 会挂起函数直到消费者请求下一个值);
  • 消费侧用 await for:
dart
await for (final v in countdown(3)) {
  print(v);   // 3 / 2 / 1,每次等 1 秒
}

预期现象:await for 是"异步 for 循环",每来一个值处理一个,天然带背压语义(处理完才请求下一个)——和 Kotlin Flow 的 suspend 背压是同一思路。

6. Stream 事件怎么排队 ​

Stream 的事件不是"并行"的,而是在事件循环里排队投递:同步 add 的事件进入 microtask / event queue,由当前 isolate 单线程依次处理。所以 Stream 上做 CPU 重活照样卡主 isolate——它只解决"多次异步投递",不解决"并发执行"。

对照点:这和"包了 Future 还会卡 UI"是同一结论:异步 ≠ 并行;真正要搬走 CPU 重活还是得上 isolate。

比喻(水管 vs 快递) :Future 是"一次性快递"(送一次就完);Stream 是"持续供水的水管"——可以开关(pause/resume)、可以分叉(broadcast 广播)、上游按你的用水节奏供水(await for / 背压)。async* / yield 则是"水厂自动产水"(生成器),而 StreamController 是"手动拧水龙头"。

Android / Flutter / Web / Backend 对照 ​

维度Dart StreamKotlin FlowJS 事件 / ObservableAndroid LiveData
冷/热默认单订阅冷,broadcast 热冷流本体 + 热化Observable 可配热(持当前值)
背压pause/resume / await for 手动suspend 自动背压无内置无(主线程投递)
取消StreamSubscription.cancel结构化取消unsubscribe自动(lifecycle)
多订阅broadcastshareInmulticast多观察者

常见场景 ​

1. 持续事件源 ​

键盘输入、传感器数据、WebSocket 消息、文件逐行读取 → Stream。

2. 多页面共享事件 ​

全局通知(登录态变化、深链)→ 广播流,各页面各自 listen。

3. UI 对高频数据降采样 ​

pause/resume 在页面不可见时暂停订阅,可见时恢复,避免后台空转。

常见误配、事故后果与排障 ​

坑现象修法
多页面监听单订阅流Bad state: Stream has already been listened to用 broadcast() / StreamController.broadcast()
广播流不缓存导致丢事件晚订阅收不到"登录态已变"需要最新值用 BehaviorSubject 思路自行缓存,或改 ValueNotifier
listen 后忘记 cancel页面销毁后回调还在跑(泄漏)StreamSubscription 存起来,dispose 里 cancel
Stream 里做 CPU 重活照样卡 UI重活搬 isolate,Stream 只做投递
先 add 后 listen(广播流)事件丢失先订阅再生产,或用带缓存的设计

与相近概念对比 ​

概念本质区别
Future单值异步;Stream 是多值异步
Iterable同步、阻塞取数据;Stream 异步、事件驱动
async / return生成单个 Future;async* / yield 生成多个事件
Kotlin Flow语义最接近;Flow 自动 suspend 背压,Dart 靠 pause/resume 手动
Rx 的 Observable概念对应;Dart 用原生 Stream 更轻

对应实验 ​

Lab说明源码
dart-stream-async-generatorasync* 生成器逐值产出、StreamController 手动推送、pause/resume 手动背压stream_async_generator_demo.dart

复习检查题 ​

  1. Future 和 Stream 的本质区别?

    答:Future 表达"将来一个值",完成一次;Stream 表达"将来一串值",onData 可多次触发直到 onDone。持续产生的事件序列必须用 Stream。

  2. 单订阅流和广播流差在哪?

    答:单订阅流只允许一个 listener(第二个 listen 抛异常),适合生产者→单个消费者;广播流允许多个 listener 同时接收,但不缓存事件,晚订阅者收不到之前的事件。

  3. pause / resume 是干什么的?

    答:手动背压——暂停流的输送(生产者感知 onPause),恢复继续(onResume);下游处理不过来时节流,配合 StreamController 的暂停回调调整生产速率。

  4. async* / yield 和 async / return 的本质区别?

    答:async 返回 Future,return 交出一个值;async* 返回 Stream,yield 逐个产出值并挂起函数直到消费者请求下一个;消费侧用 await for 逐值处理。

  5. Stream 上做 CPU 重活会卡 UI 吗?

    答:会。Stream 事件仍在事件循环里单线程排队,只解决"多次异步投递"不解决"并发执行";CPU 重活必须搬 isolate,和"包 Future 也会卡"同理。

速记 ​

  • Future = 单值;Stream = 多值异步序列
  • 单订阅 = 专用水管;广播 = 电台(晚来不补)
  • pause/resume = 手动背压
  • async* / yield 产流,await for 消费
  • Stream 事件走事件循环,异步 ≠ 并行
  • 记得 StreamSubscription.cancel 防泄漏

站点构建时间:2026/8/24 23:43:17