Kotlin 协程与 Flow 深入面经
Kotlin 协程深水区
1. 协程底层原理
CPS 变换(Continuation Passing Style)
- Kotlin 编译器将 suspend 函数变换为状态机
- 每个 suspend 点都是一个状态,函数被切分成多段
// 原始代码
suspend fun fetchUser(id: String): User {
val token = getToken() // suspend 点 1
val user = api.get(id) // suspend 点 2
return user
}
// 编译后(简化伪代码)
fun fetchUser(id: String, continuation: Continuation<User>) {
when (continuation.label) {
0 -> {
continuation.label = 1
getToken(continuation) // 挂起等待
}
1 -> {
val token = continuation.result as String
continuation.label = 2
api.get(id, continuation)
}
2 -> {
continuation.resume(continuation.result as User)
}
}
}
挂起与恢复
- suspend 函数调用后,当前线程不阻塞,返回
COROUTINE_SUSPENDED标记 - IO 完成后,通过
continuation.resume(result)恢复执行 - 恢复时可能在不同线程(取决于 Dispatcher)
2. Flow 深入
冷流 vs 热流
// 冷流:collect 时才执行,每个 collector 独立
val coldFlow = flow {
println("开始生产")
emit(1)
emit(2)
}
// 没有 collect 时,内部代码不执行
// StateFlow:热流,始终有最新值
val stateFlow = MutableStateFlow(0)
// SharedFlow:热流,replay 控制新订阅者能收到的历史消息数
val sharedFlow = MutableSharedFlow<Event>(replay = 0, extraBufferCapacity = 64)
背压处理
// buffer:生产者和消费者并发运行
flow { emit(heavyData()) }
.buffer(Channel.UNLIMITED)
.collect { process(it) }
// conflate:只保留最新值,丢弃中间值(适合 UI 状态)
stateFlow.conflate().collect { render(it) }
// collectLatest:新值到来时取消旧的处理
stateFlow.collectLatest { slowRender(it) } // 只完成最新一次
channelFlow 并发发射
// 允许在不同协程中 emit
val result = channelFlow {
launch { send(fetchFromNetwork()) } // 并发
launch { send(fetchFromCache()) } // 并发
}
3. 异常处理
// SupervisorJob:子协程独立失败
val scope = CoroutineScope(SupervisorJob() + Dispatchers.IO)
scope.launch { failingTask() } // 失败不影响其他子协程
scope.launch { normalTask() } // 继续正常运行
// CoroutineExceptionHandler
val handler = CoroutineExceptionHandler { _, e ->
Log.e("TAG", "Uncaught exception", e)
}
scope.launch(handler) { riskyWork() }
// async 的异常在 await 时抛出
val deferred = async { riskyTask() }
try {
val result = deferred.await()
} catch (e: Exception) {
// 处理异常
}
// Flow 的异常处理
flow { emit(riskyOperation()) }
.catch { e -> emit(defaultValue) } // 只捕获 catch 上游异常
.collect { }
4. Dispatcher 选择
| Dispatcher | 线程池 | 适用场景 |
|---|---|---|
| Main | 主线程 | UI 操作、轻量任务 |
| IO | 64 个线程 | 网络、文件读写 |
| Default | CPU 核数 | 计算密集 |
| Unconfined | 不指定 | 测试用 |
// limitedParallelism 限制并发度
val limitedIO = Dispatchers.IO.limitedParallelism(4)
withContext(limitedIO) { /* 最多 4 个线程并发执行 */ }
5. 结构化并发
viewModelScope.launch {
// 并发执行两个请求
val userDeferred = async { fetchUser() }
val postsDeferred = async { fetchPosts() }
// 任一失败,另一个自动取消
val user = userDeferred.await()
val posts = postsDeferred.await()
_uiState.value = Success(user, posts)
}
// coroutineScope:子协程全部完成才返回
suspend fun loadData() = coroutineScope {
val a = async { taskA() }
val b = async { taskB() }
Pair(a.await(), b.await()) // 任一失败都会取消其他
}
6. ViewModel 中的实践
class OrderViewModel(
private val repo: OrderRepository,
savedState: SavedStateHandle
) : ViewModel() {
private val orderId = savedState.get<String>("orderId")!!
// StateFlow 替代 LiveData
private val _uiState = MutableStateFlow<UiState>(Loading)
val uiState: StateFlow<UiState> = _uiState.asStateFlow()
// 一次性事件用 SharedFlow
private val _events = MutableSharedFlow<Event>()
val events: SharedFlow<Event> = _events.asSharedFlow()
init {
loadOrder()
}
private fun loadOrder() = viewModelScope.launch {
_uiState.value = Loading
try {
val order = repo.getOrder(orderId)
_uiState.value = Success(order)
} catch (e: Exception) {
_uiState.value = Error(e.message)
}
}
fun submitOrder() = viewModelScope.launch {
try {
repo.submit(orderId)
_events.emit(Event.NavigateToSuccess)
} catch (e: Exception) {
_events.emit(Event.ShowError(e.message))
}
}
}
7. 资深面试题
- suspend 函数为什么不能在普通函数中调用?
- 编译后带 Continuation 参数,需要协程上下文提供该对象
- 景默认解决:用
runBlocking桥接(但会阻塞线程,仅用于测试)
- StateFlow vs LiveData 如何选择?
- 新项目一律推荐 StateFlow:要求初始值、协程原生、功能更强完整
- LiveData 主要优势是自动处理生命周期,StateFlow 必须配合
repeatOnLifecycle
- Flow 的 catch 操作符为什么只捕上游异常?
- Flow 是流式数据,cat操作符在管道中,只能捕获它上游的异常
- 下游(collect 中)的异常不会被捕获,需单独处理
- 协程和线程的核心区别?
- 协程是语言级的轻量调度单元,运行在线程之上
- 一个线程可运行数千个协程
- 协程切换无需内核介入,开销远小于线程切换