位置:首页 > Kotlin > 三分钟搞懂 Kotlin Flow 中的背压

三分钟搞懂 Kotlin Flow 中的背压

时间:2026-08-23  |  作者:实验室老王  |  阅读:0

目录

  1. 背压到底在解决什么问题
  2. 默认模式:`collect` 会让发送方等待
  3. `buffer()`:给上下游之间加一个有限队列
  4. `conflate()`:忙不过来时,保留最新值
  5. `collectLatest {}`:新值一到,立即取消旧任务
  6. 怎么选:先问清楚你要保留什么

前言

在 Kotlin Flow 的使用中,真正难的通常不是把数据“发出来”,而是上游很快、下游很慢时该怎么处理。本文会围绕背压这个核心问题,拆开讲清默认收集、加缓冲、跳过旧值和取消旧任务四种策略的实际行为差异,帮助你根据是否要保留每条数据、是否只关心最新结果来做选择。

在 Kotlin Flow 里,所谓“背压”问题,本质上就是发送速度和处理速度不匹配:上游发得太快,下游来不及消费,轻则浪费算力,重则带来卡顿和内存压力。理解这一点后,再看 collectbuffer()conflate()collectLatest {} 的区别,就不只是记 API,而是能根据业务场景判断该保留每条数据、做有限缓冲,还是只追最新结果。

背压到底在解决什么问题

背压(Backpressure) 的作用,是防止快速的数据发送方把较慢的接收方压垮。没有背压时,系统要么在内存里堆积尚未处理的数据,要么花时间处理已经失去意义的旧数据。

这个问题在 UI 场景里尤其常见。以 Android 为例,屏幕刷新率通常是 60fps、90fps,部分设备可达到 120fps。如果某个 Flow 每秒发送超过 200 次更新,但界面最终只能按屏幕刷新节奏渲染,那么其中相当一部分数据其实不会真正显示出来。

也就是说,接收方很多时候并不需要“越快越多”,而是需要“足够快且稳定”。这时引入背压,主要有三点收益:

  • 控制内存使用
  • 避免无意义的重复处理
  • 让应用吞吐和响应更稳定

默认模式:`collect` 会让发送方等待

如果你没有额外添加缓冲或丢弃策略,Flow 默认就是最保守的串行处理模式:上游发一个,下游处理一个;下游没处理完,上游就先停下来。

suspend fun main() {
    flow {
        (1..3).forEach {
            timeLog("Send $it")
            emit(it)
            delay(100)            // 快速发送方
        }
    }.collect { value ->
        timeLog("Processing $value")
        delay(300)             // 慢速处理方
    }
}

// output:
// 20:22:19.647 Send 1
// 20:22:19.701 Processing 1
// 20:22:20.142 Send 2
// 20:22:20.143 Processing 2
// 20:22:20.564 Send 3
// 20:22:20.564 Processing 3

这里的关键点是:emit 会挂起,直到 collect 把上一个值处理完成。整个流程中没有额外队列,因此每个值都会被完整消费,而且顺序严格一致。

这种模式适合必须逐条处理、不能跳过的场景,例如按顺序写数据库、执行明确的业务流水,或者任何“每一条结果都必须保留”的逻辑。

`buffer()`:给上下游之间加一个有限队列

如果你的问题不是“旧数据毫无意义”,而是“上游偶尔会短时间发得更快”,那么更合适的做法通常是加缓冲,而不是直接丢数据。

flow {
//...
}
.buffer(capacity = 2)
.collect { value ->
//...    
}

// output:
// 20:27:03.368 Send 1
// 20:27:03.418 Processing 1
// 20:27:03.529 Send 2
// 20:27:03.639 Send 3
// 20:27:03.732 Processing 2
// 20:27:04.470 Processing 3

加上 .buffer(capacity = 2) 之后,上游最多可以先放入 2 个尚未处理的元素。这样做的直接效果,是发送方不必每次都卡在消费方后面,整体吞吐会更平滑,示例里也能看到总耗时缩短了。

它解决的是哪类问题

  • 你仍然需要处理每一个元素
  • 只是想吸收短时间的速度波动
  • 希望上游和下游有一定程度的并行

要注意什么

buffer() 不是无限提速。队列一旦满了,上游仍然会被挂起。所以它本质上是“有限缓冲”,不是“无上限堆积”。如果下游长期处理不过来,最终还是会回到等待状态。

conflate 与 collectLatest 的差异时间线图
`conflate()` 和 `collec这两种策略最容易混淆,时间线图能直接展示一个是跳过旧值,一个是取消正在执行的旧任务。

`conflate()`:忙不过来时,保留最新值

有些数据天生就不需要逐条处理,例如进度条刷新、传感器状态、滚动位置、界面上的实时数值展示。这类场景更关心“现在最新是什么”,而不是“中间每一步都发生过什么”。

flow { ... }
.conflate()
.collect { value -> ... }

// output:
// 20:30:45.745 Send 1
// 20:30:45.813 Processing 1
// 20:30:45.930 Send 2
// 20:30:46.039 Send 3
// 20:30:46.133 Processing 3

从输出可以看出,2 被跳过了,处理方在忙的时候不会排队保存所有旧值,而是只保留最新的未处理项。等消费者空下来时,直接拿最新值继续处理。

和 `buffer()` 的本质区别

buffer() 会保留所有元素,只是暂存起来;conflate() 则会丢弃中间的旧元素,只保证你尽快追上最新状态。

注意conflate() 不会打断当前已经开始的处理逻辑。它做的是“下一次取值时跳过旧值”,而不是中途取消当前任务。

`collectLatest {}`:新值一到,立即取消旧任务

如果你的需求比 conflate() 更激进,不只是跳过未开始处理的旧值,而是连正在进行中的旧任务也应该放弃,那么就该用 collectLatest {}

suspend fun main() {
    flow {
        (1..3).forEach {
            timeLog("Send $it")
            emit(it)
            delay(100)            
        }
    }.collectLatest { value ->
        timeLog("Start process $value")
        delay(300)             
        timeLog("Complete process $value")
    }
}

// output:
// 09:31:32.919 Send 1
// 09:31:32.979 Start process 1
// 09:31:33.090 Send 2
// 09:31:33.093 Start process 2
// 09:31:33.195 Send 3
// 09:31:33.195 Start process 3
// 09:31:33.498 Complete process 3

示例中只有 3 真正完成处理,原因是每当新的 emit 到来,前一个值对应的处理代码块就会被取消,处理逻辑会立刻切换到最新值。

它最适合什么场景

典型例子就是“边输入边搜索”。用户继续输入时,前一次搜索结果通常已经没有意义,此时继续等待旧请求完成只会浪费时间和资源,直接取消更合理。

和 `conflate()` 不要混淆

  • conflate():不取消当前任务,只跳过还没处理的旧值
  • collectLatest {}:连当前正在执行的旧任务也一起取消

怎么选:先问清楚你要保留什么

这几种方式没有绝对优劣,关键在于你到底想保留“每条数据”、 “短时缓冲能力”,还是“最新结果”。可以按下面的思路快速判断:

Kotlin Flow 不同背压策略的处理路径对比图
Kotlin Flow 背压策略一图看懂把默认收集、缓冲、保留最新值和取消旧任务四种策略放在同一张图里。
  • collect
    作用:发送方和处理方逐个等待,一一对应
    适用场景:你必须按顺序处理 每一个
  • buffer
    作用:创建一个容量为 n 的小队列,不丢弃元素
    适用场景:你需要缓冲突发流量,但仍然要处理全部元素
  • conflate
    作用:消费者忙时,只保留 最新的 元素
    适用场景:你关心最新状态,但希望当前任务自然执行完
  • collectLatest
    作用:新数据一到,立即取消旧任务
    适用场景:只有 最新结果 有意义,旧任务应尽快放弃

实际开发里,可以用四个问题做最后确认:

  1. 我是不是必须处理每一个值?
  2. 短时间的积压能不能通过一个小队列吸收?
  3. 中间值是否可以跳过,只保留最新状态?
  4. 新值出现时,旧任务是否应该立刻取消?

把这四件事想清楚,Kotlin Flow 的背压策略通常就不难选了。

免责声明:文中图文均来自网络,如有侵权请联系删除,心愿游戏发布此文仅为传递信息,不代表心愿游戏认同其观点或证实其描述。

相关文章

更多

精选合集

更多

大家都在玩

热门话题

大家都在看

更多