Appearance
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 Stream | Kotlin Flow | JS 事件 / Observable | Android LiveData |
|---|---|---|---|---|
| 冷/热 | 默认单订阅冷,broadcast 热 | 冷流本体 + 热化 | Observable 可配 | 热(持当前值) |
| 背压 | pause/resume / await for 手动 | suspend 自动背压 | 无内置 | 无(主线程投递) |
| 取消 | StreamSubscription.cancel | 结构化取消 | unsubscribe | 自动(lifecycle) |
| 多订阅 | broadcast | shareIn | multicast | 多观察者 |
常见场景
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-generator | async* 生成器逐值产出、StreamController 手动推送、pause/resume 手动背压 | stream_async_generator_demo.dart |
复习检查题
Future 和 Stream 的本质区别?
答:Future 表达"将来一个值",完成一次;Stream 表达"将来一串值",
onData可多次触发直到onDone。持续产生的事件序列必须用 Stream。单订阅流和广播流差在哪?
答:单订阅流只允许一个 listener(第二个 listen 抛异常),适合生产者→单个消费者;广播流允许多个 listener 同时接收,但不缓存事件,晚订阅者收不到之前的事件。
pause / resume是干什么的?答:手动背压——暂停流的输送(生产者感知
onPause),恢复继续(onResume);下游处理不过来时节流,配合StreamController的暂停回调调整生产速率。async*/yield和async/return的本质区别?答:
async返回Future,return交出一个值;async*返回Stream,yield逐个产出值并挂起函数直到消费者请求下一个;消费侧用await for逐值处理。Stream 上做 CPU 重活会卡 UI 吗?
答:会。Stream 事件仍在事件循环里单线程排队,只解决"多次异步投递"不解决"并发执行";CPU 重活必须搬 isolate,和"包 Future 也会卡"同理。
速记
- Future = 单值;Stream = 多值异步序列
- 单订阅 = 专用水管;广播 = 电台(晚来不补)
pause/resume= 手动背压async*/yield产流,await for消费- Stream 事件走事件循环,异步 ≠ 并行
- 记得
StreamSubscription.cancel防泄漏