You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.07 00:22:51