云计算百科
云计算领域专业知识百科平台

FusibleFlow:Kotlin Flow 的融合机制原理

本文以 FusibleFlow 为核心,通过源码片段解析 Flow 的\”融合\”(fuse)优化:flowOn 如何触发融合、fuse 如何把多层上下文合并为一次线程切换,以及无法融合时 ChannelFlow 兜底方案(collectTo)的工作方式。

1. 一个典型的 Flow 调用链示例

flowOf(1, 2, 3)
.map {

it * 2 }
.filter {

it > 2 }
.flowOn(Dispatchers.IO)

描述:这是一个典型的冷 Flow 链式调用,也是触发 FusibleFlow 融合的例子:

  • flowOf(1, 2, 3) 创建发射 1、2、3 的 Flow;
  • map { it * 2 } 和 filter { it > 2 } 这类中间操作符,底层都会把 Flow 包装成 TransformFlow 等实现了 FusibleFlow 接口的内部类;
  • 最后调用 flowOn(Dispatchers.IO) 时,走到源码里的 this is FusibleFlow -> fuse(context = context) 分支——由于前面的 map/filter 包装器都是 FusibleFlow,上下文会一路\”融合\”进整条链,最终不会创建任何 Channel。

也就是说,这条链的处理是:1,2,3 → 2,4,6 → 4,6,全部在 IO 线程上一次发射完成,collect 收集到的就是 4 和 6。注意 flowOn 只影响上游操作符(flowOf/map/filter)的执行线程,不影响下游 collect 所在的协程。

对比:如果上游不是 FusibleFlow(例如一个自定义的、未实现该接口的 Flow),flowOn 就只能走 else 分支,创建 ChannelFlowOperatorImpl 来做线程切换——这正是后面几节要讲的内容。

2. FusibleFlow:可融合 Flow 的内部接口

internal interface FusibleFlow<T> : Flow<T> {


fun fuse(context: CoroutineContext): Flow<T>?
}

描述:这是 kotlinx.coroutines 内部的接口(internal,不对使用者暴露)。它定义了一个 fuse 方法,表示\”实现类知道如何把一个新的上下文融合进自己\”。fuse 返回 Flow<T>?:

  • 返回融合后的新 Flow(或自身)——融合成功;
  • 返回 null —— 无法融合。

融合的意义:当多个 flowOn/操作符连续出现时,不必为每一次上下文切换都创建一个 Channel,而是把上下文合并到一起,只需最外层做一次真正的线程切换,大幅减少开销。

3. flowOn 的源码实现

public fun <T> Flow<T>.flowOn(context: CoroutineContext): Flow<T> {

赞(0)
未经允许不得转载:网硕互联帮助中心 » FusibleFlow:Kotlin Flow 的融合机制原理
分享到: 更多 (0)

评论 抢沙发

评论前必须登录!