Skip to content

源码: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 时事件不会丢,但会排队;恢复后按顺序继续消费');
}

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