Flow — 是什么,Kotlin协程中的冷流和热流

作者: IT Sectr 发布日期: 2026-03-17 阅读时间: 9 分钟

Flow — 是来自Kotlin Coroutines库的一种异步数据流类型,实现了冷语义。根据Kotlin Documentation, 2025,Flow允许使用map、filter、catch和collect操作符发射值序列。与LiveData不同,Flow构建在协程之上并支持背压。

要点

  • Flow — Kotlin Coroutines中的冷异步数据流,在收集前不发射值
  • 冷流 — 每个订阅者从头开始自己的独立发射
  • 热流 (SharedFlow, StateFlow) — 独立于订阅者发射值
  • 操作符 map, filter, catch, debounce, flatMapLatest 无阻塞地转换流
  • Flow 通过StateFlow和collectAsState()与Jetpack Compose完全兼容

Kotlin中的Flow是什么?

Flow — 是来自kotlinx.coroutines.flow包的一种类型,表示冷异步数据流。本质上,Flow是一个协程序列,通过emit()函数发射值,并以成功或异常结束。流的收集通过终端操作符collect()完成,这是一个挂起函数。

冷语义

冷流 意味着flow构建器内的代码为每个订阅者重新执行。RxJava中的Observable.fromIterable行为类似:新订阅者从头开始接收所有值。在Flow中,这是通过挂起函数collect实现的,该函数在整个数据收集期间阻塞协程。

Flow构建器

Kotlin提供了几种创建Flow的方式:flow { } — 使用emit()的基本构造,flowOf(vararg values) — 用于固定值集合,.asFlow() — 用于集合和Sequence的扩展。所有构建器都是冷的——数据仅在调用终端操作符时生成。

冷流和热流

冷流和热流的划分是响应式编程的核心概念。冷流 (Flow, Observable) 在订阅时开始生成数据。热流 (Channel, SharedFlow) 独立发射数据——订阅者只接收订阅后发生的内容,不包括序列的开头。

SharedFlow — 是一种可以拥有多个订阅者的热Flow,在设置replay时可以重放最后的值。SharedFlow适用于事件(一次性通知)。StateFlow — 是其具有固定状态值的变体,为新订阅者缓存最后一个值。

ChannelFlow 在底层使用Channel,结合了Flow和Channel的特性。它通过容量(capacity)支持缓冲和背压。当值从不同协程发射时,ChannelFlow在将回调API转换为响应式流时很有用。

冷热转换

要将冷Flow转换为热SharedFlow,使用shareIn(scope, started, replay)操作符。started参数控制启动时机:SharingStarted.WhileSubscribed() — 当有订阅者时激活,Lazily — 在第一个订阅者时启动,Eagerly — 立即启动。反向转换——热转冷:StateFlow.asFlow()返回一个冷Flow,在collect时发射StateFlow的当前值。这便于测试。

Flow操作符

Flow 提供了丰富的操作符集,作为协程内的挂起函数工作。操作符没有状态并返回新的Flow——原始流保持不变。这允许构建没有副作用的安全的转换链。

map 操作符通过异步或同步转换来转换流的每个值。filter只允许满足条件的值通过。catch在终端操作符之前捕获异常并允许恢复流。flatMapLatest在新值到达时取消先前的发射——类似于Rx中的switchMap。

Flow中的debounce 操作符将值的发布延迟指定的超时时间。如果在这段时间内到达新值——计时器重置。在Android中,debounce用于搜索:请求仅在300-400毫秒暂停后发送,这减少了3-5倍的API调用次数。

终端操作符

除了collect(),Flow还支持其他终端操作符:toList() 将所有值收集到列表中——对测试有用,first() 返回第一个元素并取消流,single() 期望恰好一个元素。fold(initial) 通过传入的函数累积值。所有终端操作符都是挂起函数,必须在协程或其他挂起函数内调用。

Flow代码示例

第一个示例——基本的Flow,通过map操作符生成数字并进行转换:

kotlin
val numberFlow = flow {
    for (i in 1..5) {
        delay(500)
        emit(i)
    }
}

scope.launch {
    numberFlow
        .map { "数字:$it" }
        .collect { value ->
            println(value)
        }
}

第二个示例——通过catch进行过滤和错误处理的流转换:

kotlin
flow {
    emit("data1")
    emit("data2")
    throw RuntimeException("network error")
}
    .catch { e ->
        emit("fallback_data")
    }
    .collect { value ->
        println(value)
    }

第三个示例——在ViewModel中使用StateFlow为Jetpack Compose提供响应式UI:

kotlin
class SearchViewModel : ViewModel() {
    private val _query = MutableStateFlow("")
    val results: StateFlow<List<Result>> = _query
        .debounce(300)
        .flatMapLatest { query ->
            repository.search(query)
        }
        .catch { emit(emptyList()) }
        .stateIn(viewModelScope, SharingStarted.WhileSubscribed(5000), emptyList())

    fun onQueryChanged(query: String) {
        _query.value = query
    }
}

StateFlow和SharedFlow

StateFlow — 是一种具有单个当前值的流。它缓存最后一个值并立即将其传递给新订阅者。StateFlow是状态的可观察容器,支持equals比较——如果新值与当前值相同,则不会发生发射。Jetpack Compose通过collectAsState()使用StateFlow。

SharedFlow — 是一种更灵活的热Flow,没有强制初始值。SharedFlow通过replay(新订阅者的值数量)、extraBufferCapacity(replay之外的缓冲区)和onBufferOverflow(溢出策略)进行配置。SharedFlow非常适合一次性事件:导航、Snackbar、分析。

Flow 在Android架构中被Google推荐为主要数据源(层:Repository → UseCase → ViewModel)。LiveData在灵活性上不如Flow:Flow支持协程、操作符、背压并且在UI层之外工作。从LiveData迁移到Flow是现代Android项目中的标准实践。

在ViewModel中使用Flow时,正确选择类型很重要。StateFlow 非常适合需要承受屏幕旋转的UI状态。SharedFlow适用于重复处理不可接受的事件——例如导航。在lifecycleScope中使用collect()的Flow提供了对执行上下文的最大控制,但在离开屏幕时需要手动取消。

Flow的测试通过kotlinx-coroutines-test进行。该库提供TestDispatcher——一种虚拟时间,可以加速延迟(delay)并控制协程的执行顺序。TestScope.runTest { } 创建了测试Flow的隔离环境。toList()操作符常用于测试中,通过超时收集所有flow值,以检查流是否发射了正确的数据序列。

Flow与Room(Android数据库库)良好集成:DAO方法可以返回Flow<List<Entity>>。Room在表发生任何更改时自动发射新值——UI无需手动触发即可更新。这是通过InvalidationTracker实现的,它在底层使用带有callbackFlow的Flow。这种方法消除了对LiveData的需求,并使数据层完全面向协程。Jetpack Compose通过collectAsState()订阅StateFlow,只重新绘制数据发生变化的组件——这提供了LiveData面向架构无法达到的性能。DataStore(SharedPreferences的替代品)也返回Flow<Preferences>,确保无需手动更新触发器即可响应式地读取应用设置。

Flow通过kotlinx-coroutines-core在JVM上支持进程间通信,无需额外库。例如,在Ktor的服务器应用程序中,Flow可以表示传入的WebSocket消息流。每条消息被发射到流中,通过操作符进行过滤和聚合,然后将结果发送给客户端。这种方法取代了Kotlin项目中的响应式库如Reactor或RxJava。

Flow与现有RxJava代码的兼容性由kotlinx-coroutines-rx3模块提供。扩展函数Flow.asObservable()将Flow转换为RxJava 3的Observable。反向转换——CompletableSource.asFlow(), Observable.asFlow()。这简化了从RxJava到协程的迁移:可以逐步重写项目,将部分层保留在RxJava上。转换时需要考虑冷热语义的差异:Observable可以是冷流或热流,而Flow对于普通Flow始终是冷流,对于SharedFlow是热流。

Flow的错误处理和测试

Flow中的错误处理有一个特点:如果异常发生在flow构建器内且在终端操作符之前,它会被传递给catch。如果异常发生在构建器之后的操作符中,该操作符之后的catch会捕获它。retryWhen 允许有条件地重试订阅:网络错误时最多重试3次,但在CancellationException时不重试。Flow消除了状态相关的错误,因为它不存储状态——与Observable中Subject存储内部状态相比,这简化了调试。

使用kotlinx-coroutines-test测试Flow使用TestDispatcher来模拟延迟。Turbine——社区流行的Flow测试库:test { } 启动Flow,awaitItem() 等待下一个值,awaitComplete() 等待完成。Turbine添加了默认超时,防止测试挂起。测试StateFlow时使用.testIn(scope)并按时间顺序检查值。

常见问题

Flow和LiveData有什么区别?

Flow — 是一个支持协程、操作符和背压的异步流,可以在任何架构层工作。LiveData — 是仅用于UI层的生命周期感知组件。Google推荐Flow用于业务逻辑和仓库,LiveData用于ViewModel中的简单观察。

什么时候使用StateFlow而不是SharedFlow?

StateFlow — 当需要存储UI状态(任务列表、搜索文本、加载标志)时 — 每个订阅者获取当前值。SharedFlow — 用于一次性事件(导航、Snackbar)。StateFlow不应用于事件,因为新值可能会被重复处理。

Flow中的背压如何工作?

Flow中,背压通过挂起机制实现:如果收集器正在处理前一个值,emit()会暂停协程。ChannelFlow中的通道(Channel)具有capacity大小的缓冲区。溢出时:suspending(等待)、drop(丢弃)或conflate(用最后一个替换)。

如何将回调转换为Flow?

使用callbackFlow — 用于回调API的Flow构建器。内部在回调中使用emit(value)调用registerCallback()。awaitClose确保在协程取消时调用unregisterCallback()。callbackFlow在底层通过Channel(UNLIMITED)支持缓冲。

Flow可以与RxJava一起使用吗?

可以,通过转换器:Flow.asObservable() 来自kotlinx-coroutines-rx3包,将Flow转换为RxJava 3的Observable。反向转换 — CompletableSource.asFlow() 用于Single/Completable/Maybe。这在大项目中从RxJava迁移到协程时很有用。

总结

  • Flow — Kotlin Coroutines中带有collect挂起函数的冷异步数据流
  • 冷流 为每个订阅者重新开始发射
  • StateFlow — 缓存最后一个值的热状态容器
  • SharedFlow — 具有replay和缓冲区配置的用于事件的热流
  • 操作符 map, filter, debounce, catch, flatMapLatest — 流转换的基础
  • Google 推荐Flow作为现代Android架构中的主要数据源
  • LiveData仅适用于UI层,Flow — 适用于应用的所有层

我们将开发一款交钥匙移动应用程序

IT Sectr自2017年以来为初创企业和企业打造iOS和Android应用程序。我们将为您提供咨询并提出最佳解决方案。

讨论项目

另请阅读