A
第 43 章KOTLIN50 分钟

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 编译器 + 调度器协作的“挂起 / 恢复”单元。一个线程上可以同时跑十万级协程,切换开销接近一次函数调用。

它解决三个痛点:

  1. 回调地狱:嵌套回调写起来痛苦,suspend 让你像写同步代码一样写异步;
  2. 取消传递:协程天然支持协作式取消,错误发生时能沿调用栈传递;
  3. 结构化并发:协程必须挂在某个 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 / onCompletion
  • transform:万能版 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() 一行搞定(下章讲)。

常见坑与最佳实践

  1. GlobalScope 满天飞:GlobalScope 没有父作用域,Activity 销毁时协程仍在跑,造成泄漏与崩溃。Android 必须用 viewModelScope / lifecycleScope。
  2. 主线程做 IO:Dispatchers.Main 中调用阻塞 API 会卡 UI;网络 / 数据库要切 IO。
  3. async 不 await:异常被吞到 Deferred 里直到 await() 才抛,往往掩盖 bug。每个 async 都应该有 await,或用 awaitAll。
  4. catch 吞了 CancellationException:catch Throwable / Exception 时务必把 CancellationException 重新抛出,否则破坏取消链。
  5. StateFlow 当一次性事件:配置变更后新观察者会再收到“最新一次”,把“导航”这类事件用 SharedFlow(replay=0) 存。
  6. Flow 用 launch 收集:直接 flow.collect { } 在 onCreate 里启动会一直阻塞协程直到流结束。在 Android 中需要 repeatOnLifecycle 或 collectAsStateWithLifecycle。
  7. flowOn 写在 collect 后:flowOn 只影响上游操作符,写在 collect 之前才有意义。
  8. flatMapLatest 里捕获异常:上游切换时旧协程被取消,会抛 CancellationException,别误当成业务错误。
  9. 默认 Flow 是冷流:每次 collect 都重新执行,不要把“网络请求”放在 cold flow 的 builder 中然后多次 collect,每次都会重新发请求。改成 stateIn / shareIn 才能复用。
  10. 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 的世界!