Keyboard shortcuts

Press ← or → to navigate between chapters

Press S or / to search in the book

Press ? to show this help

Press Esc to hide this help

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):

  1. 编译器给每个 suspend 函数隐式加一个参数 Continuation(回调), 返回值变成 Any?.
  2. 函数体被编译成一个状态机: 每个挂起点是一个状态.
  3. 挂起时函数返回特殊标记 COROUTINE_SUSPENDED, 线程被释放.
  4. 当结果就绪, 调用 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.MainUI 操作主线程
Dispatchers.IO网络 / 磁盘 IO共享弹性调度器; 默认并行度至少为 64 或可用处理器数 (取较大者)
Dispatchers.DefaultCPU 密集 (排序 / 解析)核数大小
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.

对比表

维度LiveDataStateFlowSharedFlow
库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) 时启动收集, 离开时取消, 再次进入重启, 是官方推荐的安全收集方式.

练习与掌握检查

  1. 画出一次 “页面加载后立即返回” 的 Job 图, 标出哪个 owner 取消, 哪个 Flow 停止, 哪个状态会由新 collector 得到.
  2. 在最小 ViewModel 中实现上面的 stateIn, 用 fake repository 统计订阅次数; 旋转 / 停止再恢复页面, 验证策略是否符合你选择的超时语义.
  3. 为搜索输入分别用 debounce 和 sample 记录时间线, 说明为何前者更符合 “输入停止后查询”.
  4. 制造一个手写 try/catch (Exception) 吞掉取消的版本, 按 “症状→证据→定位→修复→验证” 复盘; 能复盘并修复即达到掌握标准.