Appearance
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)
- 一句话答:Flow 利用 Kotlin 协程的
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 Flow | Dart Stream | RxJava | JS Observable |
|---|---|---|---|---|
| 冷/热 | 冷流本体 + stateIn/shareIn 热化 | 默认单订阅,broadcast() 热化 | 冷热可配 | 冷热可配 |
| 背压 | suspend 式天然背压 | pause/resume 手动 | 背压策略(Drop/Latest/Buffer) | 无内置(需手动) |
| 上下文切换 | flowOn(只切上游) | 无内置(靠 isolate/队列) | subscribeOn/observeOn | 无内置 |
| 取消 | 结构化并发自动取消 | StreamSubscription.cancel | Disposable.dispose | unsubscribe |
常见场景
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 | 上游异常直接取消 collect | catch 兜底 + 显式处理 |
与相近概念对比
| 概念 | 本质区别 |
|---|---|
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行为差异。
复习检查题
Flow 为什么是"冷"的?冷对工程意味着什么?
答:Flow 的 producer 只在
collect时才执行,每次 collect 独立运行;意味着它是"数据生成的描述"而非"一直运行的源",适合分页 / 查询 / 请求这种按需产生的数据。Flow 的背压是怎么实现的?
答:
emit是 suspend 函数,下游 collect 处理慢时上游emit挂起等待,形成天然"同步接力"背压,不积压不丢失;需要放行才用buffer/conflate。flowOn和withContext在 Flow 里的差别?答:
flowOn只切换链条中它之前的操作符到指定上下文(上游),collect 仍在调用上下文;withContext切整个块。Flow 内切上游应使用flowOn,避免破坏上下文保序。什么时候把 Flow 热化为 StateFlow?
答:当需要"多订阅者共享一份执行 + 提供当前值"时(如 UI 状态),用
stateIn(配合SharingStarted.WhileSubscribed)把冷流热化为 StateFlow;数据源层保持冷流。Flow 和 RxJava 的核心差异是什么?
答:Flow 天然 suspend 式背压、结构化并发取消(随 scope 自动取消)、不需要额外背压策略;RxJava 默认 Observable 无内置背压易积压,且取消依赖手动 dispose。
速记
- Flow = 冷流:collect 才执行,每次独立
emit= suspend,天然背压,慢消费者不爆flowOn切上游,collect定下游buffer/conflate/collectLatest:按"能不能丢中间值"选stateIn/shareIn:冷 → 热的桥- Flow 是数据流描述,StateFlow 是状态寄存器