跳到主要内容

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 操作、轻量任务
IO64 个线程网络、文件读写
DefaultCPU 核数计算密集
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 中)的异常不会被捕获,需单独处理
  • 协程和线程的核心区别?
    • 协程是语言级的轻量调度单元,运行在线程之上
    • 一个线程可运行数千个协程
    • 协程切换无需内核介入,开销远小于线程切换