Dart广播流监听后立即添加事件丢失,如何无延迟捕获所有事件?
Dart广播流监听后立即添加的事件被忽略的原因与解决方法
问题描述
在监听Dart广播流后立即向其中添加事件,发现最初添加的事件(1、2、3)被忽略,只有延迟后添加的事件(4、5、6)会被监听器打印。但移除.broadcast使用非广播流,或直接用streamController.stream替换yieldedStream()时,所有数字都能正常打印。代码示例如下:
import 'dart:async'; final streamController = StreamController<int>.broadcast(); Stream<int> yieldedStream() async* { yield 0; yield* streamController.stream; } Future<void> main(List<String> arguments) async { final subscription = yieldedStream().listen(print); streamController.add(1); streamController.add(2); streamController.add(3); await Future.delayed(Duration(seconds: 1)); streamController.add(4); streamController.add(5); streamController.add(6); await Future.delayed(Duration(seconds: 1)); subscription.cancel(); }
原因分析
1. Async*生成器的异步执行特性
调用yieldedStream().listen(print)时,async*生成器函数不会立即同步执行,而是被调度到微任务队列中等待执行。也就是说,yield 0和yield* streamController.stream这两行代码,要等main函数里的同步代码(add(1)/add(2)/add(3))全部执行完毕后才会运行。
2. 广播流的无缓存特性
广播流(StreamController.broadcast()创建)不会缓存未被订阅者接收的事件。当main函数同步执行add(1)等操作时,yieldedStream()生成的流还没完成对广播流的订阅(因为yield*还没执行),此时这些事件没有任何订阅者接收,直接被丢弃。
3. 非广播流/直接订阅的差异
- 非广播流:属于单订阅流,订阅会立即触发流的绑定,并且会缓存事件直到订阅者准备就绪,所以同步添加的事件不会丢失。
- 直接订阅控制器流:
streamController.stream.listen(print)的订阅操作是同步完成的,广播流在add(1)之前已经有了订阅者,因此所有事件都能被捕获。
解决方法
方法1:等待微任务队列完成订阅
通过await Future.microtask(() {})让当前同步代码执行完毕,等待微任务队列中的订阅操作完成后再发送事件,这是最稳妥的方案:
Future<void> main(List<String> arguments) async { final subscription = yieldedStream().listen(print); // 等待微任务队列执行,确保广播流的订阅已完成 await Future.microtask(() {}); streamController.add(1); streamController.add(2); streamController.add(3); await Future.delayed(Duration(seconds: 1)); streamController.add(4); streamController.add(5); streamController.add(6); await Future.delayed(Duration(seconds: 1)); subscription.cancel(); }
方法2:使用同步广播流控制器
创建广播流控制器时指定sync: true参数,让add操作同步触发订阅者回调。注意:这可能导致同步执行栈过深,仅在确定不会引发栈溢出的场景下使用:
final streamController = StreamController<int>.broadcast(sync: true);
方法3:调整生成器实现,确保订阅同步完成
如果业务允许,可在生成器内部标记订阅完成状态,等待订阅就绪后再发送事件:
final subscriptionCompleter = Completer<void>(); Stream<int> yieldedStream() async* { yield 0; final broadcastStream = streamController.stream; // 订阅广播流并标记完成状态 broadcastStream.listen((_) {}, onDone: () {}); subscriptionCompleter.complete(); yield* broadcastStream; } Future<void> main(List<String> arguments) async { final subscription = yieldedStream().listen(print); // 等待订阅完成 await subscriptionCompleter.future; streamController.add(1); streamController.add(2); streamController.add(3); await Future.delayed(Duration(seconds: 1)); streamController.add(4); streamController.add(5); streamController.add(6); await Future.delayed(Duration(seconds: 1)); subscription.cancel(); }
内容的提问来源于stack exchange,提问作者Tales Barreto
相关产品推荐
相关产品推荐

