Appearance
源码:labs/dart/host/bin/runtime_concurrency/stream_async_generator_demo.dart
- 原始路径:
labs/dart/host/bin/runtime_concurrency/stream_async_generator_demo.dart - 类型:
dart
dart
// Dart Stream / 异步生成器 / 手动背压 演示
//
// 对应 docs 01-runtime-concurrency/13-dart-stream-async-generator.md:
// 1. async* 生成器:yield 逐个产出,消费端逐条收到
// 2. StreamController:手动控制事件流
// 3. 手动背压:StreamSubscription.pause / resume 控制消费速率
// 4. await for:把 Stream 当异步序列循环消费
//
// 运行:dart run bin/runtime_concurrency/stream_async_generator_demo.dart
import 'dart:async';
// 1. async* 生成器:产出 1..n,每个间隔 200ms
Stream<int> countDown(int n) async* {
for (var i = n; i >= 1; i--) {
await Future.delayed(const Duration(milliseconds: 200));
yield i; // 产出后:消费端 listen 回调立即触发(同一事件循环)
}
}
// 2. StreamController:手动 push 事件 + 关闭
Stream<String> manualEmitter() {
final controller = StreamController<String>();
var tick = 0;
Timer.periodic(const Duration(milliseconds: 300), (timer) {
tick++;
controller.add('event#$tick');
if (tick >= 5) {
timer.cancel();
controller.close(); // 关闭后监听器收到 done
}
});
return controller.stream;
}
Future<void> main() async {
print('=== 1. async* 生成器 + listen ===');
await for (final v in countDown(3)) {
print(' countDown 收到: $v');
}
print('');
print('=== 2. StreamController 手动推送 ===');
await manualEmitter().listen((e) => print(' $e')).asFuture();
print('');
print('=== 3. 手动背压:pause 2 秒再恢复 ===');
final sub = countDown(4).listen((v) {
print(' [消费端] 收到 $v @ ${DateTime.now().millisecond}');
});
await Future.delayed(const Duration(milliseconds: 400));
print(' → 消费端暂停 2 秒(生产者 yield 到缓冲,恢复后继续)');
sub.pause();
await Future.delayed(const Duration(milliseconds: 2000));
sub.resume();
await sub.asFuture<void>().timeout(const Duration(seconds: 2), onTimeout: () {});
print('');
print('=== 4. 背压观察结论 ===');
print(' pause 时事件不会丢,但会排队;恢复后按顺序继续消费');
}