Kotlin 协程与 Flow 完整实战
协程基础、结构化并发、调度器、Job 取消与异常处理,Flow / StateFlow / SharedFlow 操作符与 ViewModel + LiveData 整合
学习目标
- 理解 suspend 函数、CoroutineScope 与结构化并发的语义
- 掌握 Dispatchers.Main / IO / Default 的取舍与切换
- 会用 Job 取消、CoroutineExceptionHandler 与 try/catch 处理异常
- 熟练使用 Flow 的 emit / collect 与常用操作符
- 理解 StateFlow / SharedFlow 的差异,并与 ViewModel + LiveData 整合
学习目标
异步是 Android 永恒的话题:网络请求、数据库读写、定时器、动画回调,本质都是“在合适的线程做合适的事”。Java 时代我们用过 Thread / Handler / AsyncTask / RxJava,而 Kotlin 给出了官方答案 —— 协程 + Flow。
读完本章你将能够:
- 理解
suspend函数与CoroutineScope的设计动机; - 用结构化并发管理协程生命周期,避免泄漏;
- 正确选择
Dispatchers.Main / IO / Default; - 用
Job.cancel()与异常处理机制让代码在错误时优雅退化; - 用
Flow处理“多值异步流”,理解它与Sequence/LiveData的差别; - 用
StateFlow/SharedFlow与 ViewModel 整合,替代 LiveData 也能与 LiveData 互转。
协程基础
1. 什么是协程
协程是一种“轻量级线程”:它不是操作系统线程,而是 JVM 上由 Kotlin 编译器 + 调度器协作的“挂起 / 恢复”单元。一个线程上可以同时跑十万级协程,切换开销接近一次函数调用。
它解决三个痛点:
- 回调地狱:嵌套回调写起来痛苦,
suspend让你像写同步代码一样写异步; - 取消传递:协程天然支持协作式取消,错误发生时能沿调用栈传递;
- 结构化并发:协程必须挂在某个
CoroutineScope下,作用域销毁时所有子协程自动取消。
2. suspend 函数
suspend 修饰的函数可以“挂起”协程而不阻塞线程,等就绪后自动恢复:
suspend fun fetchUser(id: Long): User {
delay(500) // 非阻塞等待,相当于 Thread.sleep 的协程版
return User(id, "User-$id", 20)
}
suspend fun main() {
val user = fetchUser(1L)
println(user)
}
suspend 函数只能在协程或其它 suspend 函数中调用。编译器会把 suspend 函数翻译成带 Continuation 参数的“状态机”,本质是回调,但语法是顺序的。
3. 启动协程:runBlocking / coroutineScope / launch / async
import kotlinx.coroutines.*
fun main() = runBlocking { // 阻塞当前线程直到协程完成,常用于 main / 测试
launch { // 启动新协程,不等待返回值
delay(200)
println("A")
}
val deferred = async { // 启动新协程,返回 Deferred<T>
delay(100)
42
}
println(deferred.await()) // 42
println("done")
}
// 输出:42 done A
四个核心 API:
| API | 作用 | 返回 |
|---|---|---|
runBlocking |
阻塞当前线程直到完成 | T |
coroutineScope |
创建子作用域并等待所有子协程 | T |
launch |
启动“火并忘记”的协程 | Job |
async |
启动带返回值的协程 | Deferred<T> |
结构化并发
1. CoroutineScope 与父子关系
每个协程必须属于一个 CoroutineScope。launch / async 创建的协程会成为当前作用域的“子协程”,父协程会等待所有子完成;父被取消时子也连带取消。这就是结构化并发:
suspend fun loadAll() = coroutineScope {
launch { fetchUsers() } // 子 1
launch { fetchPosts() } // 子 2
// coroutineScope 阻塞到两个子都完成
}
2. supervisorScope:异常隔离
普通 coroutineScope 中任何一个子协程抛异常会取消所有兄弟。supervisorScope 则允许兄弟继续运行:
suspend fun loadDashboard() = supervisorScope {
launch { fetchProfile() } // 即使这里抛异常
launch { fetchFeed() } // 这里仍会执行完
}
典型场景:仪表盘多个独立卡片,一个失败不应影响其它。
3. Android 中的作用域
AndroidX 提供了 viewModelScope 与 lifecycleScope:
class ProfileViewModel : ViewModel() {
fun load() {
viewModelScope.launch { // 与 ViewModel 同生命周期
val user = fetchUser(1L)
// ...
}
}
override fun onCleared() {
super.onCleared()
// viewModelScope 自动 cancel,无需手动
}
}
class MainActivity : AppCompatActivity() {
override fun onCreate(savedInstanceState: Bundle?) {
super.onCreate(savedInstanceState)
lifecycleScope.launch { // 与 Activity 生命周期绑定
// 在 DESTROYED 时自动取消
}
}
}
依赖(用 bash 标记,Shiki 不识别 gradle):
dependencies {
implementation "org.jetbrains.kotlinx:kotlinx-coroutines-android:1.8.1"
implementation "androidx.lifecycle:lifecycle-viewmodel-ktx:2.8.4"
implementation "androidx.lifecycle:lifecycle-runtime-ktx:2.8.4"
}
调度器 Dispatchers
协程默认不绑定特定线程,需要用 withContext 或 launch(Dispatchers.X) 显式切换:
suspend fun loadUser(): User = withContext(Dispatchers.IO) {
// 网络 / 数据库等阻塞 IO
api.fetchUser()
}
| Dispatchers | 适合场景 | 实现 |
|---|---|---|
Main |
UI 操作、LiveData.setValue |
Android 主线程调度器 |
IO |
网络、磁盘、数据库 | 64 线程的共享池 |
Default |
CPU 密集:排序、解析、图片解码 | Runtime.availableProcessors() 个线程 |
Unconfined |
不推荐日常使用 | 调用方线程直接执行 |
切换调度器后协程并不会创建新“任务”,withContext 是 suspend 函数,编译成状态机:
suspend fun showUser() {
val user = withContext(Dispatchers.IO) { api.fetch() } // IO
withContext(Dispatchers.Main) { textView.text = user.name } // Main
}
在 Compose 中要用
rememberCoroutineScope()获取绑定 Composable 的作用域,里面默认是Main调度器。
Job 与取消
1. 协作式取消
Job.cancel() 只是设置取消标志,协程在下一次 suspend 点(如 delay / withContext / yield)才会真正抛 CancellationException。如果协程内全是 CPU 计算没有 suspend 点,会一直跑下去:
val job = launch {
var i = 0
while (isActive) { // 注意用 isActive 而非 true
i++
if (i % 1000 == 0) yield() // 主动让出,给取消机会
}
}
delay(100)
job.cancel()
2. 取消的清理
被取消的协程会在下一个 suspend 点抛 CancellationException,如果协程内有资源需在 finally 中释放:
suspend fun readFile() {
val input = openFile()
try {
input.use { readAll(it) }
} finally {
// 即便被取消也会执行
// 注意:finally 里不能调用其它 suspend,否则会立即抛 CancellationException
// 如必须调用,用 NonCancellable 作用域:
withContext(NonCancellable) { closeResource() }
}
}
3. 超时
val result = withTimeoutOrNull(3000) {
api.fetchBig() // 3 秒内返回结果,否则返回 null
}
withTimeout 会在超时时抛 TimeoutCancellationException,withTimeoutOrNull 则优雅返回 null。
异常处理
1. try / catch
普通 launch 抛出的异常会传播到父作用域,要在协程内 try/catch:
viewModelScope.launch {
try {
val user = fetchUser()
_state.value = UiState.Success(user)
} catch (e: IOException) {
_state.value = UiState.Error(e.message ?: "网络错误")
} catch (e: CancellationException) {
throw e // 必须重新抛出,否则破坏取消语义
}
}
坑:catch 块捕获
CancellationException后必须重新抛出,否则会破坏取消链。
2. CoroutineExceptionHandler
给 launch 加 CoroutineExceptionHandler 作为“最后兜底”:
val handler = CoroutineExceptionHandler { _, e ->
Log.e("App", "未捕获异常", e)
}
viewModelScope.launch(handler) {
riskyCall()
}
注意它只对 launch 有效,async 抛出的异常会在 await() 时再次抛出。SupervisorJob 下才能让异常真正交到 handler。
3. SupervisorJob 与 Job
默认 Job 父子之间会传递取消;SupervisorJob 不传递。Android 中 viewModelScope 内部就是 SupervisorJob + Main,所以子协程崩溃不会连带取消其它兄弟。
Flow 入门
Flow<T> 是冷流:只有调用方 collect 时才会执行流定义的代码,且每个收集者都跑一遍。它对应“多值异步序列”,介于 Sequence 与 LiveData 之间。
1. 创建与收集
fun numberFlow(): Flow<Int> = flow {
for (i in 0..3) {
delay(100)
emit(i) // 发射值
}
}
suspend fun main() {
numberFlow().collect { value ->
println(value) // 0 1 2 3
}
}
flow { } builder 提供挂起上下文。常用 builder:
| Builder | 含义 |
|---|---|
flow { } |
任意挂起逻辑 |
flowOf(1, 2, 3) |
固定值 |
asFlow() |
把 List<T> / IntRange 转流 |
channelFlow { } |
允许在多个协程并发发送 |
2. 中间操作符
numberFlow()
.map { it * it } // 转换
.filter { it > 1 } // 过滤
.onEach { println("发射 $it") } // 副作用
.collect { println(it) }
常用操作符:
map/filter/onEach/onStart/onCompletiontransform:万能版 map,可发多个值flatMapLatest/flatMapConcat/flatMapMerge:展平嵌套流debounce/distinctUntilChanged/conflate:去抖 / 去重 / 合并combine/zip:组合多个流take/drop/takeWhile
3. 三种 flatMap 的区别
flow {
emit(1); emit(2); emit(3)
}.flatMapConcat { id ->
flow { emit("$id-a"); delay(50); emit("$id-b") }
}.collect { println(it) }
// 输出:1-a 1-b 2-a 2-b 3-a 3-b(顺序执行)
| 操作符 | 行为 |
|---|---|
flatMapConcat |
串行:上一个流结束再开始下一个 |
flatMapMerge |
并发:所有流同时跑,乱序输出 |
flatMapLatest |
切换:新值来时取消上一个流,适合搜索框 |
4. 异常与取消
flow {
emit(1)
throw RuntimeException("boom")
}.catch { e ->
emit(-1) // 错误时回退发射一个默认值
}.onCompletion {
println("完成")
}.collect { println(it) }
catch 操作符只捕获“上游”异常,且只对 collect 阶段的异常无效。Flow 取消通过收集者的协程取消,与协程取消语义一致。
5. flowOn 切换调度器
flow { emit(api.fetch()) }
.flowOn(Dispatchers.IO) // 上游切到 IO
.collect { updateUi(it) } // 下游在调用方线程
flowOn 影响上游所有操作符;下游保持调用方协程上下文。
StateFlow 与 SharedFlow
冷流每次 collect 都会重新执行,并不适合“状态持有 + 共享”。这就是 StateFlow / SharedFlow 的定位 —— 它们是 热流:始终存活,多个观察者共享同一份数据。
1. StateFlow
class CounterViewModel : ViewModel() {
private val _count = MutableStateFlow(0)
val count: StateFlow<Int> = _count.asStateFlow()
fun inc() { _count.value++ }
fun set(v: Int) { _count.value = v }
}
特点:
- 必须有初值;
value属性可在任意线程安全读写;- 与
LiveData类似的“粘性”行为:新收集者会立即收到最新值; distinctUntilChanged内置:相同值不会重复发射。
收集(在 Compose 中常用 collectAsStateWithLifecycle,下章讲):
lifecycleScope.launch {
viewModel.count.collect { value ->
textView.text = value.toString()
}
}
2. SharedFlow
SharedFlow 更通用,可以配置缓存与回放:
val events = MutableSharedFlow<String>(
replay = 0, // 新订阅者收到 0 个历史值
extraBufferCapacity = 16,
onBufferOverflow = BufferOverflow.DROP_OLDEST
)
// 发射
viewModelScope.launch { events.emit("Clicked") }
// 收集
events.collect { event -> handle(event) }
| 类型 | replay | 用途 |
|---|---|---|
StateFlow |
固定 1 | 状态持有 |
SharedFlow(replay=0) |
0 | 一次性事件(导航、Toast、SnackBar 触发) |
SharedFlow(replay=N) |
N | 历史回放 |
典型坑:把“显示一次 Toast”这种一次性事件用
StateFlow存,会导致配置变更后再次弹一次。改用SharedFlow(replay=0)即可。
3. LiveData 与 Flow 互转
AndroidX 提供了双向桥:
// Flow -> LiveData
val state: LiveData<User> = userFlow.asLiveData()
// LiveData -> Flow(推荐用 asFlow)
val flow: Flow<User> = liveData.asFlow()
实际项目中 新代码建议直接用 StateFlow,仅在需要与老 LiveData API 互操作时桥接。
与 ViewModel 整合:完整示例
下面是一个完整的“搜索用户”示例,融合了 Flow 操作符、StateFlow、ViewModel 与生命周期:
data class User(val id: Long, val name: String)
sealed class UiState<out T> {
object Loading : UiState<Nothing>()
data class Success<T>(val data: T) : UiState<T>()
data class Error(val message: String) : UiState<Nothing>()
}
class SearchViewModel(
private val repo: UserRepository,
) : ViewModel() {
private val _query = MutableStateFlow("")
val query: StateFlow<String> = _query.asStateFlow()
val result: StateFlow<UiState<List<User>>> = _query
.debounce(300) // 防抖 300ms
.distinctUntilChanged() // 相同值不重复查
.filter { it.length >= 2 } // 至少 2 字符
.flatMapLatest { q -> // 新输入时取消旧请求
flow {
emit(UiState.Loading)
try {
emit(UiState.Success(repo.search(q)))
} catch (e: IOException) {
emit(UiState.Error(e.message ?: "网络错误"))
}
}
}
.stateIn(
scope = viewModelScope,
started = SharingStarted.WhileSubscribed(5000),
initialValue = UiState.Loading,
)
fun onQueryChange(q: String) { _query.value = q }
}
Activity 端:
class SearchActivity : AppCompatActivity() {
private val vm by viewModels<SearchViewModel>()
override fun onCreate(savedInstanceState: Bundle?) {
super.onCreate(savedInstanceState)
setContentView(R.layout.activity_search)
val edit = findViewById<EditText>(R.id.et_query)
val list = findViewById<RecyclerView>(R.id.rv_users)
// 监听输入
edit.doOnTextChanged { text, _, _, _ ->
vm.onQueryChange(text?.toString() ?: "")
}
// 收集状态:必须用 lifecycleScope.launchWhenStarted 或 repeatOnLifecycle
lifecycleScope.launch {
repeatOnLifecycle(Lifecycle.State.STARTED) {
vm.result.collect { state ->
when (state) {
is UiState.Loading -> showLoading()
is UiState.Success -> adapter.submit(state.data)
is UiState.Error -> showError(state.message)
}
}
}
}
}
}
repeatOnLifecycle(STARTED) { ... }会在STARTED后开始收集,STOPPED时自动取消收集,避免后台浪费 CPU。Compose 中改用collectAsStateWithLifecycle()一行搞定(下章讲)。
常见坑与最佳实践
- GlobalScope 满天飞:
GlobalScope没有父作用域,Activity 销毁时协程仍在跑,造成泄漏与崩溃。Android 必须用viewModelScope/lifecycleScope。 - 主线程做 IO:
Dispatchers.Main中调用阻塞 API 会卡 UI;网络 / 数据库要切IO。 async不await:异常被吞到Deferred里直到await()才抛,往往掩盖 bug。每个async都应该有await,或用awaitAll。- catch 吞了
CancellationException:catchThrowable/Exception时务必把CancellationException重新抛出,否则破坏取消链。 StateFlow当一次性事件:配置变更后新观察者会再收到“最新一次”,把“导航”这类事件用SharedFlow(replay=0)存。Flow用launch收集:直接flow.collect { }在onCreate里启动会一直阻塞协程直到流结束。在 Android 中需要repeatOnLifecycle或collectAsStateWithLifecycle。flowOn写在 collect 后:flowOn只影响上游操作符,写在collect之前才有意义。flatMapLatest里捕获异常:上游切换时旧协程被取消,会抛CancellationException,别误当成业务错误。- 默认
Flow是冷流:每次collect都重新执行,不要把“网络请求”放在 cold flow 的 builder 中然后多次 collect,每次都会重新发请求。改成stateIn/shareIn才能复用。 - LiveData 与 StateFlow 双重暴露:ViewModel 同时暴露
LiveData与StateFlow,UI 不知该听谁。新代码统一用StateFlow,老代码用asLiveData()桥接。
章节小结
- 协程 = 轻量级挂起 / 恢复单元,
suspend让异步像同步一样写; - 结构化并发通过
CoroutineScope+ 父子 Job 管理生命周期,Android 中用viewModelScope/lifecycleScope; Dispatchers.Main / IO / Default各司其职,withContext切换;Job.cancel()协作式取消要在 suspend 点生效;- 异常用
try/catch+CoroutineExceptionHandler,SupervisorJob/supervisorScope隔离兄弟; Flow是冷流,用flow { emit }+collect;常用操作符map / filter / debounce / flatMapLatest / catch;StateFlow持有状态,SharedFlow(replay=0)表达一次性事件;与 LiveData 通过asLiveData/asFlow互转;- 在 Android 中用
repeatOnLifecycle或 Compose 中的collectAsStateWithLifecycle安全收集。
下一章预告
掌握了协程与 Flow,下一步是 Jetpack Compose 入门与环境搭建:Compose 与传统 View 的范式差异、@Composable 注解、函数式 UI 思维、Preview 注解、第一个 Greeting 组件、setContent 接入 Activity,以及常见的调试技巧。准备好进入声明式 UI 的世界!