Kotlin Flow 与响应式数据流

Flow 是 Android 上连接数据库、网络与界面的数据管道,但冷热流选错、背压处理不当就会写出丢事件或卡顿的代码。本文讲透冷流与热流的语义差别、StateFlow 与 SharedFlow 的选型依据、常用操作符与背压策略、flatMapLatest 的并发语义、生命周期感知收集、异常重试,以及 Room 与 Retrofit 的集成写法。

引言

在 Android 上,Flow 承担的是「数据从哪来、怎么变、怎么到界面」这条管道的全部职责:Room 查询返回 Flow<List<T>>,Retrofit 结果被包成流,界面用 collectAsStateWithLifecycle() 订阅,中间靠操作符做去抖、合并与重试。它比 LiveData 表达力强得多,也比手写回调清晰得多。

问题出在语义细节上。Flow 是冷流,每个收集者都会触发一次独立的生产过程,把它当热流用会重复请求;StateFlow 会做去重与合并,用它承载「一次性导航事件」会丢事件;SharedFlow 的缓冲区满了以后默认挂起生产者,配错 onBufferOverflow 就会把上游卡死。这些都不是 API 记不住的问题,而是对「冷热」「背压」「取消」三条语义理解不到位。

本文按「语义 → 操作符 → 热流 → 并发 → 集成」的顺序展开。协程的调度与结构化并发基础不在这里重复,需要补课可以看 Kotlin 协程与结构化并发 ,本文只讲数据流本身的行为。


目录

  1. 冷流语义与构造方式
  2. flowOn 与上下文规则
  3. 常用操作符
  4. 背压与缓冲策略
  5. StateFlow 与 SharedFlow 选型
  6. shareIn 与 stateIn 的启动策略
  7. flatMapLatest 与并发控制
  8. collectLatest 与生命周期感知收集
  9. 异常处理与重试
  10. 与 Room、Retrofit 的集成
  11. 测试与调试

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 是两种热流,选型错误是丢事件的常见根因。

维度StateFlowSharedFlow
当前值有(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 + stateInSharedFlow需要当前值,去重避免无谓重组
一次性事件SharedFlow(replay = 0) 或 ChannelStateFlow合并语义会丢连续相同事件
事件需重放SharedFlow(replay = 1)Channel多订阅者需要各自收到
搜索类请求debounce + flatMapLatestflatMapConcat串行会让响应越来越滞后
批量并发请求flatMapMerge(concurrency)flatMapLatest每个请求都需要结果
有序处理flatMapConcatflatMapMerge顺序不能乱
高频滚动埋点conflate / sample无缓冲中间值无意义

选型的第一步永远是问「这是状态还是事件」:状态用 StateFlow(有当前值、可重放、去重),事件用 SharedFlow 或 Channel(无当前值、不合并、可丢弃)。把事件塞进 StateFlow 是 Android 上丢事件的头号原因。


常见坑清单

坑现象规避方式
flow { } 内 withContext 后 emitFlow 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 的语义可以收成四句话:

  1. 冷热是根本区别:冷流每次收集都重跑上游,热流生产独立于消费;多个订阅者共享上游必须 shareIn/stateIn。
  2. 状态与事件分开建模:StateFlow 承载状态(有当前值、去重、可重放),SharedFlow/Channel 承载事件(不合并、可丢弃)。
  3. 并发语义决定操作符:flatMapLatest 取消旧内层流,flatMapMerge 并发,flatMapConcat 串行,选错就是性能或正确性问题。
  4. 收集必须绑生命周期:repeatOnLifecycle 是唯一正确写法,Compose 里用 collectAsStateWithLifecycle()。

把这四条落到 Code Review 清单上,「重复请求、事件丢失、后台更新界面、去抖失效」这几类线上问题基本能在合入前拦住。

继续阅读

探索更多技术文章

浏览归档,发现更多关于系统设计、工具链和工程实践的内容。

全部文章 返回首页

「Android 开发」更多文章

  1. Kotlin Multiplatform 跨平台共享
  2. Android R8 混淆与 Baseline Profile
  3. Android 测试体系:单元测试、Espresso 与 Compose 测试