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

背压:Flow 的默认答案是「我等你

上一章的结论有一个直接推论:既然 emit 只是一次普通函数调用,那下游没返回,emit 就返回不了。Flow 的背压不是一个被加进去的功能,是「什么都不加」的自然结果。

背压bufferconflatecollectLatest

先看四条时间线

场景:上游每 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 天生就是 conflated

你不需要给 StateFlowconflate()——它本来就是。一个慢的 collector 只会看到「它有空的时候的最新值」,中间的全跳过了。

这也是为什么 StateFlow 不适合做事件总线:事件是不能丢的,而它天天在丢。第 21 章会把这笔账算清楚。

collectLatest:新的来了就掐掉旧的

source.collectLatest { slow(it) }

这个和前两个的机制完全不同:它不丢值,每个值都会开始处理;但只要下一个值来了,上一个还没处理完的就被取消

结果:只有 8 号真正跑完(221 毫秒),前 7 个都处理到一半被掐了。

看 demo 的日志,会看到一连串 cancel——那是真的子协程取消,第 11 章那套机制。

适合:处理过程本身可以被安全放弃的场景。

  • 搜索框:用户还在打字,上一个关键词的请求就该取消
  • 图片预览:滑到下一张,上一张的解码没必要继续
  • 任何「只有最后一次的结果有意义」的异步操作
⚠ collectLatest 的两个前提

① 你的处理逻辑必须真的能被取消。第 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(UNLIMITED) 当万金油

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

消费者不慢的时候,这四个选择根本不存在。这是这一章最实用的一条:先量一量你的下游到底慢不慢,再决定加不加东西。

⌗ 到你手上

这一章的一句话

背压是 emit 会挂起的自然结果,不是加上去的功能。buffer 用内存换吞吐,conflate 丢中间值保最新,collectLatest 掐掉没做完的活——三种取舍,先量清楚下游到底慢不慢再选。

下一章:这一卷最后一块——热流SharedFlow 有四个旋钮,而它的默认值会让你悄悄丢事件。理解了这四个旋钮,你会发现 StateFlow 只是其中一组特定的取值。