Flow — 是来自Kotlin Coroutines库的一种异步数据流类型,实现了冷语义。根据Kotlin Documentation, 2025,Flow允许使用map、filter、catch和collect操作符发射值序列。与LiveData不同,Flow构建在协程之上并支持背压。
要点
Flow — 是来自kotlinx.coroutines.flow包的一种类型,表示冷异步数据流。本质上,Flow是一个协程序列,通过emit()函数发射值,并以成功或异常结束。流的收集通过终端操作符collect()完成,这是一个挂起函数。
冷流 意味着flow构建器内的代码为每个订阅者重新执行。RxJava中的Observable.fromIterable行为类似:新订阅者从头开始接收所有值。在Flow中,这是通过挂起函数collect实现的,该函数在整个数据收集期间阻塞协程。
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——原始流保持不变。这允许构建没有副作用的安全的转换链。
map 操作符通过异步或同步转换来转换流的每个值。filter只允许满足条件的值通过。catch在终端操作符之前捕获异常并允许恢复流。flatMapLatest在新值到达时取消先前的发射——类似于Rx中的switchMap。
Flow中的debounce 操作符将值的发布延迟指定的超时时间。如果在这段时间内到达新值——计时器重置。在Android中,debounce用于搜索:请求仅在300-400毫秒暂停后发送,这减少了3-5倍的API调用次数。
除了collect(),Flow还支持其他终端操作符:toList() 将所有值收集到列表中——对测试有用,first() 返回第一个元素并取消流,single() 期望恰好一个元素。fold(initial) 通过传入的函数累积值。所有终端操作符都是挂起函数,必须在协程或其他挂起函数内调用。
第一个示例——基本的Flow,通过map操作符生成数字并进行转换:
val numberFlow = flow {
for (i in 1..5) {
delay(500)
emit(i)
}
}
scope.launch {
numberFlow
.map { "数字:$it" }
.collect { value ->
println(value)
}
}
第二个示例——通过catch进行过滤和错误处理的流转换:
flow {
emit("data1")
emit("data2")
throw RuntimeException("network error")
}
.catch { e ->
emit("fallback_data")
}
.collect { value ->
println(value)
}
第三个示例——在ViewModel中使用StateFlow为Jetpack Compose提供响应式UI:
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 — 是一种具有单个当前值的流。它缓存最后一个值并立即将其传递给新订阅者。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构建器内且在终端操作符之前,它会被传递给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 — 是仅用于UI层的生命周期感知组件。Google推荐Flow用于业务逻辑和仓库,LiveData用于ViewModel中的简单观察。
StateFlow — 当需要存储UI状态(任务列表、搜索文本、加载标志)时 — 每个订阅者获取当前值。SharedFlow — 用于一次性事件(导航、Snackbar)。StateFlow不应用于事件,因为新值可能会被重复处理。
在Flow中,背压通过挂起机制实现:如果收集器正在处理前一个值,emit()会暂停协程。ChannelFlow中的通道(Channel)具有capacity大小的缓冲区。溢出时:suspending(等待)、drop(丢弃)或conflate(用最后一个替换)。
使用callbackFlow — 用于回调API的Flow构建器。内部在回调中使用emit(value)调用registerCallback()。awaitClose确保在协程取消时调用unregisterCallback()。callbackFlow在底层通过Channel(UNLIMITED)支持缓冲。
可以,通过转换器:Flow.asObservable() 来自kotlinx-coroutines-rx3包,将Flow转换为RxJava 3的Observable。反向转换 — CompletableSource.asFlow() 用于Single/Completable/Maybe。这在大项目中从RxJava迁移到协程时很有用。
总结
我们将开发一款交钥匙移动应用程序
IT Sectr自2017年以来为初创企业和企业打造iOS和Android应用程序。我们将为您提供咨询并提出最佳解决方案。