引言
在 Android 上,Flow 承担的是「数据从哪来、怎么变、怎么到界面」这条管道的全部职责:Room 查询返回 Flow<List<T>>,Retrofit 结果被包成流,界面用 collectAsStateWithLifecycle() 订阅,中间靠操作符做去抖、合并与重试。它比 LiveData 表达力强得多,也比手写回调清晰得多。
问题出在语义细节上。Flow 是冷流,每个收集者都会触发一次独立的生产过程,把它当热流用会重复请求;StateFlow 会做去重与合并,用它承载「一次性导航事件」会丢事件;SharedFlow 的缓冲区满了以后默认挂起生产者,配错 onBufferOverflow 就会把上游卡死。这些都不是 API 记不住的问题,而是对「冷热」「背压」「取消」三条语义理解不到位。
本文按「语义 → 操作符 → 热流 → 并发 → 集成」的顺序展开。协程的调度与结构化并发基础不在这里重复,需要补课可以看 Kotlin 协程与结构化并发 ,本文只讲数据流本身的行为。
目录
- 冷流语义与构造方式
- flowOn 与上下文规则
- 常用操作符
- 背压与缓冲策略
- StateFlow 与 SharedFlow 选型
- shareIn 与 stateIn 的启动策略
- flatMapLatest 与并发控制
- collectLatest 与生命周期感知收集
- 异常处理与重试
- 与 Room、Retrofit 的集成
- 测试与调试
1. 冷流语义与构造方式
Flow<T> 只是一个「可以按需生成多个值的挂起函数集合」。它本身不持有数据,每次 collect 都会重新执行一遍上游代码。
fun ticker(): Flow<Int> = flow {
var i = 0
while (true) {
emit(i++) // emit 是挂起函数,会随收集者节奏让步
delay(1000)
}
}
// 两个收集者 → 两份独立的计时器
launch { ticker().collect { log("A: $it") } }
launch { ticker().collect { log("B: $it") } }
构造方式有四种,选择依据是「数据源是什么形态」:
| 构造器 | 适用场景 | 关键点 |
|---|---|---|
flow { emit(...) } | 冷流、按需生成 | 块内不得切换上下文 |
flowOf(a, b, c) / asFlow() | 固定序列、集合转换 | 最简单的热身前练习 |
channelFlow { send(...) } | 需要并发发送 | 支持多协程并发 send |
callbackFlow { ... awaitClose { } } | 回调式 API 适配 | 必须 awaitClose 注销监听 |
callbackFlow 是最容易写错的一个,因为 awaitClose 不是可选项:
fun locationUpdates(client: LocationClient): Flow<Location> = callbackFlow {
val listener = LocationListener { location -> trySend(location) }
client.register(listener)
awaitClose { client.unregister(listener) } // 少了这行,流永远不结束、监听器泄漏
}
awaitClose 保证收集取消时执行清理;缺了它,callbackFlow 会立刻抛出 IllegalStateException 或者让协程永远挂起。
2. flowOn 与上下文规则
Flow 有一条硬规则:flow { } 的块内不允许直接切换上下文,也就是不能在块里写 withContext(Dispatchers.IO) 再 emit,否则会抛 Flow invariant is violated。原因是这样会破坏「发射与收集在同一协程」的保证,导致缓冲区语义错乱。
切换上游执行线程的唯一正确方式是 flowOn:
fun readFromDisk(): Flow<String> = flow {
emit(file.readText()) // 在 flowOn 指定的调度器上执行
}.flowOn(Dispatchers.IO)
flowOn 的语义边界要记牢:它只影响上游(flowOn 之前的操作符与发射端),对下游操作符无效,collect 始终在收集者的上下文里执行;连续写多个 flowOn 时只有最后一个对上游生效。
flow { emit(load()) }
.map { parse(it) } // 在 IO 上
.flowOn(Dispatchers.IO)
.map { render(it) } // 回到收集者上下文(通常是 Main)
这与 withContext 在协程里的行为完全不同:withContext 会切换并恢复,flowOn 只改变上游的上下文且不阻塞收集者。
3. 常用操作符
操作符按职责可以分成四组,记住分组比背 API 有用。
| 分组 | 代表操作符 | 作用 |
|---|---|---|
| 转换 | map、filter、transform、scan | 一对一或一对多改写 |
| 组合 | zip、combine、merge | 多流合成一条 |
| 限制 | take、debounce、sample、distinctUntilChanged | 控制发射节奏 |
| 生命周期 | onStart、onEach、onCompletion、catch | 插入副作用 |
其中最容易混淆的是 zip 与 combine:
// zip:严格配对,等两个流都有新值才发射,发射次数 = min(上游次数)
flowOf(1, 2, 3).zip(flowOf("a", "b")) { i, s -> "$i$s" } // 1a, 2b
// combine:任一新值都触发,发射次数 = max(上游次数)
flowOf(1, 2, 3).combine(flowOf("a", "b")) { i, s -> "$i$s" } // 1a, 2a, 3a, 3b(示意)
combine 的第一次发射要等所有上游都产生过至少一个值;zip 则是一一配对,任一上游耗尽即结束。表单校验、多源状态聚合用 combine,两个接口结果合并用 zip。transform 是 map 与 filter 的通用形式,允许在一次调用中发射零到多个值(不发射即相当于过滤)。
4. 背压与缓冲策略
Flow 的天然背压机制是「挂起」:emit 是挂起函数,收集者处理慢时生产者会等。这在多数情况下够用,但有两种场景需要显式干预——生产端有并发(channelFlow、buffer 之后)时挂起会造成整体吞吐下降;消费端只想看最新值时等待毫无意义。
flow {
repeat(1000) { emit(it) }
}
.buffer(capacity = 64) // 生产与消费解耦,生产端可以跑在前面
.collect { slowProcess(it) }
四种典型策略:
| 策略 | 写法 | 语义 | 适用 |
|---|---|---|---|
| 默认挂起 | 不加操作符 | 生产者等消费者 | 数据不能丢 |
| 缓冲 | buffer(64) | 队列缓冲,满了再挂起 | 吞吐优先 |
| 合并 | conflate() | 只保留最新值,丢弃中间值 | UI 状态刷新 |
| 取最新 | collectLatest { } | 新值到来时取消上一次处理 | 搜索、渲染 |
buffer 的默认容量是 64(Channel.BUFFERED)。注意一个细节:buffer 会把上游放进独立的协程执行,因此上游的 flowOn 与 buffer 顺序会影响实际线程——flowOn 之后 buffer 意味着缓冲在收集者线程,反之亦然。
conflate() 与 collectLatest 的区别在于「丢的是谁」:conflate 丢的是未处理的值,collectLatest 丢的是未完成的处理。滚动埋点这类「只关心最新位置」的场景用 conflate 更省资源。
5. StateFlow 与 SharedFlow 选型
热流的核心特征是「生产独立于消费」。StateFlow 与 SharedFlow 是两种热流,选型错误是丢事件的常见根因。
| 维度 | StateFlow | SharedFlow |
|---|---|---|
| 当前值 | 有(value 属性) | 无 |
| 初始值 | 必须提供 | 不需要 |
| 去重 | 相同值不重复发射 | 不去重 |
| 缓冲 | 恒定合并,只保留最新 | 由 replay 与 extraBufferCapacity 决定 |
| 订阅时行为 | 立即收到当前值 | 只收到订阅后的发射(受 replay 影响) |
| 典型用途 | UI 状态 | 一次性事件、广播 |
StateFlow 的去重语义来自 equals:写入一个与当前值结构相等的新值不会触发下游。这既是优化也是陷阱——把可变对象写进 MutableStateFlow 后原地修改,equals 返回 true,界面就不会更新。
class CartViewModel : ViewModel() {
private val _state = MutableStateFlow(CartState())
val state: StateFlow<CartState> = _state.asStateFlow() // 对外只读
fun add(item: Item) = _state.update { old -> // update 是原子操作
old.copy(items = old.items + item)
}
}
update { } 基于 compare-and-set 重试,是并发写入的唯一正确姿势;_state.value = _state.value.copy(...) 在多线程下会丢更新。
一次性事件(Toast、导航、Snackbar)不应该用 StateFlow,因为合并语义会丢掉连续两次相同事件。正确做法是用 MutableSharedFlow<UiEvent>(replay = 0, extraBufferCapacity = 1, onBufferOverflow = BufferOverflow.DROP_OLDEST) 并对外暴露 asSharedFlow();也可以直接用 Channel<UiEvent>(Channel.BUFFERED) 在收集端 receiveAsFlow(),它的语义更贴近「队列」,代价是只能被一个收集者消费。
6. shareIn 与 stateIn 的启动策略
把冷流变成热流的目的是「多个订阅者共享同一次上游执行」。Room 查询直接暴露给多个页面时,每个页面各查一次数据库显然是浪费。
val articles: StateFlow<List<Article>> = repository.observeAll()
.stateIn(
scope = viewModelScope,
started = SharingStarted.WhileSubscribed(5_000),
initialValue = emptyList(),
)
started 的三种取值决定了上游何时启动、何时停止:
| 取值 | 启动时机 | 停止时机 | 适用 |
|---|---|---|---|
Eagerly | 立即 | 永不 | 必须在后台持续更新的状态 |
Lazily | 首个订阅者 | 永不 | 一次性加载后长期缓存 |
WhileSubscribed(5_000) | 首个订阅者 | 无订阅 5 秒后 | Android 界面(推荐) |
WhileSubscribed(5_000) 里的 5 秒是给配置变更留的缓冲:旋转屏幕时订阅会短暂断开,若立刻停上游就会重新发一次请求,5 秒窗口把这个抖动吸收掉。这是官方在 Android 上的推荐默认值。
shareIn 与 stateIn 的差别只有两点:shareIn 返回 SharedFlow,用 replay 控制重放;stateIn 返回 StateFlow,必须给 initialValue。对 UI 状态用 stateIn,对事件广播用 shareIn。
7. flatMapLatest 与并发控制
flatMap 系列决定了「内层流如何并发」,是三兄弟里最需要理解的:
| 操作符 | 并发语义 | 典型场景 |
|---|---|---|
flatMapConcat | 串行,前一个完成才订阅下一个 | 有序上传、分页拉取 |
flatMapMerge(concurrency = N) | 并发 N 个 | 批量请求合并结果 |
flatMapLatest | 新值到来时取消上一个内层流 | 搜索、按 id 切换详情 |
搜索框是 flatMapLatest 的标准场景:
@OptIn(FlowPreview::class, ExperimentalCoroutinesApi::class)
val results: StateFlow<List<Article>> = query
.debounce(300) // 停止输入 300ms 后才算一次有效输入
.distinctUntilChanged() // 内容没变不重复请求
.flatMapLatest { q -> repository.search(q) } // 新查询取消旧查询
.stateIn(viewModelScope, SharingStarted.WhileSubscribed(5_000), emptyList())
debounce 必须在 flatMapLatest 之前:前者压缩输入频率,后者保证只有最后一次的结果被采用。顺序颠倒的话,每个字符都会触发一次请求再被取消,等于没做去抖。
并发请求需要控制上限时用 flatMapMerge,它不会因为某个内层流慢而阻塞其他流:
ids.asFlow()
.flatMapMerge(concurrency = 4) { id -> flow { emit(api.detail(id)) } }
.toList()
需要注意:flatMapLatest 取消的是内层流的协程,如果内层流内部有不可取消的操作(如 withContext(NonCancellable) 或阻塞 IO),取消不会立刻生效。
8. collectLatest 与生命周期感知收集
collectLatest 在收集端做与 flatMapLatest 对称的事:新值到来时取消上一次 collect 块的执行。
viewModel.state
.collectLatest { state ->
render(state) // 如果 render 是挂起函数,新状态到来会取消它
}
关键是「取消的是块内的挂起调用」。collectLatest 里若只有同步代码,取消无从生效;块内必须存在挂起点(delay、withContext、awaitXxx)才有意义。用它做「快速切换列表项时的详情加载」比在 UI 层手动维护 job 干净得多。
Android 上收集 Flow 必须与生命周期绑定,否则页面销毁后收集仍在运行,造成泄漏或「更新已销毁的 View」。唯一推荐写法是 repeatOnLifecycle:
viewLifecycleOwner.lifecycleScope.launch {
viewLifecycleOwner.repeatOnLifecycle(Lifecycle.State.STARTED) {
viewModel.state.collect { render(it) }
}
}
repeatOnLifecycle(STARTED) 会在进入 STARTED 时启动收集、退到 STOPPED 时取消,重新可见时再启动,天然符合「后台不更新 UI」的要求。在 Compose 中对应的是 collectAsStateWithLifecycle():
val state by viewModel.state.collectAsStateWithLifecycle()
它内部就是 repeatOnLifecycle + produceState,依赖 androidx.lifecycle:lifecycle-runtime-compose:2.8.7。相比直接 collectAsState(),它能在应用退到后台时停止收集,省电也省内存——订阅位置与重组范围的关系可参考 Compose 状态管理与重组优化
。
9. 异常处理与重试
Flow 的异常处理有两个必须记住的边界:catch 只捕获上游异常,且不会捕获由自身协程取消产生的 CancellationException。
repository.observeAll()
.map { it.filterNot(Article::hidden) } // 这里的异常能被 catch 到
.catch { e ->
log(e)
emit(emptyList()) // 可以发射兜底值
}
.collect { render(it) } // 这里的异常 catch 不到
collect 块内的异常需要用 try/catch 自己包住;catch 之后还可以继续接操作符(如 onEach、retry),但不能再回到上游。catch 里 emit 兜底值是合法且常用的写法,等价于「降级」。
重试用 retry 或 retryWhen:
repository.refresh()
.retryWhen { cause, attempt ->
if (cause is IOException && attempt < 3) {
delay(2.0.pow(attempt.toInt()).toLong() * 1000) // 指数退避
true
} else false
}
.catch { emit(RefreshResult.Failed) }
retryWhen 返回 true 表示重新订阅上游——这意味着上游的副作用(网络请求、数据库查询)会重新执行一次,幂等性要自己保证。对 stateIn 出来的热流做 retry 要格外小心:重试会重启整个上游,多个订阅者共享的缓存会被清空。
清理逻辑放 onCompletion { cause -> ... }:cause 为 null 表示正常完成,非 null 表示异常或取消,适合做资源释放与日志上报。
10. 与 Room、Retrofit 的集成
Room 从 2.6 起支持直接在 DAO 上返回 Flow,写法与具体实现见 Android 数据持久化与 Room
:
@Dao
interface ArticleDao {
@Query("SELECT * FROM article ORDER BY publishedAt DESC")
fun observeAll(): Flow<List<Article>>
@Query("SELECT * FROM article WHERE category = :category")
fun observeByCategory(category: String): Flow<List<Article>>
}
Room 的 Flow 有两个特性值得单独说:一是数据表变更时自动重新查询并发射新结果(基于 InvalidationTracker),不需要手动刷新;二是查询在 Dispatchers.IO 上执行,但收集本身要在合适的作用域里。切换查询参数时应把参数做成 Flow 再用 flatMapLatest 组合,而不是每次都新建一个 Flow:
val articles: Flow<List<Article>> = categoryFlow
.flatMapLatest { category -> dao.observeByCategory(category) }
Retrofit 的 suspend 接口本身已经是挂起函数,包成 Flow 时用 flow { emit(api.article(id)) } 即可——注意这里的 flowOn(Dispatchers.IO) 是冗余的,Retrofit 的 suspend 已经切到 IO 线程,而对 Room 的 Flow 查询与文件 IO 而言它则是必需的。网络层的拦截器、超时与重试策略属于 Retrofit 侧的话题,见 Android 网络层与 Retrofit 实践
。
需要「远端触发刷新、本地持续提供数据」的双源模式时,用 onStart 触发一次刷新,再用 map 对本地结果做加工:
fun observeFeed(): Flow<List<Article>> =
dao.observeAll()
.onStart { runCatching { api.refreshFeed() } } // 订阅时先刷新一次
.map { it.sortedByDescending(Article::publishedAt) }
11. 测试与调试
Flow 的测试核心是「虚拟时间 + 可控调度器」,runTest 会跳过 delay 的真实等待:
@Test
fun search_debounces() = runTest {
val vm = SearchViewModel(FakeRepository(), StandardTestDispatcher(testScheduler))
vm.onQueryChanged("k")
vm.onQueryChanged("ko")
vm.onQueryChanged("kot")
advanceTimeBy(301) // 推进虚拟时间越过 debounce 窗口
assertEquals(listOf("kot"), vm.queries)
}
断言流的值序列建议引入 Turbine(app.cash.turbine:turbine:1.1.0),它把「等待下一个值」变成一句 awaitItem():
@Test
fun emits_loading_then_data() = runTest {
viewModel.state.test {
assertEquals(Loading, awaitItem())
assertEquals(Data(listOf(a1)), awaitItem())
cancelAndIgnoreRemainingEvents()
}
}
test { } 会启动一个收集协程并按顺序暴露发射值,awaitItem() 在超时未发射时直接失败——这比用 toList() 收集(遇到永不结束的流会挂死)安全得多。调试时最常用的一招是在操作符链中间插 onEach { log(it) },观察每个阶段的发射频率。排查「界面不更新」时按这个顺序检查:上游是否真的发射(onEach 日志)、StateFlow 是否因 equals 去重被吞、stateIn 的 started 是否让上游根本没启动、收集是否被生命周期提前取消。
权衡取舍
| 场景 | 推荐方案 | 不推荐 | 原因 |
|---|---|---|---|
| UI 状态 | StateFlow + stateIn | SharedFlow | 需要当前值,去重避免无谓重组 |
| 一次性事件 | SharedFlow(replay = 0) 或 Channel | StateFlow | 合并语义会丢连续相同事件 |
| 事件需重放 | SharedFlow(replay = 1) | Channel | 多订阅者需要各自收到 |
| 搜索类请求 | debounce + flatMapLatest | flatMapConcat | 串行会让响应越来越滞后 |
| 批量并发请求 | flatMapMerge(concurrency) | flatMapLatest | 每个请求都需要结果 |
| 有序处理 | flatMapConcat | flatMapMerge | 顺序不能乱 |
| 高频滚动埋点 | conflate / sample | 无缓冲 | 中间值无意义 |
选型的第一步永远是问「这是状态还是事件」:状态用 StateFlow(有当前值、可重放、去重),事件用 SharedFlow 或 Channel(无当前值、不合并、可丢弃)。把事件塞进 StateFlow 是 Android 上丢事件的头号原因。
常见坑清单
| 坑 | 现象 | 规避方式 |
|---|---|---|
flow { } 内 withContext 后 emit | Flow invariant is violated | 用 flowOn 切换上游上下文 |
callbackFlow 缺 awaitClose | 流不结束、监听器泄漏 | 在 awaitClose 中注销 |
用 StateFlow 发一次性事件 | 事件偶发丢失 | 换 SharedFlow(replay = 0) 或 Channel |
MutableStateFlow 原地修改集合 | 界面不更新 | 用 update { it.copy(...) } 造新对象 |
直接赋值 _state.value = _state.value.copy() | 并发下丢更新 | 用 update { } 的 CAS 语义 |
stateIn 用 Eagerly | 后台仍在持续查询 | 用 WhileSubscribed(5_000) |
debounce 写在 flatMapLatest 之后 | 去抖失效 | 先 debounce 再 flatMapLatest |
catch 写在 collect 之后 | 捕不到下游异常 | 下游用 try/catch |
直接 collect 不绑生命周期 | 后台更新 UI、泄漏 | repeatOnLifecycle / collectAsStateWithLifecycle |
retryWhen 无上限 | 无限重试打爆接口 | 限制次数并加指数退避 |
| 把冷流当热流给多个收集者 | 重复请求、结果不一致 | shareIn / stateIn 收敛上游 |
小结
Flow 的语义可以收成四句话:
- 冷热是根本区别:冷流每次收集都重跑上游,热流生产独立于消费;多个订阅者共享上游必须
shareIn/stateIn。 - 状态与事件分开建模:
StateFlow承载状态(有当前值、去重、可重放),SharedFlow/Channel承载事件(不合并、可丢弃)。 - 并发语义决定操作符:
flatMapLatest取消旧内层流,flatMapMerge并发,flatMapConcat串行,选错就是性能或正确性问题。 - 收集必须绑生命周期:
repeatOnLifecycle是唯一正确写法,Compose 里用collectAsStateWithLifecycle()。
把这四条落到 Code Review 清单上,「重复请求、事件丢失、后台更新界面、去抖失效」这几类线上问题基本能在合入前拦住。
继续阅读
探索更多技术文章
浏览归档,发现更多关于系统设计、工具链和工程实践的内容。