背压:Flow 的默认答案是「我等你」
上一章的结论有一个直接推论:既然 emit 只是一次普通函数调用,那下游没返回,emit 就返回不了。Flow 的背压不是一个被加进去的功能,是「什么都不加」的自然结果。
先看四条时间线
场景:上游每 20 毫秒发一个值,一共 8 个;下游每个要处理 60 毫秒。生产快,消费慢——这是背压问题的标准形状。
默认参数下的结果:
写法 处理完的值 个数 没走到终点的 总耗时 什么都不加(默认) 1,2,3,4,5,6,7,8 8/8 0 642 ms buffer(4) 1,2,3,4,5,6,7,8 8/8 0 501 ms conflate() 1,3,6,8 4/8 4 个被丢掉 261 ms collectLatest 8 1/8 7 个被掐掉 221 ms
四行是四种完全不同的取舍,没有哪个是「更好的」。先看第一行为什么是那样。
为什么默认是「等」
回想上一章那条调用栈:emit(v) 直接调用下游的处理块。所以:
t=20 上游想 emit(1) → 调用 collect 块 → 它要跑 60 ms t=80 collect 块返回 → emit(1) 才返回 → 上游才能继续 t=100 上游 emit(2)(本该在 t=40 发的,晚了 60 ms) …
上游被下游拖着走。总耗时 642 毫秒,而不是「生产 160 毫秒 + 最后一个处理 60 毫秒」。
这个行为有个名字叫背压(backpressure):下游的压力反向传导给上游。RxJava 为了做这件事专门造了 Flowable 和一整套请求策略;Flow 什么都没做,它天然就是这样。
Flow 的背压不是功能,是 suspend 的副产品。
emit 是一个挂起函数,它调用的下游也可能挂起。挂起意味着「我还没完成」,于是上游自然就等在那儿了。
而剩下三种写法,做的都是同一件事:在中间插一个东西,把上下游解耦。它们的差别只在于——中间那个东西满了怎么办。
三种解耦,三种态度
buffer(n):不丢,用内存换吞吐
source.buffer(4).collect { slow(it) }
buffer 开一个子协程去 collect 上游,中间架一个容量 4 的 Channel。上游往里塞,下游从里取,两边并行跑。
结果:总耗时从 642 降到 501 毫秒,一个值都不丢。
但也别指望它能降到「生产时间」——总耗时仍然被消费速度决定(8 × 60 = 480 毫秒是下限)。buffer 能消除的是「等待」,消除不了「慢」。缓冲满了之后,上游照样会被卡住。
适合:每个值都必须处理,而且生产是一阵一阵的。比如读文件的行、处理一批上传任务。
conflate():只要最新的
source.conflate().collect { render(it) }
conflate() 就是 buffer(1, onBufferOverflow = DROP_OLDEST)。缓冲里永远只留最新的一个,下游一忙完就拿最新的那个。
结果:只处理了 4 个(1、3、6、8),丢了 4 个,总耗时 261 毫秒。
注意最后一个值 8 一定会被处理——这是 conflate 很重要的性质:它丢中间值,但不丢最终值。
适合:UI 状态。用户只关心「现在是什么样」,不关心中间经过了哪些状态。滚动位置、进度条、传感器读数、股票价格——全都是这一类。
你不需要给 StateFlow 加 conflate()——它本来就是。一个慢的 collector 只会看到「它有空的时候的最新值」,中间的全跳过了。
这也是为什么 StateFlow 不适合做事件总线:事件是不能丢的,而它天天在丢。第 21 章会把这笔账算清楚。
collectLatest:新的来了就掐掉旧的
source.collectLatest { slow(it) }
这个和前两个的机制完全不同:它不丢值,每个值都会开始处理;但只要下一个值来了,上一个还没处理完的就被取消。
结果:只有 8 号真正跑完(221 毫秒),前 7 个都处理到一半被掐了。
看 demo 的日志,会看到一连串 cancel——那是真的子协程取消,第 11 章那套机制。
适合:处理过程本身可以被安全放弃的场景。
- 搜索框:用户还在打字,上一个关键词的请求就该取消
- 图片预览:滑到下一张,上一张的解码没必要继续
- 任何「只有最后一次的结果有意义」的异步操作
① 你的处理逻辑必须真的能被取消。第 11 章那条规矩在这儿一模一样适用:如果 slow(it) 是一段纯 CPU 循环,collectLatest 一个都取消不掉,它会退化成「一个一个排队跑完」,而且比什么都不加还慢(多了协程开销)。
② 处理过程不能有不可回退的副作用。写数据库写到一半被取消,你就有了一条半截的记录。这种时候要么用 NonCancellable 保护关键段,要么根本别用 collectLatest。
一张选择表
| 你的情况 | 用 | 代价 |
|---|---|---|
| 下游其实不慢 | 什么都不加 | 零。别提前优化 |
| 每个值都必须处理完 | buffer(n) | 内存。n 要有上限,别写 UNLIMITED |
| 只关心最新状态 | conflate() | 中间值丢了 |
| 只有最后一次有意义,而且处理能被取消 | collectLatest | 处理到一半的工作白做了 |
| 用户在快速输入,想等他停下来 | debounce(t) | 多了 t 毫秒延迟 |
| 想定期采样,不管发了多少 | sample(t) | 可能漏掉最后一个 |
最后两行是「按时间」而不是「按能力」来减流量的,经常和上面几个搭配用:
// 搜索框的标准写法:三个操作符各管一件事
searchQueryFlow
.debounce(300) // 等他打完
.distinctUntilChanged() // 一样的词不重复搜
.flatMapLatest { q -> repo.search(q) } // 新词来了就取消上一个请求
.collect { results = it }
这三行是 Flow 在 Android 上最有说服力的一段代码:换成回调式写法,「防抖 + 去重 + 取消上一个请求」大概要写四五十行,还容易出竞态。
buffer(Channel.UNLIMITED) 确实能让上游永远不被卡住。代价是:如果上游持续比下游快,缓冲会无限增长,直到 OOM。
而且它把问题藏起来了——本来上游被卡住是一个明确的信号(「你的下游太慢」),无限缓冲让这个信号消失了,然后在几分钟后变成一次内存崩溃。
一个有界的 buffer 加上明确的溢出策略,永远好过无界缓冲。
把滑杆拉到另一头
demo 上面有两个滑杆。把「每个处理要」拖到比「生产间隔」还小,你会看到四行数字变得几乎一样:
生产间隔 20 ms,处理 10 ms: 什么都不加 8/8 约 240 ms buffer(4) 8/8 约 240 ms conflate() 8/8 约 240 ms collectLatest 8/8 约 240 ms
消费者不慢的时候,这四个选择根本不存在。这是这一章最实用的一条:先量一量你的下游到底慢不慢,再决定加不加东西。
怎么判断自己有没有背压问题:
1. 在 collect 块的开头和结尾各打一个时间戳,看单次耗时
2. 在 emit 前后各打一个,看上游有没有被拖慢
两个时间差得多 → 有背压
常见的三个真实场景:
· 传感器 / 位置回调(每秒几十次)→ conflate
· 上传队列(每个几秒) → buffer + 有界
· 搜索 / 联想 → debounce + flatMapLatest
一条经验:凡是最终目的地是 UI 的流,几乎都该 conflate
(而如果你已经在用 StateFlow,它自带了)。
这一章的一句话
背压是 emit 会挂起的自然结果,不是加上去的功能。buffer 用内存换吞吐,conflate 丢中间值保最新,collectLatest 掐掉没做完的活——三种取舍,先量清楚下游到底慢不慢再选。
下一章:这一卷最后一块——热流。SharedFlow 有四个旋钮,而它的默认值会让你悄悄丢事件。理解了这四个旋钮,你会发现 StateFlow 只是其中一组特定的取值。