Skip to content

Kotlin Flow:冷流数据管道 / 背压 / flowOn ​

回到总览:Kotlin 协程并发总览
相关对照:Kotlin Scope / Job / 取消传播 / 状态流(StateFlow / SharedFlow 是热流,本页讲冷流 Flow 本体)

一句话定义 ​

Flow 是 Kotlin 协程生态里的冷流数据管道:生产者在被收集(collect)时才执行(冷),天然支持 suspend 式背压(慢消费者自动反压、不积压不丢失),flowOn 可切换上游执行上下文,再配合丰富操作符做异步数据流的组合与变换。

代码索引 ​

当前主题暂无 lab,对应实验见文末「计划补齐」。

主题Lab 说明源码
Flow 冷流 / flowOn / buffer / 操作符计划补齐—

为什么需要 ​

很多工程师会用 StateFlow / SharedFlow 做 UI 状态,但真正问到下面这些问题时就开始含糊:

  • Flow 和 StateFlow / SharedFlow 到底差在哪?
    • 一句话答:普通 Flow 是冷流(只有收集时才生产数据);StateFlow(持状态快照)和 SharedFlow(广播事件)是热流(无订阅者也在内存中流转)。(详见下文 §1、§4)
  • 为什么说 Flow 是"冷"的,冷和热对工程意味着什么?
    • 一句话答:冷流每次 collect 都会重新触发生产者执行逻辑;这意味着冷流不浪费资源,但多次收集会重复触发上游副作用。(详见下文 §1)
  • Flow 的背压是怎么实现的?为什么"慢消费者"不会爆内存?
    • 一句话答:Flow 利用 Kotlin 协程的 suspend 挂起机制实现背压;当消费者处理慢时,上游 emit 自动挂起暂停生产,无需额外的缓冲队列。(详见下文 §2)
  • flowOn 和 withContext 在 Flow 里有什么差别?
    • 一句话答:flowOn 改变的是它上游代码的执行 Dispatcher;而 Flow 严格禁止在内部直接使用 withContext 切换线程(打破上下文守恒)。(详见下文 §3)
  • Flow 和 RxJava、和 Sequence 有什么本质区别?
    • 一句话答:Sequence 是纯同步阻塞流;RxJava 依赖复杂操作符与线程池;Flow 依托协程挂起与结构化并发,更轻量且无缝契合挂起函数。(详见下文 §1)

如果说不清,工程上就很容易出现:

  • 用 StateFlow 表达"分页拉取/传感器数据/DB 查询"这种一次性流,语义拧巴
  • 以为 flowOn(IO) 会像 withContext 一样切回上游上下文,结果线程行为对不上
  • 在 Flow 里用 withContext 切线程,违反 Flow 上下文保序约定
  • 把大流量数据流直接 collect 到 UI,不懂 buffer / conflate 该何时用

底层机制 ​

1. 冷流 vs 热流:决定"什么时候执行" ​

  • 冷流(Flow) :生产者代码只在 collect 时才运行,每次 collect 都独立执行一遍(像函数调用)。
  • 热流(StateFlow / SharedFlow / Channel) :生产者独立于消费者运行(如 state 一直持有值、事件一直广播),与是否有人收集无关。
kotlin
val cold = flow {
    println("produce")          // collect 时才会打印
    emit(1)
    emit(2)
}

// 每次 collect 都重新执行 producer —— 冷
cold.collect { println("A: $it") }   // 输出 produce / A:1 / A:2
cold.collect { println("B: $it") }   // 再次输出 produce / B:1 / B:2

预期现象:Flow 更像"数据生成的函数描述",StateFlow 更像"当前值的寄存器"。分页、DB 查询、网络轮询这种"每次要新数据"的场景用冷流;"永远记住最新状态 / 广播给所有订阅者"用热流。

2. suspend 式背压:慢消费者不会爆 ​

Flow 的每个上游 emit 都是 suspend 函数:如果下游 collect 处理得慢,emit 会挂起等待,直到下游准备好——这就是"背压"(backpressure),不需要缓冲区也不会丢数据。

kotlin
flow {
    repeat(5) { i ->
        emit(i)                    // 如果下游处理慢,这里会挂起等下游
    }
}.collect { value ->
    delay(100)                     // 下游很慢
    println(value)
}

预期现象:emit 与 collect 之间是"同步接力":上游发一个,下游处理一个,处理完上游才发下一个。不会像 RxJava 的 Observable 默认那样在无背压时无限积压(这是 Flow 相对 RxJava 的最大工程优势之一)。

当你想"允许丢中间值 / 让下游快点"时,才用:

  • buffer(capacity):给管道加缓冲区,上游不必等下游
  • conflate():只保留最新值,中间值可丢(适合"进度快照"类 UI)
  • collectLatest:下游还在处理时来新值,取消旧处理(适合搜索防抖)

3. flowOn:切换"上游"执行上下文 ​

flowOn 把它之前的操作符切到指定上下文执行,且只影响上游,下游仍在 collect 所在上下文:

kotlin
flow {
    emit(fetchFromDb())            // 跑在 IO
}
.map { parse(it) }                 // 也跑在 IO(在 flowOn 之前的链上)
.flowOn(Dispatchers.IO)            // 上游(emit + map)切到 IO
.collect { render(it) }            // 下游仍在主线程 collect

对照点:withContext(IO) 切"整个块"(块内外的上下文都被改写);flowOn(IO) 切"flow 链条中它之前的那一段"。在 Flow 里不要用 withContext 切线程——它会违反 Flow 的上下文保序约定,应该用 flowOn。

4. 操作符组合:数据流的"构建器语法" ​

kotlin
flowOf(1, 2, 3)
    .map { it * 2 }                       // 1 对 1 变换
    .filter { it > 2 }                    // 过滤
    .flatMapLatest { loadDetail(it) }     // 最新值驱动异步请求(搜索场景)
    .catch { e -> emit(fallback) }        // 上游异常兜底
    .collect { println(it) }

预期现象:Flow 的操作符大多是 suspend 的(map/filter 内部可挂起),与 Sequence(同步、阻塞)有本质区别;与 RxJava 操作符语义对齐,但天然挂起式、可取消、可背压。

5. 冷 → 热桥接:stateIn / shareIn ​

冷流想共享给多个订阅者、或变成"当前值"时,用 stateIn / shareIn 热化:

kotlin
val uiState: StateFlow<UiState> = repository.dataFlow
    .stateIn(
        scope = viewModelScope,
        started = SharingStarted.WhileSubscribed(5000),
        initialValue = UiState.Loading,
    )

预期现象:stateIn 把冷流变成 StateFlow(有当前值、多订阅者共享一份执行),是"DB/网络冷流 → UI 状态热流"的标准桥接。这也解释了为什么 UI 层看到的多是 StateFlow,而数据源层是冷 Flow。

比喻(点餐 vs 看电视) :Flow(冷流)像点餐——你下单(collect)厨师才开始做(producer 执行),每次点都重新做一份(每次 collect 独立执行),厨房按你的节奏上菜(suspend 背压:你没吃完不上下一道)。StateFlow(热流)像一直播的电视节目——无论有没有人看都在播,谁打开电视看到的都是当前画面(最新值),来晚了看不到前面(历史事件不重放)。

Android / Flutter / Web / Backend 对照 ​

维度Kotlin FlowDart StreamRxJavaJS Observable
冷/热冷流本体 + stateIn/shareIn 热化默认单订阅,broadcast() 热化冷热可配冷热可配
背压suspend 式天然背压pause/resume 手动背压策略(Drop/Latest/Buffer)无内置(需手动)
上下文切换flowOn(只切上游)无内置(靠 isolate/队列)subscribeOn/observeOn无内置
取消结构化并发自动取消StreamSubscription.cancelDisposable.disposeunsubscribe

常见场景 ​

1. 分页 / 搜索 ​

flatMapLatest + 搜索关键字 Flow,自动取消上一次请求,天然防抖。

2. DB 查询 / 网络轮询转 UI 状态 ​

冷流查询 + stateIn 热化为 StateFlow,UI 层只 collect 状态,数据源层保持冷流语义。

3. 传感器 / 位置 / 大文件逐块读取 ​

高频率数据流用 buffer / conflate 控制消费节奏,避免 UI 被淹没。

常见误配、事故后果与排障 ​

坑现象修法
把"一次性数据获取"也用 StateFlow语义拧巴、初值处理复杂冷数据源用 Flow,UI 再 stateIn
Flow 里用 withContext 切线程上下文保序被破坏、行为诡异用 flowOn 切上游
大流量不处理背压下游 UI 卡顿、队列堆积buffer / conflate / collectLatest 按需选
flowOn 放在链中间只切了"它之前"的段,预期不符记住:flowOn 影响上游,collect 影响下游
忘记 catch上游异常直接取消 collectcatch 兜底 + 显式处理

与相近概念对比 ​

概念本质区别
Sequence同步、阻塞、非 suspend;Flow 是异步、suspend、可取消
RxJava Observable无内置背压易积压;Flow 天然 suspend 背压 + 结构化取消
Channel热管道、可广播可背压,偏向"通信";Flow 偏向"数据流描述"
StateFlow / SharedFlow热流:持有当前值 / 广播事件;由 Flow stateIn / shareIn 热化而来

对应实验 ​

计划补齐:

  • labs/kotlin/host/src/main/kotlin/runtime_concurrency/flow_cold_stream/README.md — 观察:冷流每次 collect 独立执行、suspend 背压(下游慢则 emit 挂起)、flowOn 只切上游、buffer / conflate 行为差异。

复习检查题 ​

  1. Flow 为什么是"冷"的?冷对工程意味着什么?

    答:Flow 的 producer 只在 collect 时才执行,每次 collect 独立运行;意味着它是"数据生成的描述"而非"一直运行的源",适合分页 / 查询 / 请求这种按需产生的数据。

  2. Flow 的背压是怎么实现的?

    答:emit 是 suspend 函数,下游 collect 处理慢时上游 emit 挂起等待,形成天然"同步接力"背压,不积压不丢失;需要放行才用 buffer / conflate。

  3. flowOn 和 withContext 在 Flow 里的差别?

    答:flowOn 只切换链条中它之前的操作符到指定上下文(上游),collect 仍在调用上下文;withContext 切整个块。Flow 内切上游应使用 flowOn,避免破坏上下文保序。

  4. 什么时候把 Flow 热化为 StateFlow?

    答:当需要"多订阅者共享一份执行 + 提供当前值"时(如 UI 状态),用 stateIn(配合 SharingStarted.WhileSubscribed)把冷流热化为 StateFlow;数据源层保持冷流。

  5. Flow 和 RxJava 的核心差异是什么?

    答:Flow 天然 suspend 式背压、结构化并发取消(随 scope 自动取消)、不需要额外背压策略;RxJava 默认 Observable 无内置背压易积压,且取消依赖手动 dispose。

速记 ​

  • Flow = 冷流:collect 才执行,每次独立
  • emit = suspend,天然背压,慢消费者不爆
  • flowOn 切上游,collect 定下游
  • buffer / conflate / collectLatest:按"能不能丢中间值"选
  • stateIn / shareIn:冷 → 热的桥
  • Flow 是数据流描述,StateFlow 是状态寄存器

站点构建时间:2026/8/24 23:43:17