Kotlin 协程与 Flow ★
这是你的头号短板, 也是中级面试的高频考点. 协程是现代 Android 异步编程的基石, 几乎必问. 本篇从原理到实战到面试题完整覆盖.
前置知识与示例说明: 先掌握 Kotlin 空安全, lambda 和异常 (见 Kotlin 语言核心).除明确标注 “可运行示例” 的代码外, 本文为 Android / 协程依赖下的上下文片段; 其中 api, repo, render 等名称由所在工程提供. 伪代码只用于解释状态机, 不能直接编译.
学习目标: 能画出一个请求的挂起 / 恢复与 Job 父子关系, 能在 ViewModel 中把冷流转换为可恢复的 UI 状态, 并能说明取消, 异常和缓冲的边界.
一, 协程是什么 (先建立心智模型)
协程不是线程. 它是一种可挂起 / 恢复的计算, 由编译器 + 运行时在用户态调度.
- 挂起 (suspend): 不阻塞线程. 协程挂起时, 它所在的线程被释放去做别的事, 等条件满足再恢复执行.
- 核心价值: 用同步顺序的写法表达异步逻辑, 消灭回调地狱.
// 回调写法
api.login(user) { token ->
api.getProfile(token) { profile -> updateUI(profile) }
}
// 协程写法
val token = api.login(user) // suspend,挂起不阻塞
val profile = api.getProfile(token) // suspend
updateUI(profile)
二, suspend 原理 (高频追问: 协程为什么不阻塞线程?)
suspend 函数由 Kotlin 编译器做 CPS 变换 (Continuation-Passing Style):
- 编译器给每个 suspend 函数隐式加一个参数
Continuation(回调), 返回值变成Any?. - 函数体被编译成一个状态机: 每个挂起点是一个状态.
- 挂起时函数返回特殊标记
COROUTINE_SUSPENDED, 线程被释放. - 当结果就绪, 调用
continuation.resumeWith(result), 状态机从上次挂起点恢复执行.
所以协程 “挂起不阻塞” 的本质:把后续代码包装成回调, 挂起时退出函数让出线程, 恢复时再回来. 这是面试最爱的深挖点.
一次请求可以按下列顺序追踪. 这里的 “恢复” 由挂起 API 和 dispatcher 决定, 不能假定回到同一物理线程:
viewModelScope.launch (Job A, Main)
-> api.fetch() 发起异步 I/O
-> 返回 COROUTINE_SUSPENDED,Main 可处理其他消息
-> I/O 完成,Continuation 恢复 Job A
-> dispatcher 安排后续代码,更新 _uiState
// 你写的:
suspend fun login(user: User): Token
// 编译后近似:
fun login(user: User, cont: Continuation<Token>): Any?
三, 结构化并发 (Structured Concurrency)
协程必须在 CoroutineScope 中启动. 核心:协程有父子层级, 父协程等待所有子协程完成, 父被取消则子全部取消, 避免协程泄漏.
- CoroutineScope: 协程作用域, 持有
CoroutineContext. - CoroutineContext: 元素集合, 关键有
Job(生命周期),CoroutineDispatcher(线程),CoroutineName,CoroutineExceptionHandler. - Job: 协程生命周期句柄, 可
cancel()/join(), 状态有 New/Active/Completing/Cancelling/Cancelled/Completed. - Job vs Deferred:
Deferred<T>是携带结果的 Job, 继承cancel()/join()并增加await();async返回Deferred. 取消一个async的任务后,await()会抛出CancellationException; 若该Deferred已成功完成, 取消不再影响其结果. - 父子关系: 协程内启动的子协程, 其 Job 是父 Job 的子节点.
Android 常用现成 scope:
viewModelScope: 绑定 ViewModel, onCleared 自动取消.lifecycleScope: 绑定 Lifecycle.GlobalScope:不推荐, 脱离结构化并发, 易泄漏.
父子 Job 图 (伪图) 应当能在面试时画出:
ViewModel.onCleared()
└─ viewModelScope 的 SupervisorJob
└─ launch load()(子 Job)
├─ async loadProfile()(子 Job)
└─ async loadFeed()(子 Job)
viewModelScope 可理解为 SupervisorJob, Dispatchers.Main.immediate 和 ViewModel 清理回调组合出的作用域, 具体实现随 lifecycle 版本演进. 父 scope 取消会向下取消所有子 Job; 其中的 supervisor 语义避免一个直接子任务失败自动取消其他直接子任务, 但你仍要在各任务处把错误归约为明确状态. 不要把 “使用 viewModelScope” 误解为所有异常都会自动显示给用户.
coroutineScope 与 supervisorScope 都等块内子协程全部完成才返回, 区别在失败传播: coroutineScope 下任一子协程失败会取消整个作用域 (包括所有兄弟) 并向外抛出异常; supervisorScope 隔离失败, 直接子协程的异常不会取消兄弟或自身, 但该异常仍要由对应子协程处理或 await() 观察. 选型: 多个步骤是同一事务, 失败即整体回滚 → coroutineScope; 多个相互独立的卡片 / 接口并行加载, 一个失败不拖垮其他 → supervisorScope.
class MyViewModel : ViewModel() {
fun load() = viewModelScope.launch { // ViewModel 销毁自动取消
val data = repo.fetch() // suspend
_state.value = data
}
}
四, Dispatchers 与 withContext
| Dispatcher | 用途 | 线程池 |
|---|---|---|
Dispatchers.Main | UI 操作 | 主线程 |
Dispatchers.IO | 网络 / 磁盘 IO | 共享弹性调度器; 默认并行度至少为 64 或可用处理器数 (取较大者) |
Dispatchers.Default | CPU 密集 (排序 / 解析) | 核数大小 |
Dispatchers.Unconfined | 不限定 (测试 / 特殊) | 当前线程, 恢复后随挂起点 |
viewModelScope.launch { // Main
val data = withContext(Dispatchers.IO) { // 切到 IO
api.fetch()
}
textView.text = data // 自动回到 Main
}
withContext 切换上下文并挂起等待结果返回, 是最常用的线程切换方式. 当前 kotlinx.coroutines 文档的 Dispatchers.IO 默认并行度为 “至少 64 或可用处理器数” 的较大值, 仍可由系统属性配置; IO.limitedParallelism(n) 创建的弹性视图可拥有独立并行度配额, 因此不能把总线程数简单等同于表中的默认值. IO 与 Default 共享底层线程资源, 两者间切换不一定真正换线程; 版本升级时应以目标 kotlinx.coroutines 文档和实际负载重新核验. 其中 “64” 是默认值, 随版本调整.
limitedParallelism 与 Semaphore 隔离层次不同: 前者在调度器层给视图划独立并发配额, 该视图上的任务最多 n 个同时执行, 不会与他处代码互相挤占线程; Semaphore 只在应用逻辑层限制在途任务数, 底层仍共享同一 IO 池, 不预留资源. 需要 “为某模块独立线程预算” 用 limitedParallelism, 需要 “限制打某下游的并发请求数” 用 Semaphore.
五, launch vs async
launch: 返回Job, 不返回结果, 适合副作用; 未处理异常会按父 Job / 根协程规则传播.async: 返回Deferred<T>,await()用于取得结果; 失败会存为 Deferred 的完成结果,await()在观察它时重新抛出该异常.
// 并发请求,总耗时 = max 而非 sum
val a = async { api.getA() }
val b = async { api.getB() }
val result = a.await() + b.await()
await() 不是异常传播的唯一开关. 普通父 Job 中, 一个子 async 失败会立即取消父及兄弟, 即使尚未 await(); 随后 await 只是在调用点重新抛出该失败. supervisorScope 中失败的 async 不会取消兄弟或 supervisor, 但其异常仍保存在对应 Deferred, 调用方必须 await()/检查或以其他方式处理. 根 async(没有结构化父 Job, 例如不推荐的独立 scope) 若既没有 await 又没有显式观察 / 处理, 异常可能只停留在 Deferred, 因而容易被遗漏; 这不是 “所有不 await 都一定静默”.
// 上下文片段:三种异常边界,省略 dispatcher 与错误呈现.
coroutineScope {
async { error("ordinary child failure") } // 立即取消本 coroutineScope
}
supervisorScope {
val independent = async { error("independent failure") }
// 兄弟可继续;此处必须在合适位置 independent.await() 并处理异常.
}
六, 取消机制 (协作式)
协程取消是协作式的: 取消只是把 Job 标记为 Cancelling, 真正停止需要协程代码配合检查.
- 可取消的 suspend 挂起点 (delay 等) 在恢复时会检查取消状态, 抛出
CancellationException;withContext(NonCancellable)显式排除或未真正挂起即返回的挂起点除外. - 纯 CPU 循环不会自动响应取消, 需手动检查
isActive/ensureActive()/yield().
val job = scope.launch {
while (isActive) { // 配合检查,否则取消无效
doHeavyWork()
}
}
job.cancel() // 标记取消
易错点:
CancellationException是正常的取消信号, 不要在 catch 中吞掉它, 否则破坏取消.try/catch (e: Exception)会误捕获它, 应catch (e: CancellationException) { throw e }或只 catch 具体异常.- 取消后想做清理用
try/finally, 但 finally 中若要再调挂起函数, 需用withContext(NonCancellable).
七, 异常处理
- launch: 异常会立即向上传播给父 Job, 触发父及所有兄弟取消.
- async: 失败会记录在
Deferred; 普通父 Job 下子 async 失败会立刻取消父与兄弟,await()只在调用处重新抛出该结果.supervisorScope下失败不自动取消兄弟, 但 await 仍会重新抛出; 根 async 通常由 await 或其他显式观察结果的方式暴露失败. - CoroutineExceptionHandler: 只对 launch 的根协程生效, 作为 “兜底” 处理未捕获异常.
- SupervisorJob / supervisorScope: 子协程失败不影响兄弟和父协程 (单向传播).适合多个独立任务的场景 (如同时加载多个卡片, 一个失败不拖垮其他).
supervisorScope {
launch { riskyA() } // A 失败不影响 B
launch { riskyB() }
}
对比记忆: 普通 Job 一个子失败全家取消; SupervisorJob 子失败各自负责.
八, Flow (冷流)
Flow 是协程版的 “异步数据流”, 可以发射多个值 (协程的 suspend 函数只返回一个值).
- 冷流 (Cold): Flow 默认是冷的: 没有收集者就不执行, 每个 collect 都重新触发上游. 类比 “按需播放的录像”.
- 构建:
flow { emit(x) },flowOf(),asFlow(). - 操作符:
map/filter/transform(中间操作, 惰性),collect/first/toList(末端操作, 触发执行). - 线程切换:
flowOn(Dispatchers.IO)只影响上游; 收集所在线程由 collect 处的 scope 决定. 不要在 flow {} 里用 withContext 切线程 (会报错), 要用 flowOn. - 背压:
buffer()(并发缓冲, 让上游发射和下游收集可在缓冲容量内解耦),conflate()(下游慢时只保留最新值, 适合 UI 进度 / 状态),collectLatest(新值到来取消上个收集块, 适合搜索联想 / 列表刷新).边界是: 这些操作符优化的是 “生产快, 消费慢” 的处理方式, 不应掩盖下游耗时任务本身; 被取消的收集块仍要遵守协作式取消.
flow {
emit(fetchFromNetwork()) // 上游在 IO
}.flowOn(Dispatchers.IO)
.map { it.toUiModel() }
.collect { render(it) } // 收集在调用方线程(如 Main)
九, StateFlow 与 SharedFlow (热流)
热流 (Hot): 不管有没有收集者都 “活着”, 发射独立于收集者. 用于状态 / 事件, 是 Android MVVM 中替代 LiveData 的主力.
StateFlow
- 持有一个最新值, 新收集者立即拿到当前值 (类似 LiveData).
- 必须有初始值,
value可读可写 (MutableStateFlow). - 去重: 值相等 (equals) 时不发射 (conflate 语义).
- 适合表达 UI State.
private val _uiState = MutableStateFlow(UiState.Loading)
val uiState: StateFlow<UiState> = _uiState.asStateFlow()
_uiState.value = UiState.Success(data)
SharedFlow
- 可配置
replay,extraBufferCapacity和onBufferOverflow, 适合广播允许丢失或明确允许重放的信号. replay = 0在没有订阅者时不会保留值, 因此不能保证 UI 在 STOPPED, 旋转, 切后台或进程重建期间一定收到事件.- 对于必须保证被处理的业务结果, ViewModel 应将其立即转换为可恢复的 UI state; UI 观察状态并在处理后回调 ViewModel 清除或推进状态.
- 纯 UI 局部且允许丢失的瞬时效果, 才可以在明确语义并有测试的前提下使用 SharedFlow/Channel.
对比表
| 维度 | LiveData | StateFlow | SharedFlow |
|---|---|---|---|
| 库 | Jetpack | 协程 | 协程 |
| 初始值 | 否 | 必须 | 否 |
| 生命周期感知 | 自带 | 需 repeatOnLifecycle | 需 repeatOnLifecycle |
| 去重 | 否 | 是 | 否 |
| 粘性 | 是 | 是 | 可配 replay |
| 适用 | 老项目状态 | 新项目状态 | 事件 |
Android 正确收集姿势 (避免后台浪费):
lifecycleScope.launch {
repeatOnLifecycle(Lifecycle.State.STARTED) {
viewModel.uiState.collect { render(it) }
}
}
repeatOnLifecycle 由 lifecycle-runtime-ktx 2.4.0 引入. 旧写法 lifecycleScope.launch { viewModel.uiState.collect{...} } 会一直收集直到整个 scope 取消, 退后台不停止, 回前台不重启; repeatOnLifecycle 在进入指定状态时开启收集, 离开时取消, 再次进入重新启动.
进阶补充: Channel, callbackFlow 与测试
Channel 与 callbackFlow
Channel 更像协程里的阻塞队列, 适合一对一事件传递; callbackFlow 用于把回调式 API 包装成 Flow.
fun locationFlow(): Flow<Location> = callbackFlow {
val listener = LocationListener { location -> trySend(location) }
locationClient.addListener(listener)
awaitClose { locationClient.removeListener(listener) }
}
关键点: awaitClose 必须释放监听器, 否则会泄漏.
Mutex / Semaphore: 限流与串行化
协程版 Mutex 提供互斥, Semaphore 限制并发进入数量; 与线程锁不同, 竞争时挂起协程而不是阻塞线程.
// 上下文片段: 上报串行化, 事件必须按产生顺序出网, 避免乱序与重复.
private val reportMutex = Mutex()
suspend fun enqueueReport(event: RiskEvent) {
reportMutex.withLock { // 排队: 同一时刻只有一个上报在途
sender.send(event) // send 可取消时, withLock 自动释放
}
}
// 限流: 最多 3 个并发请求打采集服务
private val permits = Semaphore(3)
suspend fun rateLimitedCollect() {
permits.withPermit { api.collect() }
}
风控场景用它保证 “上报串行 + 限流”: 事件队列消费者逐个上报, 采集失败重试时也不会乱序; 需要节流并发请求时用 Semaphore, 而不是无限开协程.
combine / zip / flatMapLatest
| 操作符 | 行为 | 场景 |
|---|---|---|
| combine | 任一上游变化就用最新值组合 | 表单状态, 多个配置源 |
| zip | 等两边一一配对 | 两个一次性结果合并 |
| flatMapLatest | 新请求来时取消旧请求 | 搜索框, 筛选条件变化 |
stateIn / shareIn
stateIn 把冷 Flow 转成有当前值的 StateFlow; shareIn 把冷 Flow 共享给多个订阅者. 面试要说明 scope, SharingStarted 策略, replay, 以及订阅者消失后上游是否继续运行. UI 常用 WhileSubscribed(stopTimeoutMillis=...), 但超时时间必须依据返回前台频率, 上游成本和数据新鲜度测量, 不能机械套固定值.
上下文片段 (Android ViewModel, 省略 UiState, repository, 依赖注入与 import):
class FeedViewModel(
repository: FeedRepository
) : ViewModel() {
val uiState: StateFlow<UiState> = repository.observeFeed()
.map<Feed, UiState> { UiState.Content(it) }
.onStart { emit(UiState.Loading) }
.catch { error -> emit(UiState.Error(error.message ?: "Unknown error")) }
.stateIn(
scope = viewModelScope,
started = SharingStarted.WhileSubscribed(stopTimeoutMillis = 5_000),
initialValue = UiState.Loading
)
val sharedRefreshes: SharedFlow<RefreshResult> = repository.observeRefreshes()
.shareIn(
scope = viewModelScope,
started = SharingStarted.WhileSubscribed(stopTimeoutMillis = 5_000),
replay = 1
)
}
上例的 5_000 是待测量的示例值, 不是推荐常量: 若短暂旋转或切换导航页很常见, 上游建立昂贵且数据允许短暂复用, 可从短超时开始; 若数据敏感, 上游成本低或后台保持订阅有风险, 应选择 0 或更严格策略. 测量前后台切换频率, 上游重连 / 查询次数, 内存和数据陈旧窗口, 再调整并记录依据. Eagerly 会立刻启动且在 scope 存活期间持续; Lazily 首次订阅后持续; WhileSubscribed 在订阅者为零后按超时停止并可配置 replay 过期. 选择错误的症状是离屏仍持续请求或回前台频繁闪烁加载; 用日志/网络拦截和生命周期事件取证, 定位 started/replay/scope, 调整后重复前后台与旋转步骤验证.
高频 Flow 的错误, 重试与时间语义
catch 能捕获它上游抛出的异常: 它可捕获 flow {}, 上游操作符和其上游内部流抛出的异常, 不能捕获位于它下游的 collect 块或后续操作符异常. 它同样不应把取消转成业务错误. retry/retryWhen 只重试暂时性且幂等安全的失败, 并必须限制次数或总时限. 搜索中应把 retry 放在每次 flatMapLatest 产生的内部请求流内, 否则外层重试会重新订阅输入流, 语义不是 “重试当前请求”. debounce 表示一段静默后才发送, 适合文本搜索; sample 表示按周期取期间最新值, 适合允许丢失的高频展示. 二者都改变事件语义, 不能用于每一条都必须处理的交易, 消息或写入.
// 上下文片段:query 是用户输入 Flow;search 为 suspend 且可取消的 repository 方法.
val results: Flow<SearchState> = query
.debounce(300)
.map(String::trim)
.distinctUntilChanged()
.flatMapLatest { keyword ->
flow { emit(repository.search(keyword)) }
.retry(retries = 2) { error -> error is IOException }
}
.catch { error -> emit(SearchState.Error(error.message ?: "search failed")) }
val sampledProgress: Flow<Int> = downloadProgress.sample(250)
300 和 250 同样只是假设值. 用输入事件时间戳, 请求数, 首个结果延迟和用户可感知刷新频率测量; 当请求仍过多提高 debounce, 当等待感明显则降低, 前提是业务允许. 以下才是会吞掉取消的反例:
// 上下文片段:不要这样写;Exception 包含 CancellationException.
suspend fun loadSafely(): UiState = try {
UiState.Content(repository.load())
} catch (error: Exception) {
UiState.Error(error.message ?: "load failed")
}
失败链路: 症状是取消页面后仍显示错误或不断重试; 证据是上述手写 try/catch (Exception) 将 CancellationException 归约为错误状态, 或 retry 未限制; 定位异常类型与 owner Job; 修复是先 catch (error: CancellationException) { throw error }, 再处理具体业务失败, 并限定可重试错误; 验证是开始请求后销毁 owner, 确认无新请求和无错误 UI 更新.
竞态叙事 (双请求竞态, 事件上报去重核心): 症状是刷新后偶发旧数据覆盖新结果; 证据是同一 key 两次请求的完成顺序与 _uiState 最终值矛盾; 定位为旧请求后完成, 覆盖了新值; 修复是启动新请求时取消旧请求 (协程取消或 flatMapLatest), 并把响应经版本号 update 归约进 StateFlow; 验证是注入两个不同延迟的响应, 断言最终 StateFlow 值为新请求结果.
协程测试
使用 runTest, StandardTestDispatcher 和虚拟时间, 不要在测试里真实 delay.
runTest 用虚拟时间执行: delay 不真实等待, 而是把任务排入虚拟时钟. 用 advanceTimeBy(ms) 快进虚拟时钟触发定时任务, 用 runCurrent() 立即执行当前虚拟时刻上已排队的任务; 断言 delay(300) 之后的逻辑时, 应先 advanceTimeBy(300) 再检查. 虚拟时间不要与真实 Thread.sleep 混用, 否则测试变慢且时序不稳.
需要测 Dispatchers.Main 相关逻辑时, 用 Dispatchers.setMain(testDispatcher) 替换主线程调度器, 测试结束用 Dispatchers.resetMain() 还原; 配合 StandardTestDispatcher 可手动控制协程的调度时机, 精确复现 “挂起后恢复” 的时序.
追问: 为什么 callbackFlow 里要写 awaitClose? 因为 Flow 被取消时需要注销外部回调, 否则回调继续持有对象导致泄漏.
高频面试题
Q1: 协程和线程的区别? 线程是 OS 调度的重量级资源; 协程是用户态的轻量 “可挂起计算”, 多个协程可复用少量线程. 协程挂起不阻塞线程. 一个线程能跑成千上万协程.
Q2: suspend 关键字做了什么? 协程凭什么不阻塞线程? 编译器对 suspend 函数做 CPS 变换, 加 Continuation 参数, 函数体编译成状态机. 挂起时返回 COROUTINE_SUSPENDED 并释放线程, 结果就绪后通过 resumeWith 从挂起点恢复. 所以是 “挂起协程” 而非 “阻塞线程”.
Q3: launch 和 async 区别? async 不 await 会怎样? launch 返回 Job 无结果, async 返回 Deferred 有结果. 普通父 Job 的子 async 失败会立刻取消父和兄弟, await 负责在调用点重新抛出结果; supervisorScope 下失败不取消兄弟, 但仍应 await / 处理对应 Deferred; 只有没有结构化父级且从未观察的根 async 才容易遗漏异常, 不能概括为 “不 await 一定静默”.
Q4: 协程的取消是怎样的? 为什么有时 cancel 无效? 协作式取消. cancel 只标记状态, suspend 函数恢复时检查并抛 CancellationException. 纯 CPU 循环不调用 suspend 函数, 不会响应取消, 需手动 isActive/ensureActive/yield.
Q5: Job 和 SupervisorJob 区别? 普通 Job: 子协程异常会取消父和所有兄弟. SupervisorJob: 子异常只影响自己, 不向上传播取消兄弟. 多个独立任务用 supervisorScope.
Q6: CoroutineScope, CoroutineContext, Job 的关系? Scope 持有 Context; Context 是元素集合 (Job, Dispatcher 等); Job 管理生命周期与父子层级. 三者共同实现结构化并发.
Q7: Flow 冷热的区别? StateFlow 和 SharedFlow 怎么选? 冷流无收集者不执行, 每次 collect 重新触发; 热流的生产生命周期由其 scope 和 started 策略决定. 持久页面状态通常用 StateFlow; 必须保证处理的业务结果先归约为 UI state. SharedFlow/Channel 只适合已经明确允许丢失, 重放和多订阅语义的信号, 不能靠 replay=0 获得 “恰好一次” 保证.
Q8: 为什么用 StateFlow 替代 LiveData? LiveData 有什么坑? StateFlow 不依赖 Android, 可在纯 Kotlin 层用, 操作符丰富, 与协程统一. LiveData 的坑: 粘性事件 (新观察者收到旧值), 只能主线程 setValue, observeForever 易泄漏. 但 StateFlow 不感知生命周期, 需 repeatOnLifecycle.
Q9: flowOn 和 withContext 区别? 能在 flow {} 里 withContext 切线程吗? 不能. flow {} 内 emit 必须在收集协程的上下文, 直接 withContext 会抛异常 (违反上下文保留).要切上游线程用 flowOn, 它只影响上游操作符.
Q10: repeatOnLifecycle 解决什么问题? 解决 “App 退后台时仍在收集 Flow 浪费资源 / 可能崩溃”.它在进入指定状态 (如 STARTED) 时启动收集, 离开时取消, 再次进入重启, 是官方推荐的安全收集方式.
练习与掌握检查
- 画出一次 “页面加载后立即返回” 的 Job 图, 标出哪个 owner 取消, 哪个 Flow 停止, 哪个状态会由新 collector 得到.
- 在最小 ViewModel 中实现上面的
stateIn, 用 fake repository 统计订阅次数; 旋转 / 停止再恢复页面, 验证策略是否符合你选择的超时语义. - 为搜索输入分别用
debounce和sample记录时间线, 说明为何前者更符合 “输入停止后查询”. - 制造一个手写
try/catch (Exception)吞掉取消的版本, 按 “症状→证据→定位→修复→验证” 复盘; 能复盘并修复即达到掌握标准.