第093篇 Flow 入门:冷流、热流与背压

简介: 本文深入解析 Kotlin Flow 的核心机制:它本质是**冷流 + 一对一发射 + 结构化协作**的组合。厘清 Flow 与 RxJava Flowable 的差异、冷流“按需生产”特性、`collect` 的挂起本质及天然背压原理(依托 `emit` 挂起),并结合 Android 工程实践(`repeatOnLifecycle`、`flowOn` 误区、异常转数据流)与高频面试题,助你真正掌握 Flow 设计哲学与避坑要点。

协程把"异步"这件事讲清楚了,Flow 解决的是另一个维度:异步数据流。上一节收在异常处理,本节开始一个新主题:为什么需要 Flow、它和 RxJava 的 Flowable 差在哪、冷流为什么"有订阅才生产"、以及 collect 与"直接在协程里 while 循环读队列"到底差在哪里。这题常被当成热身,但面试官往往在最后一句追问里放一个能筛掉人的问题:同一个 Flow 连续 collect 两次会怎样。

先把结论放在前面:Flow 是冷流 + 一对一发射 + 结构化协作三件事的组合。冷流意味着每个收集者都独立触发一次上游生产(和热流相对);一对一意味着 collect 是 suspend 函数、发射不能并发;结构化意味着 Flow 的收集受收集者所在协程的取消与作用域约束。把这三条讲清,Flow 的坑位就基本定了。

机制拆解

flow { } 构造的是一个对象,不是数据。这个对象里装着一个 FlowCollector<T> 类型的 lambda,只有在有人 collect 时才被调用。这就是"冷流"的含义所在——生产是被消费驱动的,不是被创建驱动的。

顺着这条路径看,对比"直接 while 循环读队列"就清楚了:

// 传统写法:自己管队列、自己管唤醒、自己管结束标志
while (isActive) {
   
    val item = queue.take()      // 阻塞点
    handle(item)
}
// 问题:queue 满了怎么办?take 怎么被唤醒?页面退出谁来打断?

// Flow 写法:生产逻辑挂在 collect 上,结束由协程作用域管
flow {
   
    producer().forEach {
    emit(it) }   // emit 是 suspend,天然背压
}.collect {
    handle(it) }

emit 是挂起函数,这一点是 Flow 背压机制的秘密所在。上游 emit 时若下游正在处理,下游会挂起发射者,从而形成"拉不动就别拉"的自然反压。对比 RxJava 的 Flowable,Flow 靠的是 Kotlin 挂起函数天然提供的这套机制,不需要额外的 BackpressureStrategy 枚举——这也是很多从 Rx 迁过来的人最常问的一点。

再讲 collect 的挂起性质。collect 是 suspend 函数,两个 collect 不能并发调用同一个 Flow 的同一个实例,除非上游做了多发射者兼容(比如 channelFlow)。

工程落地:真实项目里怎么用

场景一:Room 的 Flow<List<User>> 首次查询发射一次,此后数据变更时自动重发射。ViewModel 里 val users = repo.observeUsers(),Fragment viewLifecycleOwner.lifecycleScope.launch { repeatOnLifecycle(STARTED) { viewModel.users.collect { render(it) } } }。这四行代码里藏着三个考点:为什么用 viewLifecycleOwner(避免视图销毁后残留收集)、为什么 launch 而不是 async(不需要返回值)、为什么 repeatOnLifecycle 而不是裸 collect(页面不可见时取消收集,省电省流量)。

场景二:把 Flow 放进 lifecycleScope 但没用 repeatOnLifecycle。APP 退到后台时收集仍在跑,传感器/数据库回调持续触发,耗电且可能触发后台限制。修法就是补上 repeatOnLifecycle。

场景三:flowOf(1, 2, 3).map { delay(1000); it }.flowOn(Dispatchers.Default)。不加 flowOn 时 delay 会阻塞收集协程的线程;加上后上游在 Default 池运行。下游的 flowOn 影响的是上游,这是常见误解点。

这些坑的正确绕法

在 flow 构建器里切换上下文直接 emit,触发异常提示应用 flowOn。 报错的原文是 Flow invariant is violated: Emission from another coroutine is detected。原因正是 emit 只能在构建协程里调用。修法:不在 flow {} 里写 withContext,而用 flowOn 把上游整体挪到别的调度器。

其次是用 collect 拉数据却忘了处理失败路径,网络异常直接把协程崩掉。修法:上游用 catch { e -> emit(ErrorState(e)) },把异常转成发射出去的数据流,而不是让它逃逸。

还有一个更隐蔽的坑:在同一个 Flow 实例上并发 collect,或把一个单消费者通道写成热流共享。修法:单消费者场景各自 collect(冷流本就支持多收集者);确需共享时用 shareIn/stateIn,并显式给 SharingStarted 策略。

代码里见真章

看一段能直接跑的代码,把上面的机制落到具体写法上:

// 1) 冷流:每次 collect 独立执行一遍上游
val counter = flow {
   
    println("上游开始")     // collect 两次会打印两次
    repeat(3) {
    emit(it) }
}

// 2) flowOn 切上游调度器(不是切下游)
val data = flow {
   
    api.queryBlocking()     // 假装是阻塞 IO
        .forEach {
    emit(it) }
}.flowOn(Dispatchers.IO)
    .flowOn(Dispatchers.Default)  // 多个 flowOn 各自成一段

// 3) 异常转数据流:catch 之后还能正常 collect
val uiState: Flow<UiState> = repo.observe()
    .map<List<User>, UiState> {
    UiState.Loading -> UiState.Content(it) }
    .catch {
    e -> emit(UiState.Error(e.message ?: "未知")) }
    .flowOn(Dispatchers.IO)

// 4) Android 落地:视图可见才收集
class UserFragment : Fragment() {
   
    private val vm: UserViewModel by viewModels()
    override fun onViewCreated(v: View, s: Bundle?) {
   
        v.findViewById<RecyclerView>(R.id.list).let {
    rv ->
            vm.users.flowWithLifecycle(v.lifecycle)   // 或用 repeatOnLifecycle 手写
                .onEach {
    rv.adapter.submitList(it) }
                .launchIn(v.lifecycleScope)
        }
    }
}

关键行解读:flow { } 体只在 collect 时执行,这就是冷流;flowOn 影响上游段,可多次串联形成不同调度器流水线;catch 之后用 emit 发出错误态,避免异常逃逸;Android 端用 flowWithLifecycle/launchIn 简化写法,本质就是 repeatOnLifecycle + collect。

这题在面试里怎么问、怎么答

"冷流和热流什么区别?Flow 属于哪种?" 答:冷流每次订阅重新生产、无状态、可多方订阅;热流按时间推送、共享状态、错过即丢失。Flow 默认冷,shareIn/stateIn 可转热。

"Flow 怎么实现背压?" 答:emit 是挂起函数,下游处理慢时上游发射者被挂起(BufferOverflow 决定缓冲策略),形成天然的拉模式反压,不需要额外策略枚举。

"flowOn 和 withContext 的区别?" 答:flowOn 只影响上游段的执行上下文,是 Flow 官方推荐方式;withContext 写在 flow {} 里换上下文再 emit 会破坏 Flow 不变式直接抛异常。

"Flow 能替换 RxJava 吗?" 答:可覆盖大部分场景(Flow + flowOn + retryWhen + stateIn),但 Rx 的 observeOn 精确到下游、share() 复杂度、成熟算子生态仍有差异;协程目前没有完全对标的 BackpressureStrategy,最接近的是 MutableSharedFlow(replay=0, extraBufferCapacity=N, onBufferOverflow=SUSPEND)。

给正在准备面试的你

1. 项目里统一约定:UI 状态用 StateFlow,一次性事件用 SharedFlow 或 Channel,别用裸 Flow 直接驱动界面。下一节会专门讲这两者的分工。 2. Android 侧把 repeatOnLifecycle 写进模板代码(BaseFragment 的 collectFlow 扩展),避免每人各写一遍。 3. 排查"数据不刷新"先看是冷流漏了 collect,还是热流的 SharingStarted 设成了 WhileSubscribed(5000) 导致上游过早停。 4. 自动化预防:给 Flow 链写 Turbine 单测,断言发射顺序与异常路径,这比手工点界面可靠得多。

把冷热流画成两条并排的竖线:左边"冷流(Flow 默认)"顶端标 collect A 和 collect B 两条箭头指向同一个 flow{} 节点,节点下分出两条独立的生产线,标"每次订阅重新生产";右边"热流(shareIn/stateIn)"画一个居中的方块 SharedFlow,三条收集箭头都从它出发,方块旁标"单一生产源,状态共享"。这张图讲完,冷热差异基本不用再背。

再补一个容易被追问的点:channelFlow 与 callbackFlow 什么时候用。当需要"多个协程并发 emit 到同一条流"时,普通 flow {} 因为发射者约束不允许,应改用 channelFlow {}(内部是 Channel,天然支持多生产者单消费者);需要把回调 API 包装成 Flow 时用 callbackFlow,并在 awaitClose { unregister() } 里解绑,否则回调持有导致泄漏。

复习时别孤立刷题:协程异常处理——上一节讲清异常在 Job 树里的传播,本节换到数据流层面,Flow 的 catch 操作符正是那套异常语义在流上的投影;取消协程会让 collect 停止,emit 途中取消同样是协作式的。


如果这篇文章对你有帮助,欢迎点赞、在看、转发三连。你的支持就是这个系列持续更新的动力。

「Android软件开发面试·从入门到精通」连载系列

上一篇:协程异常处理:从-try-catch-到-CoroutineExceptionHandler

下一篇预告:StateFlow-与-SharedFlow:UI-状态分发的标配

有任何问题欢迎在评论区留言交流。

相关文章
|
1天前
|
缓存 编译器 PHP
第113篇 Compose 与 Kotlin 特性:为什么 Compose 离不开 Kotlin
本文深入解析 Jetpack Compose 的底层机制,揭示其“声明式 UI”背后的三大支柱:带接收者 Lambda、编译器插件(KCP)与稳定性注解(`@Composable`/`@Stable`/`@Immutable`)。重点剖析编译器如何重写函数、实现智能跳过重组,以及常见性能陷阱(如列表无 key、lambda 不稳定)的根因与解法,助你面试直击本质。
27 3
|
1天前
|
API Android开发 开发者
第127篇ViewPager2 与 Fragment:懒加载的正确实现
ViewPager2 内核是 RecyclerView,由 FragmentStateAdapter 托管 Fragment 生命周期,依赖 `getItemId`/`containsItem` 实现稳定复用。旧版 `setUserVisibleHint` 已废弃,应统一用 `repeatOnLifecycle(STARTED)` 控制可见性逻辑。关键在理解“复用机制+生命周期托管”,避免数据错乱、重复加载与内存泄漏。
22 2
|
1天前
|
安全 Java 编译器
第098篇 空安全与 Java 混编:平台类型的风险控制
本文深度剖析Kotlin与Java混编中空安全的“信任边界”:指出平台类型、注解失效与反射泛型三大漏洞,提出“类型收口+运行时断言+架构约定”三层防线,强调用`requireNotNull`替代`!!`、用`sealed`封装多态结果,真正实现NPE可控可追溯。
15 0
|
1天前
|
自然语言处理 Java Android开发
第072篇 中缀表达式与运算符重载:可读性的双刃剑
Kotlin运算符重载比Java更彻底:不仅支持`+ - * /`等映射为`plus`/`minus`等约定函数,还允许任意单参函数通过`infix`声明实现中缀调用(如`a to b`)。核心原则是——符号语义必须与原始含义一致(`+`即相加,`-`即取反或相减),滥用将损害可读性。
27 1
|
1天前
|
并行计算 Java 测试技术
第042篇 CountDownLatch 与 CyclicBarrier:等待与协作
本文深入解析`CountDownLatch`、`CyclicBarrier`与`Semaphore`的核心差异与实战陷阱:前者为一次性事件门闩,后者为可复用集合点,`Semaphore`则基于独占模式限流。重点剖析“计数未归零阻塞”“屏障损坏”“许可泄漏”等生产级坑点,并给出`finally`防护、超时兜底、参与方一致性等落地解法,助你面试直击采分关键。
20 0
第042篇 CountDownLatch 与 CyclicBarrier:等待与协作
|
1天前
|
Java 编译器 Android开发
第106篇 data class 的 equals 与 copy:深浅拷贝的坑
Kotlin `data class` 自动生成 `equals`/`hashCode`/`toString`/`copy`/`componentN`,但**非深拷贝、数组引用比较、可变集合共享引用**——易致UI状态错乱、Diff失效、哈希查找丢失。核心原则:状态类只用不可变类型(`List`/`Set`),数组需覆写 `equals`,`copy` 后禁用就地修改。
14 0
|
1天前
|
存储 Java 编译器
第118篇 const 与编译期常量:延迟初始化之外的第三个选择
`const val` 是 Kotlin 编译期常量,值在编译时内联到各使用点,字节码中无字段、零运行时开销,但**不二进制兼容**——改值后依赖方必须重编译,否则仍用旧值。适用于功能开关、注解参数等编译期确定场景;跨模块配置、敏感信息、需热更新者禁用。
19 0
|
1天前
|
JSON Java API
第102篇 Kotlin 反射与 KClass:运行时元编程
Kotlin/Java反射面试常考深层差异:KClass提供结构化元信息(属性名、泛型、注解默认值),但需额外引入`kotlin-reflect`,增大Android方法数与包体积;Java反射零依赖但元信息原始。关键权衡:能用`::class.java`或编译期方案(如KSP、kotlinx.serialization)就不用Kotlin反射。
14 0
|
2天前
|
监控 IDE Java
第023篇 自定义异常与异常链:生产代码怎么设计错误
本文深入剖析自定义异常与异常链的实战要点,提炼出“动机—上下文—链路—反例”四步公式,直击面试高频追问(如“为何不用RuntimeException?”)。通过典型踩坑案例(丢cause、泛化捕获、序列化失败)和可运行代码,讲清如何科学分类错误、精准携带订单号等上下文、完整保留根因链。强调:异常是可观测性的第一现场,设计不当反增维护成本。
21 0
|
2天前
|
安全 Java 编译器
第013篇 final 的四种用法:变量、方法、类与参数
本文深度解析 Java 中 `final` 的四大语义:变量(引用/值仅赋值一次)、参数(方法内不可重赋值)、方法(禁止重写)、类(禁止继承)。厘清“引用不可变≠内容不可变”“final 不保证线程安全”“this 逸出破坏可见性”等高频误区,结合源码、面试话术与工程实践,助你透彻掌握 final 的原理、坑点与最佳用法。
19 0

热门文章

最新文章