Flutter 基础体系 · 第 29/80 篇。示例基于当前稳定 Flutter 与 Dart 3 语言能力;Android、iOS、桌面和 Web 差异会明确说明。

Dart Stream:单订阅、广播、转换、背压边界和取消

Stream<T> 表示一组按时间到达的异步事件,事件类型为 T。与 Future<T> 只完成一次不同,Stream 可能依次产生多个数据事件,也可能产生错误事件,最后以一个 done 事件结束。

一个 Stream 的基本事件序列可以抽象为:

E=d1,d2,,dn,doneE = d_1, d_2, \ldots, d_n, \text{done}

其中 did_i 是数据事件。错误事件不一定替代 done;除非底层实现或订阅者主动取消,一个 Stream 通常可以在错误之后继续发送数据。

Stream 的关键难点不在于“如何收到一个值”,而在于以下状态和边界:

  • 一个 Stream 能否被多个订阅者订阅;
  • 新订阅者能否收到订阅前已经发生的事件;
  • 事件转换是否异步、是否保持顺序;
  • 订阅者处理变慢时,事件在哪里缓存;
  • 取消订阅后,底层资源何时真正释放;
  • 错误、暂停、关闭和取消之间如何相互影响。

一、Stream、StreamSubscription 与 StreamController

1. Stream 是事件来源的只读视图

Stream<T> 对外暴露的是事件序列。调用 listen 后,会得到一个 StreamSubscription<T>

final subscription = stream.listen(
  (value) {
    print('data: $value');
  },
  onError: (Object error, StackTrace stackTrace) {
    print('error: $error');
  },
  onDone: () {
    print('done');
  },
);

三类回调分别对应:

  1. 数据事件:onData
  2. 错误事件:onError
  3. done 事件:onDone

StreamSubscription 是订阅本身的控制对象,负责:

  • pause():暂时不向该订阅者交付数据;
  • resume():恢复交付;
  • cancel():取消订阅;
  • asFuture():把“done 或错误”转换为一个 Future。

因此,Stream 描述“有什么事件”,StreamSubscription 描述“某个订阅者如何接收这些事件”。

2. StreamController 是常见的事件生产端

自定义 Stream 时,通常使用 StreamController<T>

final controller = StreamController<int>();

controller.stream.listen(print);

controller.add(1);
controller.add(2);
await controller.close();

这里的数据流是:

controller.add(1)
        │
        ▼
controller.stream
        │
        ▼
StreamSubscription 的 onData(1)

controller.add 只是提交一个数据事件,并不等价于“订阅者已经处理完这个事件”。生产者和消费者之间仍然可能存在异步排队。

生产端还可以注册生命周期回调:

final controller = StreamController<int>(
  onListen: () {
    print('第一次建立订阅');
  },
  onPause: () {
    print('订阅被暂停');
  },
  onResume: () {
    print('订阅恢复');
  },
  onCancel: () {
    print('订阅被取消');
  },
);

这些回调适合连接外部资源,例如定时器、Socket、文件句柄或平台事件通道。但它们不是所有 Stream 都存在的统一能力;只有控制器或底层实现明确提供时,生产端才能据此暂停或释放资源。


二、单订阅 Stream:一次生命周期只能有一个订阅

1. 定义与行为

单订阅 Stream(single-subscription stream)最多允许建立一个有效订阅。第二次调用 listen 通常会抛出 StateError

Future<void> main() async {
  final stream = Stream.fromIterable([1, 2, 3]);

  await stream.listen((value) {
    print('first: $value');
  }).asFuture<void>();

  // 这里会抛出 StateError:
  // stream.listen((value) {
  //   print('second: $value');
  // });
}

关键点是:单订阅限制作用于这个 Stream 实例,而不是 Stream 类型本身。以下代码创建了两个独立的 Stream 实例,因此可以分别订阅:

final first = Stream.fromIterable([1, 2, 3]);
final second = Stream.fromIterable([1, 2, 3]);

first.listen((value) => print('A: $value'));
second.listen((value) => print('B: $value'));

2. 为什么默认倾向单订阅

单订阅适合有明确开始和结束边界的资源:

  • 读取一个文件;
  • 发起一次分页加载;
  • 消费一个 HTTP 响应体;
  • 解析一条 Socket 连接;
  • 执行一次异步计算流程。

这类资源通常有一个自然的所有权关系:

建立资源 → 产生事件 → 结束或取消 → 释放资源

如果允许多个订阅者随意加入,就必须额外定义:

  • 每个订阅者是否从头开始;
  • 数据是否复制给每个人;
  • 一个订阅者取消是否影响其他订阅者;
  • 底层资源何时关闭。

单订阅 Stream 通过限制订阅数量,避免了这些语义歧义。

3. 单订阅 Stream 是否能“重放”

单订阅并不意味着自动缓存历史事件。下面的代码中,事件只会发送给已经建立的那个订阅:

final controller = StreamController<int>();

controller.add(1);
controller.add(2);

controller.stream.listen(print);

通常不会打印 12。因为这些事件发生时没有订阅者,普通 StreamController 不会把它们当作历史记录保存下来。

如果业务需要“后来加入的订阅者也能看到历史值”,需要显式设计缓存或使用具有回放语义的其他抽象,不能从“单订阅”这个名称推导出回放能力。


三、广播 Stream:多个订阅者与实时事件

1. 广播的核心语义

广播 Stream(broadcast stream)允许多个订阅者同时监听:

final controller = StreamController<int>.broadcast();

controller.stream.listen((value) {
  print('A: $value');
});

controller.stream.listen((value) {
  print('B: $value');
});

controller.add(10);

事件 10 会分别交付给 A 和 B。每个订阅者都有自己的 StreamSubscription,可以独立暂停或取消。

广播 Stream 的典型来源包括:

  • 用户输入事件;
  • 传感器事件;
  • 网络连接状态变化;
  • Flutter 生命周期事件;
  • 全局状态或消息通知。

2. 广播不是历史缓存

广播 Stream 通常是“实时分发”而不是“事件回放”。如果事件产生时没有监听者,事件可能直接丢失:

final controller = StreamController<int>.broadcast();

controller.add(1); // 没有订阅者,事件通常被丢弃

controller.stream.listen((value) {
  print(value);
});

controller.add(2); // 订阅者可以收到

所以广播 Stream 更接近:

当前有哪些监听者,就向当前监听者发送事件

而不是:

先保存所有事件,未来每个监听者从头读取

这也是广播 Stream 不适合直接承载“必须可靠消费”的任务队列的原因。支付指令、数据库写入任务、文件上传任务等需要可靠处理时,应使用具有明确持久化、确认或重试语义的机制。

3. asBroadcastStream 的作用与限制

已有一个单订阅 Stream 时,可以调用 asBroadcastStream 创建广播视图:

final source = Stream.fromIterable([1, 2, 3]);
final broadcast = source.asBroadcastStream();

broadcast.listen((value) {
  print('A: $value');
});

broadcast.listen((value) {
  print('B: $value');
});

它的作用是让多个订阅者共享同一个底层订阅,而不是为每个监听者重新执行 source

可以通过回调控制底层订阅的暂停和取消:

final broadcast = source.asBroadcastStream(
  onListen: (subscription) {
    print('广播流第一次有监听者');
  },
  onCancel: (subscription) {
    print('广播流没有监听者了');
    // 可以选择 subscription.cancel(),但必须确认底层资源确实应被关闭。
  },
);

这里有一个重要风险:多个广播订阅者共享底层订阅。某个订阅者的取消,不等价于底层数据源应当立即取消;只有广播层不再需要底层资源时,才应关闭共享资源。

如果需要严格控制“有几个监听者、最后一个监听者取消时如何释放资源”,自己使用 StreamController.broadcast 并维护资源生命周期通常更容易验证。


四、事件顺序、同步与异步交付

1. Stream 保证什么顺序

对于同一个订阅者,正常情况下,数据事件按源 Stream 发出的顺序交付:

source: 1 → 2 → 3
listener: 1 → 2 → 3

但“按顺序交付”不表示回调中的异步操作会自动串行完成:

stream.listen((value) async {
  await saveToDatabase(value);
  print('saved: $value');
});

listen 不会等待这个 async 回调返回的 Future。可能出现:

收到 1
收到 2
save(2) 先完成
save(1) 后完成

如果必须按事件顺序等待异步操作,应使用 asyncMapawait for 等明确的转换方式。

2. sync 控制器与异步控制器

final asyncController = StreamController<int>();

final syncController = StreamController<int>.sync();

普通 StreamController 默认异步交付事件;sync: true 允许在 add 调用过程中同步调用订阅者回调。

同步控制器的风险是重入。比如订阅者在处理一个事件时,又调用了生产端的 add,可能形成嵌套回调:

add(1)
└─ onData(1)
   └─ add(2)
      └─ onData(2)

同步交付不是性能保证,也不是“更正确”的模式。它改变了调用栈和重入时序,只有在生产者和消费者都能接受这种时序时才应使用。


五、Stream 转换:从事件序列到新的事件序列

Stream 转换不是简单的“把一个值改成另一个值”。转换还必须定义:

  • 一个输入事件产生几个输出事件;
  • 转换是否异步;
  • 输入错误如何传播;
  • 输入结束后输出何时结束;
  • 转换过程中是否暂停上游;
  • 订阅取消时中间资源如何清理。

1. map:同步的一对一转换

final numbers = Stream.fromIterable([1, 2, 3]);

final strings = numbers.map((value) => 'number=$value');

await for (final value in strings) {
  print(value);
}

输出:

number=1
number=2
number=3

map 的典型形式是:

dif(di)d_i \mapsto f(d_i)

每个输入事件产生一个输出事件。map 回调返回普通值,不适合直接等待异步操作。

2. where:过滤事件

final evenNumbers = Stream.fromIterable([1, 2, 3, 4])
    .where((value) => value.isEven);

await for (final value in evenNumbers) {
  print(value);
}

输出:

2
4

where 的输出数量可能少于输入数量,但仍然保持剩余事件的相对顺序。

3. expand:一个输入产生多个同步输出

final result = Stream.fromIterable([1, 2, 3]).expand(
  (value) => [value, value * 10],
);

await for (final value in result) {
  print(value);
}

输出:

1
10
2
20
3
30

它实现的是:

di[fi1,fi2,]d_i \mapsto [f_{i1}, f_{i2}, \ldots]

4. asyncMap:异步的一对一转换

Future<String> loadName(int id) async {
  await Future<void>.delayed(const Duration(milliseconds: 10));
  return 'user-$id';
}

Future<void> main() async {
  final names = Stream.fromIterable([1, 2, 3]).asyncMap(loadName);

  await for (final name in names) {
    print(name);
  }
}

输出顺序为:

user-1
user-2
user-3

asyncMap 的重要语义是:异步转换完成后,才继续向下游发送对应结果;通常会对上游形成暂停边界,因此不会像直接在 listenasync 回调中那样失去等待关系。

但这不意味着所有外部副作用都自动变成事务。如果 loadName(1) 已经向服务器发出请求,之后订阅被取消,请求是否能撤回取决于 HTTP 客户端或底层资源本身。

5. asyncExpand:一个输入产生一个异步 Stream

Stream<int> childrenOf(int value) async* {
  yield value;
  await Future<void>.delayed(const Duration(milliseconds: 1));
  yield value * 10;
}

Future<void> main() async {
  final result = Stream.fromIterable([1, 2]).asyncExpand(childrenOf);

  await for (final value in result) {
    print(value);
  }
}

结果按输入顺序展开:

1
10
2
20

asyncExpand 适合分页、嵌套查询和“一个事件对应一段异步事件流”的场景。它与并发合并不同:如果业务希望多个输入同时发起并尽快合并结果,应显式设计并发策略,不能把 asyncExpand 当作无限并发操作。

6. transform:使用 StreamTransformer 表达完整协议

final transformer = StreamTransformer<int, String>.fromHandlers(
  handleData: (value, sink) {
    if (value < 0) {
      sink.addError(
        ArgumentError('value must be non-negative'),
      );
      return;
    }
    sink.add('ok:$value');
  },
);

final result = Stream.fromIterable([1, -1, 2]).transform(transformer);

result.listen(
  print,
  onError: (Object error) {
    print('error: $error');
  },
);

transform 可以同时控制数据、错误和 done,是封装协议解析、分帧、校验、状态机转换的主要工具。

转换器必须正确处理错误和关闭。例如一个“按行拆分”的转换器不能只处理普通数据,还必须考虑:

  • 一次输入是否包含半行;
  • 多次输入如何拼接;
  • 输入结束时残留内容是否算一行;
  • 输入错误是否透传;
  • 下游取消时缓存是否释放。

六、错误事件不是异常抛出的一种简单替代

1. 错误可以被 Stream 传递

final controller = StreamController<int>();

final subscription = controller.stream.listen(
  (value) => print('data=$value'),
  onError: (Object error, StackTrace stackTrace) {
    print('error=$error');
  },
  onDone: () => print('done'),
);

controller.add(1);
controller.addError(StateError('temporary failure'));
controller.add(2);
await controller.close();
await subscription.asFuture<void>();

概念上的事件顺序是:

data=1
error=temporary failure
data=2
done

错误是否终止 Stream,取决于数据源和订阅策略。错误事件本身不必然意味着 done。

2. await for 的错误处理

Future<void> consume(Stream<int> stream) async {
  try {
    await for (final value in stream) {
      print(value);
    }
  } catch (error, stackTrace) {
    print('捕获错误:$error');
  }
}

await for 会依次等待数据,遇到错误时抛出异常。若不捕获,错误会从当前异步函数继续向上传播。

循环中的 break 会结束消费,并触发订阅取消:

await for (final value in stream) {
  if (value == 10) {
    break;
  }
}

取消是异步清理过程,底层资源是否已经完成释放,需要等待相应的 Future,而不能仅凭循环退出就假定清理完成。


七、背压:Dart Stream 的边界不是完整的需求拉取协议

1. 背压是什么

设生产者产生事件的速率为 λp\lambda_p,消费者处理事件的速率为 λc\lambda_c

当:

λp>λc\lambda_p > \lambda_c

并且生产者持续运行时,未处理事件数量会增长:

B(t)=B(0)+0t(λp(τ)λc(τ))dτB(t) = B(0) + \int_0^t \left(\lambda_p(\tau)-\lambda_c(\tau)\right)d\tau

其中:

  • B(t)B(t) 是时刻 tt 的缓存数量;
  • λp\lambda_p 是生产速率;
  • λc\lambda_c 是消费速率。

如果没有暂停生产、丢弃事件、限流或持久化队列,B(t)B(t) 就可能持续增大,最终造成延迟、内存压力甚至进程崩溃。

这就是背压(backpressure):消费者向生产者反馈“当前无法继续按原速接收”的机制。

2. pause 是订阅级暂停,不是全局 demand 计数

final controller = StreamController<int>();

final subscription = controller.stream.listen((value) {
  print('consume $value');
});

subscription.pause();

controller.add(1);
controller.add(2);

await Future<void>.delayed(Duration.zero);
print('暂停期间不会调用 onData');

subscription.resume();

暂停只针对这个订阅者。它不表示整个 Stream 只有一个消费者,也不表示所有生产者都停止。

对于单订阅 Stream,如果底层资源支持暂停,暂停可能沿链路传播:

订阅 pause
   ↓
StreamSubscription 暂停交付
   ↓
上游 Stream 或 Controller 收到 onPause
   ↓
底层定时器、Socket 或读取操作可能暂停

但这条链路不是无条件成立的。

3. Controller 的缓存边界

如果生产者继续调用 controller.add,而订阅者处于暂停状态,事件通常会在 Stream 实现内部缓存,直到订阅恢复。

这只是“缓存”,不是无限容量的背压协议。StreamController 没有一个统一的、可查询的无限安全容量,也没有类似某些响应式协议中的逐条 demand:

消费者请求 1 个 → 生产者只发送 1 个

因此,暂停只能解决“暂时停止向订阅者交付”,不自动解决:

  • 生产者仍在高速生成事件;
  • 事件体积很大;
  • 暂停持续时间不可控;
  • 上游外部系统无法暂停;
  • 多个广播订阅者速度不同。

4. 广播 Stream 的背压更复杂

广播 Stream 有多个订阅者:

             ┌─ subscriber A:快速
source ──────┼─ subscriber B:暂停
             └─ subscriber C:慢速

每个订阅者可能有自己的暂停状态和缓存。一个订阅者暂停,不应自动阻塞其他订阅者;否则最慢的监听者会把整个广播源拖住。

因此,广播模式常见的边界是:

  • 每个订阅者独立暂停;
  • 暂停期间的数据可能按订阅者分别缓存;
  • 没有监听者时,广播事件通常丢弃;
  • 底层源是否暂停,取决于广播实现是否显式管理。

这说明广播 Stream 天生更适合“实时通知”,而不是“每条事件都必须被所有消费者可靠处理”。

5. 生产者不可暂停时怎么办

假设一个外部 SDK 持续回调传感器数据,而且没有暂停 API:

传感器持续产生
        ↓
Dart 回调持续进入
        ↓
订阅者 pause 只能暂停交付
        ↓
数据在内存中积累或被实现丢弃

此时必须在边界处选择策略:

  1. 降采样:只保留固定时间窗口内的最后一个值;
  2. 合并:把多个事件聚合成一个统计结果;
  3. 丢弃:只保留最新值或只在空闲时处理;
  4. 有界队列:超过容量后拒绝、覆盖或报错;
  5. 外部持久化:将可靠任务交给数据库或消息队列;
  6. 改变源 API:使用真正支持暂停或按需读取的接口。

例如“只保留最新值”的语义可以用 StreamTransformer 或状态层实现,但它改变了业务语义:消费者不再得到每一个原始事件,而是得到采样后的状态。


八、取消:停止交付不等于立即中止所有外部工作

1. cancel() 是异步操作

final subscription = stream.listen(print);

await subscription.cancel();
print('订阅清理完成');

cancel() 返回 Future<void>,因为取消可能需要异步释放资源,例如:

  • 关闭 Socket;
  • 停止文件读取;
  • 解除平台事件监听;
  • 等待原生对象释放;
  • 关闭控制器或内部任务。

如果后续逻辑依赖资源已经释放,应等待 cancel() 完成。

2. Controller 的 onCancel

Timer? timer;

final controller = StreamController<int>(
  onListen: () {
    var value = 0;
    timer = Timer.periodic(const Duration(milliseconds: 100), (_) {
      controller.add(value++);
    });
  },
  onCancel: () async {
    timer?.cancel();
    timer = null;
    print('timer stopped');
  },
);

这个示例的意图是:

第一次监听 → 创建 Timer
取消监听   → 停止 Timer

实际代码中不能在变量初始化表达式内部直接引用尚未完成初始化的 controller。更完整、可运行的写法是先声明可空变量:

import 'dart:async';

Future<void> main() async {
  StreamController<int>? controller;
  Timer? timer;

  controller = StreamController<int>(
    onListen: () {
      var value = 0;
      timer = Timer.periodic(const Duration(milliseconds: 100), (_) {
        controller!.add(value++);
      });
    },
    onCancel: () {
      timer?.cancel();
      timer = null;
      print('timer stopped');
    },
  );

  final subscription = controller.stream.listen((value) {
    print(value);
    if (value == 2) {
      unawaited(subscription.cancel());
    }
  });

  await Future<void>.delayed(const Duration(seconds: 1));
  await controller.close();
}

不过这个示例在 subscription 的初始化闭包中引用自身,虽然回调通常在初始化完成后才执行,但可读性较差。更稳妥的生产代码应将订阅变量设为可空,并在回调中判空:

import 'dart:async';

Future<void> main() async {
  late StreamController<int> controller;
  StreamSubscription<int>? subscription;
  Timer? timer;

  controller = StreamController<int>(
    onListen: () {
      var value = 0;
      timer = Timer.periodic(const Duration(milliseconds: 100), (_) {
        if (!controller.isClosed) {
          controller.add(value++);
        }
      });
    },
    onCancel: () {
      timer?.cancel();
      timer = null;
      print('timer stopped');
    },
  );

  subscription = controller.stream.listen((value) async {
    print(value);
    if (value == 2) {
      await subscription?.cancel();
    }
  });

  await Future<void>.delayed(const Duration(seconds: 1));
  await controller.close();
}

这里还要注意 controller.close()subscription.cancel() 的职责不同:

  • close():向订阅者发送 done,结束 Stream;
  • cancel():当前订阅者不再接收事件,并触发取消清理;
  • 底层资源是否释放:取决于 onCancel 或具体 Stream 实现。

3. cancelOnError

stream.listen(
  print,
  onError: (Object error) {
    print('error=$error');
  },
  cancelOnError: true,
);

设置后,第一次错误事件会导致该订阅自动取消。它适合“一次失败即不可继续”的流程,例如读取一个必须完整解析的文件。

如果错误是可恢复的网络抖动或单条数据校验失败,cancelOnError: true 可能会过早终止整个流,应该在错误回调中按业务决定是否继续。

4. 转换链的取消传播

final result = source
    .where(isValid)
    .asyncMap(load)
    .map(format);

final subscription = result.listen(print);
await subscription.cancel();

理想的数据流关系是:

下游 cancel
   ↓
map 取消上游
   ↓
asyncMap 取消上游
   ↓
where 取消 source
   ↓
source 释放底层资源

标准 Stream API 会为转换链提供取消传播,但外部副作用是否可撤回仍取决于实现。例如已经发出的网络请求不一定会因为 Stream 订阅取消而自动中止。需要取消 HTTP 请求时,必须把请求取消能力接入实际的客户端 API。


九、await for 与 Flutter 生命周期

1. await for 适合单一消费流程

Future<void> consume(Stream<String> messages) async {
  await for (final message in messages) {
    print('message=$message');
  }
}

它把 Stream 看作异步迭代器:

等待下一个事件
   ↓
执行循环体
   ↓
循环体完成后再请求下一个事件

因此,循环体中的 await 会自然形成消费节奏:

await for (final message in messages) {
  await save(message);
}

但这仍然不是持久化队列。应用进程被系统杀死、页面退出或网络断开时,内存中的事件可能消失。

2. Flutter Widget 中必须管理订阅生命周期

如果手动订阅 Stream,应保存并取消订阅:

class ExampleState extends State<Example> {
  StreamSubscription<int>? _subscription;

  @override
  void initState() {
    super.initState();

    _subscription = widget.values.listen((value) {
      if (!mounted) return;
      setState(() {
        // 更新界面状态
      });
    });
  }

  @override
  void dispose() {
    _subscription?.cancel();
    super.dispose();
  }
}

这里有两个独立问题:

  1. dispose 后不能再调用 setState
  2. 取消订阅可以阻止后续事件继续进入页面逻辑。

if (!mounted) return 只能保护 UI 更新,不能替代取消订阅。如果不取消,订阅仍可能持有 State、闭包或其他资源,形成生命周期泄漏。

如果使用 StreamBuilder,Flutter 会负责订阅和取消,但仍需理解其语义:更换 stream、Widget 重建和页面销毁都会影响订阅生命周期;StreamBuilder 也不会替单订阅 Stream 创建多个独立数据源。

3. 移动端、桌面和 Web 的差异

Dart Stream 的核心语义在 Android、iOS、桌面和 Web 上基本一致:单订阅、广播、暂停、取消和错误传播都属于 Dart Stream 抽象。

差异主要来自事件源:

  • Android、iOS 常通过插件接收原生回调、传感器、定位或平台通道事件;
  • 桌面平台可能连接文件系统、进程、Socket 或系统通知;
  • Web 受浏览器事件循环、页面后台限制和浏览器网络 API 约束;
  • Web 不具备移动端原生线程模型中的所有能力,某些插件事件源的生命周期也不同。

因此不能从 Dart Stream 保证推出“平台资源一定可暂停”或“取消一定能中止底层 I/O”。这些行为必须查看具体插件或平台 API 的实现。


十、一个完整的可运行示例:分页流、过滤、异步转换与取消

下面的示例模拟分页加载:

  • 上游产生页码;
  • asyncExpand 将每个页码展开为该页的记录;
  • where 过滤无效记录;
  • asyncMap 模拟异步详情加载;
  • 消费到某个值后取消。
import 'dart:async';

Future<List<int>> loadPage(int page) async {
  await Future<void>.delayed(const Duration(milliseconds: 50));
  return [page * 10, page * 10 + 1, page * 10 + 2];
}

Future<String> loadDetail(int id) async {
  await Future<void>.delayed(const Duration(milliseconds: 20));
  return 'detail-$id';
}

Stream<String> buildStream() {
  return Stream.fromIterable([1, 2, 3])
      .asyncExpand((page) => Stream.fromFuture(loadPage(page)).expand(
            (items) => items,
          ))
      .where((id) => id.isOdd)
      .asyncMap(loadDetail);
}

Future<void> main() async {
  final stream = buildStream();
  final subscription = stream.listen(
    (value) => print(value),
    onError: (Object error, StackTrace stackTrace) {
      print('error: $error');
    },
    onDone: () {
      print('done');
    },
  );

  await subscription.asFuture<void>();
}

预期输出:

detail-11
detail-21
detail-31
done

逐步推导如下。

第一步,页码事件为:

1, 2, 3

第二步,每个页码加载三条记录:

1 → 10, 11, 12
2 → 20, 21, 22
3 → 30, 31, 32

第三步,where((id) => id.isOdd) 保留:

11, 21, 31

第四步,asyncMap(loadDetail) 将它们转换为:

detail-11, detail-21, detail-31

第五步,源 Stream 正常结束,转换链向下游发送 done。

这个例子中的 asyncMap 维护了转换结果的事件顺序,但“顺序发送”不代表底层请求可以被撤销。若在 detail-21 处理期间取消订阅,已经发出的请求是否停止,必须由 loadDetail 所使用的 HTTP 或数据库 API 支持取消。


十一、常见误解与失败表现

1. 把广播 Stream 当作消息队列

错误假设:

所有事件都会被保存,任何时候加入的订阅者都能补收。

实际广播 Stream 通常不回放历史事件。订阅前产生的事件可能已经丢失,零订阅者期间产生的事件也可能被丢弃。

诊断方法是记录事件产生时间和订阅建立时间:

print('emit at ${DateTime.now()}');
print('listen at ${DateTime.now()}');

如果只在订阅建立前发送事件,监听者收不到并不一定是 Stream 异常,而可能正是广播语义。

2. 以为 asynclisten 回调会自动串行

错误代码:

stream.listen((value) async {
  await write(value);
});

这里 listen 不等待回调返回的 Future。需要顺序消费时,改为:

await for (final value in stream) {
  await write(value);
}

或:

stream.asyncMap(write).listen(print);

3. 以为 pause 一定让生产者停止

pause 首先保证的是:该订阅者暂时不接收数据。它是否向上游传播,取决于 Stream 实现和底层资源。

如果生产者不可暂停,暂停期间就可能积累缓存。需要验证生产者是否真正停止,可以在源头记录计数:

var produced = 0;

timer = Timer.periodic(const Duration(milliseconds: 10), (_) {
  produced++;
  controller.add(produced);
});

暂停订阅后,如果 produced 仍高速增长,说明暂停没有阻止源头生产。

4. 以为 cancel() 能撤销已经执行的副作用

取消 Stream 订阅只能控制 Stream 事件交付和取消传播。它不能自动回滚:

  • 已经写入数据库的数据;
  • 已经发送的 HTTP 请求;
  • 已经提交的原生平台操作;
  • 已经执行的文件写入。

如果操作具有不可逆副作用,应单独设计幂等键、事务、补偿或显式取消 API。

5. 只关闭 Controller,不取消外部资源

await controller.close();

这会结束 Stream,但如果 Timer、Socket 或平台监听仍然运行,生产端可能继续尝试发送事件,或者外部资源继续占用。

正确的资源模型应明确:

建立监听 → 创建外部资源
取消订阅 → 停止外部资源
关闭 Stream → 通知下游 done

关闭和取消不是同一个动作,除非具体实现明确把二者绑定。


十二、如何选择 Stream 形态

可以用以下因果条件判断,而不是只记 API 名称。

选择单订阅 Stream

当数据源具有一次性生命周期,并且需要完整消费:

文件读取、HTTP 响应体、一次解析流程、单连接消费

它能防止多个消费者无意间共享一个一次性资源。

选择广播 Stream

当多个观察者都只关心实时变化,且允许错过订阅前事件:

界面状态变化、连接状态、用户输入、实时通知

要特别评估慢订阅者、暂停缓存和零监听者期间的数据丢失。

选择转换链

当数据处理可以表达为事件级操作:

过滤 → 同步映射 → 异步加载 → 展开 → 错误处理

转换链的优势是取消、错误和结束信号可以沿链路传播;随意在 listen 中启动未等待的异步任务,则会脱离这条控制链。

不把 Stream 当作完整背压系统

Dart Stream 提供了订阅级 pause/resume,但没有统一的无限安全缓存,也没有跨所有数据源生效的需求计数协议。

当速率关系满足:

λp>λc\lambda_p > \lambda_c

必须额外决定“暂停、缓存、丢弃、合并、限流或持久化”中的哪一种。若这个选择没有被明确建模,系统只是把问题推迟到了内存、延迟或数据丢失上。


Dart Stream 的完整理解可以归结为一条事件生命周期:

数据源建立
   ↓
Stream 产生数据或错误
   ↓
订阅者接收、暂停或继续
   ↓
转换链保持或改变事件结构
   ↓
订阅取消或 Stream 发送 done
   ↓
底层资源完成清理

单订阅决定“谁拥有这条数据流”,广播决定“当前观察者如何共享实时事件”,转换决定“事件如何变形”,暂停决定“消费变慢时交付如何受限”,取消决定“订阅和资源如何结束”。只有把这几个边界分别验证,才能准确判断一个 Stream 方案是否真的满足业务语义。


系列导航与关联阅读

官方资料

本文依据 Flutter 与 Dart 官方文档重新梳理;正文与示例由 WR BLOG 编写。