Dart Futures实现生产者消费者:能否替代轮询等待数据结构非空?
用Dart Stream替代轮询实现生产者-消费者模式
当然可以不用轮询,Dart Stream就是为这种生产者主动推送、消费者被动接收的场景设计的,比你现在的轮询方式更高效、更简洁。
为什么轮询不是最优解?
你当前的100ms轮询方式有两个明显问题:
- 资源浪费:不管pipe有没有数据,消费者都会定期唤醒检查
- 延迟不确定:如果数据刚被生产者加入pipe,消费者可能要等接近100ms才能处理,实时性差
用Stream实现的核心思路
用StreamController创建一个可双向操作的管道:
- 生产者通过
streamController.sink.add()向管道推送数据 - 消费者通过
await for循环或者stream.listen()监听管道,有数据时自动触发处理,无需主动轮询
而且因为Dart是单线程事件循环模型,Stream的操作天然是线程安全的,不用额外加同步锁。
代码示例
基础实现(await for 顺序处理)
import 'dart:async'; void main() async { // 创建Stream控制器,作为数据管道 final controller = StreamController<String>(); // 生产者函数:持续生成数据 Future<void> producer() async { var count = 0; while (true) { await Future.delayed(const Duration(milliseconds: 500)); final data = "消息${count++}"; controller.sink.add(data); print("生产者推送:$data"); } } // 消费者函数:等待并处理数据 Future<void> consumer() async { // await for会自动等待Stream的新数据,不用轮询 await for (final data in controller.stream) { print("消费者处理:$data"); // 模拟处理耗时 await Future.delayed(const Duration(milliseconds: 100)); } } // 启动生产者和消费者 producer(); await consumer(); // 实际场景记得关闭控制器,避免内存泄漏 // await controller.close(); }
用listen回调的方式(非顺序处理)
如果不需要在async函数里按顺序处理数据,也可以用listen注册回调:
void main() { final controller = StreamController<String>(); // 生产者 Future<void> producer() async { var count = 0; while (true) { await Future.delayed(const Duration(milliseconds: 500)); controller.sink.add("消息${count++}"); } } // 消费者 void consumer() { controller.stream.listen((data) { print("消费者处理:$data"); }); } producer(); consumer(); }
额外优化点
- 如果需要限制缓冲区大小(防止生产者速度远快于消费者导致内存占用过高),可以创建带缓冲区的
StreamController,比如StreamController<String>(bufferingCapacity: 10) - 当不再需要管道时,一定要调用
controller.close()关闭,避免内存泄漏
内容的提问来源于stack exchange,提问作者David Tonhofer
相关产品推荐
相关产品推荐

