Kotlin 协程与 Flow 实战
Kotlin 协程与 Flow 实战
摘要:协程是 Kotlin 处理异步与并发的主力方案,Flow 则是冷数据流。本文系统讲解协程基础、结构化并发、调度器、异常处理,以及 Flow 的操作符与背压,帮助你在 Android/后端中正确使用。
一、协程是什么
协程是一种可挂起的计算实例。它像轻量级线程,但并非由操作系统调度,而是由 Kotlin 运行时协作式调度。一个线程上可以跑十万级协程,因为它们在挂起时不占用线程。
核心价值:
- 用同步写法写异步代码,避免回调地狱。
- 结构化并发,父子协程生命周期绑定,不会泄漏。
- 取消传播天然支持。
二、启动协程
builders
fun main() = runBlocking {
launch {
delay(100)
println("子协程完成")
}
println("主协程继续")
}
runBlocking:阻塞当前线程直到协程完成,常用于测试与 main。launch:启动一个不返回结果的协程,返回Job。async:启动并返回Deferred<T>,用await取结果。
suspend 函数
suspend fun fetchUser(id: Int): User {
delay(500)
return User(id)
}
suspend 关键字表示函数可挂起。只能在协程或其他 suspend 函数中调用。
三、结构化并发
suspend fun loadAll() = coroutineScope {
val a = async { fetchA() }
val b = async { fetchB() }
combine(a.await(), b.await())
}
coroutineScope 会等待所有子协程完成;任一子协程失败,会取消其他子协程并向上抛异常。这是结构化并发的精髓:作用域内的任务要么全成功,要么统一失败处理。
CoroutineScope 与 Job
val scope = CoroutineScope(Job() + Dispatchers.Main)
val job = scope.launch { ... }
job.cancel() // 取消单个
scope.cancel() // 取消整个作用域
四、调度器
Dispatchers.Main:主线程,UI 操作。Dispatchers.IO:IO 密集,线程池较大。Dispatchers.Default:CPU 密集,线程数等于 CPU 核数。Dispatchers.Unconfined:不切换线程,慎用。
切换线程用 withContext:
suspend fun save(data: Data) {
val parsed = withContext(Dispatchers.Default) { parse(data) }
withContext(Dispatchers.IO) { db.write(parsed) }
}
withContext 会挂起当前协程,切到指定线程执行,完成后切回。
五、取消与超时
协程取消是协作式的:被取消的协程在挂起点(如 delay、await)会抛出 CancellationException。自定义耗时操作需检查 isActive:
while (isActive) {
doChunk()
}
超时:
withTimeout(2000) { fetch() }
或返回 null 而不抛异常:
withTimeoutOrNull(2000) { fetch() }
六、异常处理
协程异常传播遵循"子失败父失败"原则。async 在 await 时才抛异常;launch 立即由 CoroutineExceptionHandler 处理。
val handler = CoroutineExceptionHandler { _, e ->
log("捕获: ${e.message}")
}
scope.launch(handler) { throw RuntimeException("boom") }
supervisorJob 让子协程失败不影响兄弟:
val scope = CoroutineScope(SupervisorJob() + Dispatchers.Main)
七、Flow 基础
Flow 是冷流,类似 RxJava 的 Observable:
fun numbers(): Flow<Int> = flow {
for (i in 1..3) {
delay(100)
emit(i)
}
}
runBlocking {
numbers().collect { println(it) }
}
只有 collect 时才会执行。每次 collect 都重新执行。
常用操作符
numbers()
.map { it * 2 }
.filter { it > 2 }
.take(2)
.collect { println(it) }
map:转换。filter:过滤。take:取前 N 个后自动取消上游。flatMapLatest/flatMapConcat/flatMapMerge:展平,行为各异。debounce:防抖。distinctUntilChanged:去重。
八、热流 StateFlow 与 SharedFlow
StateFlow 永远有值且只保留最新值,适合表示 UI 状态:
class ViewModel : ViewModel() {
private val _uiState = MutableStateFlow(UiState.Loading)
val uiState: StateFlow<UiState> = _uiState.asStateFlow()
fun load() {
viewModelScope.launch {
_uiState.value = UiState.Loading
_uiState.value = UiState.Success(repo.fetch())
}
}
}
SharedFlow 更通用,可配置缓存与重放:
private val _events = MutableSharedFlow<Event>(extraBufferCapacity = 16)
val events = _events.asSharedFlow()
九、背压
冷流默认无缓冲,emit 会挂起直到 collect 处理完。可用 buffer 引入缓冲区,conflate 只保留最新,flowOn 切换上游线程:
flow { emitHeavy() }
.flowOn(Dispatchers.Default)
.buffer(64)
.collect { render(it) }
十、测试
用 runTest 替代 runBlocking,可虚拟时间:
@Test
fun test() = runTest {
val result = withTimeout(1000) { fetch() }
assertEquals(expected, result)
}
用 Turbine 库测试 Flow 更方便,可断言每次发射。
十一、总结
协程与 Flow 用一套统一的抽象解决了异步、并发、流式数据三大问题。记住三条原则:结构化并发(不要用 GlobalScope)、协作式取消(检查 isActive)、用 withContext 切线程而非 launch。在 Android 中配合 ViewModel 的 viewModelScope,可以让生命周期管理变得几乎透明。
- 点赞
- 收藏
- 关注作者
评论(0)