卷 IV · 流动CH 14深度 14/24

操作符只是套娃:中间那层没有线程

RxJava 留给很多人一个印象:操作符链上有队列、有调度、有中间对象。Flow 的操作符什么都没有——它就是一串直接的函数调用,嵌套得像一组俄罗斯套娃。看懂这一层,一堆「为什么 map 不切后台」之类的问题就自己散了。

操作符collector 包装flowOn上下文保持

map 的全部实现

先把它写出来,一共四行,而且没有省略:

fun <T, R> Flow<T>.map(transform: (T) -> R): Flow<R> = flow {
    collect { value ->          // 去 collect 上游
        emit(transform(value))  // 转换一下,再往下游 emit
    }
}

就这些。map 返回的是一个新的冷流,它的 collect 做的唯一一件事,就是去 collect 上游,并且把自己的处理包进那个 lambda 里

filteronEachtaketransform……全都是这个形状。所以一条链:

source.onEach { … }.map { … }.filter { … }.collect { … }

展开之后是这样一个嵌套结构:

filter.collect(你的 collect 块)
  └ map.collect(包了 filter 判断的块)
      └ onEach.collect(包了 map 转换的块)
          └ source.collect(包了 onEach 日志的块)
              └ 真正的生产逻辑在这里 emit

一个值从 emit 出来之后,走的是一条直的调用栈emit → onEach 的块 → map 的块 → filter 的块 → 你的 collect 块,然后一路返回。

看它走一遍

demo 里那个缩进就是调用栈的深度。注意被 filter 拦下的那些值——它们的路径到 filter 那一行就断了,下面的代码根本没被调用,不是「被丢进垃圾桶」。

◆ 这一章的核心

Flow 的操作符链没有队列、没有缓冲、没有中间集合、没有线程。

它就是一串函数调用。一个值从产生到被消费,整个过程发生在同一个协程、同一条调用栈上。

由此可以直接推出三件事:

  • map { } 里的代码跑在 collect 所在的线程上——它不会「切到后台」
  • 下游慢,上游的 emit卡在那儿——这就是下一章的背压
  • 操作符本身几乎零开销——真正的成本是你在 lambda 里写的东西

「map 不切后台」这件事有多常被误解

看这段代码,它在很多项目里都有:

// ❌ 以为解析跑在后台
repo.observeRawJson()                  // Room 的 Flow,冷的
    .map { parseHugeJson(it) }         // ← 这一行跑在哪儿?
    .collect { uiState = it }          // 在 Main 上 collect

答案是:跑在 Main 上。因为 collect 在 Main 的协程里,而 map 的 lambda 是被这条调用栈直接调用的。

把 demo 里那个开关打开,你会看到 flowOn(Dispatchers.IO) 加上之后的差别——但也只有一部分改变了。

flowOn 只管它上面那一段

flow { … }                     // ← IO
    .map { a() }               // ← IO
    .flowOn(Dispatchers.IO)    // ══ 分界线 ══
    .map { b() }               // ← 收集者的上下文(Main)
    .collect { c() }           // ← Main

规矩很直观:flowOn 影响它上游的一切,不影响下游。因为它改的是「上游在哪个上下文里生产」,而下游永远跟着收集者走。

这个方向和 RxJava 的 subscribeOnobserveOn 不太一样,是很多人从 Rx 转过来时的第一个坑。记法:flowOn 站在链条中间,往上看。

它内部做了什么

这里有一个和上一章那条「上下文保持」规矩相关的细节。emit 必须发生在 collect 的协程里,那 flowOn 怎么做到换上下文的?

答案:它开了一个新协程,中间架了一个 Channel。

     上游(跑在 IO 的一个新协程里)
        │  emit → 发进 Channel
        ▼
    ┌────────┐
    │ Channel│   默认容量 64
    └────────┘
        │  下游(跑在 collect 的协程里)从 Channel 收
        ▼
     下游 + 你的 collect 块

两个后果,都值得知道:

  • flowOn 自带一个缓冲(默认 64)。所以加了 flowOn 之后,上下游就解耦了——上游不会再被下游拖住,除非缓冲也满了。这个副作用经常正是你想要的,但它是副作用,不是主要功能。
  • 它是有成本的:一个协程 + 一个 Channel + 每个值一次跨线程传递。所以别在链条里到处撒 flowOn
⚠ 相邻的 flowOn 会融合,不相邻的不会
// 只产生一次上下文切换(编译期/运行时会融合)
.flowOn(Dispatchers.IO)
.flowOn(Dispatchers.Default)     // 后者被前者覆盖,只留一个

// ❌ 产生两次切换、两个 Channel、两个协程
.map { a() }.flowOn(Dispatchers.IO)
.map { b() }.flowOn(Dispatchers.Default)
.collect { }

第二种写法在语义上是对的(a 在 IO、b 在 Default),但代价翻倍。大多数情况下你只需要一个 flowOn,放在链条最上面那一段的末尾。

该在哪儿写 flowOn

回到前面那个错例,正确写法是:

// ✅ 解析跑在 Default,UI 更新在 Main
repo.observeRawJson()
    .map { parseHugeJson(it) }
    .flowOn(Dispatchers.Default)     // ← 上面这两行都去 Default
    .collect { uiState = it }        // ← 这一行仍然在 Main

更好的做法是根本不用在 UI 层写 flowOn让 Repository 自己保证主线程安全(第 12 章那条规矩)。

// Repository 里
fun observeTodos(): Flow<List<Todo>> =
    dao.observeRaw()
        .map { parseHugeJson(it) }
        .flowOn(Dispatchers.Default)     // ← 责任在这一层

// ViewModel / UI 里:什么都不用管
repo.observeTodos().collect { … }

那些「看起来像操作符但不是」的东西

有一类操作符打破了「纯函数调用」这个模型,因为它们内部会开协程

操作符内部有没有协程/Channel为什么
map filter onEach transform take没有纯函数调用,套娃
flowOn要换上下文,只能靠 Channel 中转
buffer conflate要解耦上下游,必须有个中间地带(下一章)
collectLatest mapLatest flatMapLatest要能取消上一个,就得让上一个跑在独立的子协程里
combine zip merge要同时收好几条流
debounce sample要用到定时器
channelFlow callbackFlow它们的整个用途就是「允许从别的上下文 emit」

这张表有一个很实用的读法:第一行是免费的,其余每一行都在花钱。链条上每出现一个下面几行的操作符,就多一个协程、多一次跨线程传递。

倒不是说要避免用它们——该用就用。但当你在一条 Flow 上串了七八个操作符还觉得慢的时候,先数一数有几个不在第一行。

✎ channelFlow 和 flow 的分界

什么时候必须用 channelFlow当你需要从多个地方、或者从别的协程里 emit 的时候

// ❌ flow { } 里不能这么干(上下文保持不允许)
flow {
    launch { emit(a()) }      // 编译不过 / 运行时炸
    launch { emit(b()) }
}

// ✅ channelFlow 就是为这个存在的
channelFlow {
    launch { send(a()) }
    launch { send(b()) }
}

代价:channelFlow 内部真的有 Channel 和协程,比 flow { } 贵。能用 flow { } 就别用它。

callbackFlowchannelFlow 的特化版,专门用来包回调式 API,而且它强制你写 awaitClose { }——不写编译期就报错。这是个很好的设计:它让你没法忘记注销监听器。

顺带说 combine 的一个坑

combine 每次任意一条流发出新值就会重新计算,用的是所有流的最新值

combine(userFlow, settingsFlow, notificationsFlow) { u, s, n ->
    UiState(u, s, n)
}

两个后果:

  • 必须每条流都至少发过一个值,它才会第一次输出。有一条流迟迟不发,整个 UI 就一直是初始状态——这是「界面一直在转圈」的一个常见原因。给每条流一个初始值(onStart { emit(默认值) } 或者用 StateFlow)。
  • 它会输出很多中间态。三条流各变一次,你会收到三次输出,其中前两次是「新旧混合」的。UI 上通常没关系(会很快被下一个覆盖),但如果你在这里触发副作用就要小心。
⌗ 到你手上

这一章的一句话

操作符是套娃不是流水线:map 的 collect 去调上游的 collect,只是把 collector 包了一层。整条链是一串直接函数调用,所以中间那层没有线程——想换线程只有 flowOn,而它只管它上面那一段。

下一章:既然一个值从 emitcollect 是一条直的调用栈,那如果下游处理得很慢会怎么样?答案是emit 就卡在那儿——而这就是背压。Flow 对背压的默认答案不是丢也不是缓存,是「我等你」。