操作符只是套娃:中间那层没有线程
RxJava 留给很多人一个印象:操作符链上有队列、有调度、有中间对象。Flow 的操作符什么都没有——它就是一串直接的函数调用,嵌套得像一组俄罗斯套娃。看懂这一层,一堆「为什么 map 不切后台」之类的问题就自己散了。
map 的全部实现
先把它写出来,一共四行,而且没有省略:
fun <T, R> Flow<T>.map(transform: (T) -> R): Flow<R> = flow {
collect { value -> // 去 collect 上游
emit(transform(value)) // 转换一下,再往下游 emit
}
}
就这些。map 返回的是一个新的冷流,它的 collect 做的唯一一件事,就是去 collect 上游,并且把自己的处理包进那个 lambda 里。
filter、onEach、take、transform……全都是这个形状。所以一条链:
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 的 subscribeOn/observeOn 不太一样,是很多人从 Rx 转过来时的第一个坑。记法:flowOn 站在链条中间,往上看。
它内部做了什么
这里有一个和上一章那条「上下文保持」规矩相关的细节。emit 必须发生在 collect 的协程里,那 flowOn 怎么做到换上下文的?
答案:它开了一个新协程,中间架了一个 Channel。
上游(跑在 IO 的一个新协程里)
│ emit → 发进 Channel
▼
┌────────┐
│ Channel│ 默认容量 64
└────────┘
│ 下游(跑在 collect 的协程里)从 Channel 收
▼
下游 + 你的 collect 块
两个后果,都值得知道:
flowOn自带一个缓冲(默认 64)。所以加了flowOn之后,上下游就解耦了——上游不会再被下游拖住,除非缓冲也满了。这个副作用经常正是你想要的,但它是副作用,不是主要功能。- 它是有成本的:一个协程 + 一个 Channel + 每个值一次跨线程传递。所以别在链条里到处撒
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:当你需要从多个地方、或者从别的协程里 emit 的时候。
// ❌ flow { } 里不能这么干(上下文保持不允许)
flow {
launch { emit(a()) } // 编译不过 / 运行时炸
launch { emit(b()) }
}
// ✅ channelFlow 就是为这个存在的
channelFlow {
launch { send(a()) }
launch { send(b()) }
}
代价:channelFlow 内部真的有 Channel 和协程,比 flow { } 贵。能用 flow { } 就别用它。
callbackFlow 是 channelFlow 的特化版,专门用来包回调式 API,而且它强制你写 awaitClose { }——不写编译期就报错。这是个很好的设计:它让你没法忘记注销监听器。
顺带说 combine 的一个坑
combine 每次任意一条流发出新值就会重新计算,用的是所有流的最新值:
combine(userFlow, settingsFlow, notificationsFlow) { u, s, n ->
UiState(u, s, n)
}
两个后果:
- 必须每条流都至少发过一个值,它才会第一次输出。有一条流迟迟不发,整个 UI 就一直是初始状态——这是「界面一直在转圈」的一个常见原因。给每条流一个初始值(
onStart { emit(默认值) }或者用StateFlow)。 - 它会输出很多中间态。三条流各变一次,你会收到三次输出,其中前两次是「新旧混合」的。UI 上通常没关系(会很快被下一个覆盖),但如果你在这里触发副作用就要小心。
在项目里扫一遍所有 Flow 链,对每一条问三句:
1. 有没有重活(解析、排序、加解密)写在 map/onEach 里,
而链条上没有 flowOn?
→ 它在 Main 上跑。加 flowOn,或者挪到 Repository。
2. flowOn 是不是写在了链条最下面(紧挨着 collect)?
→ 那它可能没起作用:collect 本来就在收集者的上下文里。
flowOn 要写在「重活那一段」的下面、UI 那一段的上面。
3. 有没有连续两个以上会开协程的操作符
(flowOn / buffer / *Latest / combine / debounce)?
→ 每个都是一个协程 + 一个 Channel,数一数值不值。
这一章的一句话
操作符是套娃不是流水线:map 的 collect 去调上游的 collect,只是把 collector 包了一层。整条链是一串直接函数调用,所以中间那层没有线程——想换线程只有 flowOn,而它只管它上面那一段。
下一章:既然一个值从 emit 到 collect 是一条直的调用栈,那如果下游处理得很慢会怎么样?答案是emit 就卡在那儿——而这就是背压。Flow 对背压的默认答案不是丢也不是缓存,是「我等你」。