Dart Stream中是否存在Kotlin Flow.mapLatest的等价实现?
Dart Stream 对应 Kotlin Flow.mapLatest() 的实现方式
Dart 标准库中并没有直接提供和 Kotlin Flow.mapLatest() 完全等价的函数,但我们可以借助 stream_transform 包的 switchLatest 扩展方法,快速实现同样的逻辑。
核心逻辑回顾
Kotlin 的 mapLatest() 核心是:当原流发出新元素时,立即取消前一个元素的转换计算,只保留最新元素的转换结果输出。这个逻辑可以拆解为两步:
- 将原流的每个元素通过转换函数转为一个异步的单元素流(由转换结果的 Future 生成)
- 始终只订阅最新生成的那个单元素流,忽略之前未完成的流
具体实现代码
先确保引入 stream_transform 依赖,然后给 Stream 写一个扩展方法:
import 'dart:async'; import 'package:stream_transform/stream_transform.dart'; extension StreamMapLatestExtension<T> on Stream<T> { Stream<R> mapLatest<R>(FutureOr<R> Function(T) transform) { // 把每个元素转成由转换结果生成的单元素流 return map((element) => Stream.fromFuture(Future.value(transform(element)))) // 只保留最新的流的输出 .switchLatest(); } }
使用示例
比如我们有一个输入流,每次输入触发一个异步操作(比如网络请求),用 mapLatest 就能确保只处理最新的请求:
void main() { final inputStream = Stream.periodic(const Duration(milliseconds: 500), (i) => i); inputStream .mapLatest((value) async { // 模拟异步操作,比如网络请求 await Future.delayed(const Duration(milliseconds: 800)); return "处理结果: $value"; }) .listen((result) { print(result); // 输出只会是 "处理结果: 1"、"处理结果: 3"、"处理结果: 5"... 因为每次新元素会打断前一个未完成的异步操作 }); }
不依赖第三方包的手动实现
如果不想引入 stream_transform,也可以手动实现类似逻辑:
extension StreamMapLatestManual<T> on Stream<T> { Stream<R> mapLatest<R>(FutureOr<R> Function(T) transform) { return Stream<R>.multi((controller) { StreamSubscription? activeSubscription; final sourceSubscription = listen((element) { // 取消前一个未完成的订阅 activeSubscription?.cancel(); // 处理当前元素的异步转换 final future = Future.value(transform(element)); activeSubscription = future.asStream().listen( controller.add, onError: controller.addError, onDone: () => activeSubscription = null, ); }, onError: controller.addError, onDone: controller.close); controller.onCancel = () { activeSubscription?.cancel(); return sourceSubscription.cancel(); }; }); } }
这个手动实现的逻辑和依赖包的方式一致:每次原流有新元素时,取消前一个异步操作的订阅,只处理最新的元素转换。
内容的提问来源于stack exchange,提问作者Daniel Valencia
相关产品推荐
相关产品推荐

