当前位置: 首页 > news >正文

Kotlin 冷流与热流详解

Kotlin 冷流与热流详解

核心区别

特性冷流 (Cold Flow)热流 (Hot Flow)
数据生产时机有订阅者才开始生产独立于订阅者,自行生产
订阅者接收数据每个订阅者收到完整序列订阅后才开始接收
多订阅者行为各自独立,数据重新生产共享同一数据源
典型代表flow { }StateFlow,SharedFlow
类比音乐 App 按需播放广播电台实时广播

一、冷流 (Cold Flow)

冷流是按需生产的。只有收集器(collector)开始收集时,Flow 内部的代码才会执行。

1. 基本示例

kotlin

import kotlinx.coroutines.flow.flow import kotlinx.coroutines.delay import kotlinx.coroutines.runBlocking fun coldFlow() = flow { println("Flow 开始执行") for (i in 1..3) { delay(100) emit(i) // 发射数据 } } fun main() = runBlocking { val flow = coldFlow() println("--- 第一个订阅者 ---") flow.collect { println("A: $it") } println("--- 第二个订阅者 ---") flow.collect { println("B: $it") } }

输出:

plain

--- 第一个订阅者 --- Flow 开始执行 A: 1 A: 2 A: 3 --- 第二个订阅者 --- Flow 开始执行 ← 再次执行! B: 1 B: 2 B: 3

每个collect都会触发 Flow 内部代码重新执行,两个订阅者互不影响。

2. 冷流的本质

kotlin

// 这就像调用一个 suspend 函数,每次调用都是新的执行 val result1 = fetchData() // 第一次网络请求 val result2 = fetchData() // 第二次网络请求

3. 常见冷流操作符

kotlin

flow { emit(1) } // 基础构建 flowOf(1, 2, 3) // 固定值 listOf(1,2,3).asFlow() // 集合转 Flow (1..10).asFlow()

二、热流 (Hot Flow)

热流是独立存在的,数据生产不依赖于订阅者。新订阅者只能收到订阅之后的数据。

1. StateFlow

StateFlow是一个有状态的热流,始终持有一个最新值。

kotlin

import kotlinx.coroutines.flow.MutableStateFlow import kotlinx.coroutines.flow.StateFlow import kotlinx.coroutines.launch import kotlinx.coroutines.runBlocking class ViewModel { private val _uiState = MutableStateFlow("初始状态") val uiState: StateFlow<String> = _uiState fun updateState(newState: String) { _uiState.value = newState } } fun main() = runBlocking { val vm = ViewModel() // 订阅者1:从一开始就订阅 val job1 = launch { vm.uiState.collect { println("订阅者1: $it") } } delay(50) vm.updateState("状态1") delay(50) // 订阅者2:中途订阅,只会收到当前最新值及后续值 val job2 = launch { vm.uiState.collect { println("订阅者2: $it") } } delay(50) vm.updateState("状态2") delay(100) job1.cancel() job2.cancel() }

输出:

plain

订阅者1: 初始状态 订阅者1: 状态1 订阅者2: 状态1 ← 订阅者2只收到当前最新值 订阅者1: 状态2 订阅者2: 状态2

StateFlow 特点:

  • 必须有一个初始值

  • 新订阅者立即收到当前最新值

  • 适合 UI 状态管理(如 Loading/Success/Error)

2. SharedFlow

SharedFlow是一个无状态的热流,更灵活,可配置缓存。

kotlin

import kotlinx.coroutines.flow.MutableSharedFlow import kotlinx.coroutines.launch import kotlinx.coroutines.runBlocking fun main() = runBlocking { // replay=2:新订阅者会收到最近2个值 val sharedFlow = MutableSharedFlow<Int>(replay = 2) // 发射一些数据(此时无订阅者,数据会丢失或缓存) sharedFlow.emit(1) sharedFlow.emit(2) sharedFlow.emit(3) println("--- 订阅者1加入 ---") val job1 = launch { sharedFlow.collect { println("订阅者1: $it") } } delay(50) sharedFlow.emit(4) delay(50) println("--- 订阅者2加入 ---") val job2 = launch { sharedFlow.collect { println("订阅者2: $it") } } delay(50) sharedFlow.emit(5) delay(100) job1.cancel() job2.cancel() }

输出:

plain

--- 订阅者1加入 --- 订阅者1: 2 ← replay=2,收到最近2个:2, 3 订阅者1: 3 订阅者1: 4 --- 订阅者2加入 --- 订阅者2: 3 ← replay=2,收到最近2个:3, 4 订阅者2: 4 订阅者1: 5 订阅者2: 5

SharedFlow 配置参数:

表格

参数说明示例
replay新订阅者能收到的历史值数量replay=0不缓存,replay=1类似 StateFlow
extraBufferCapacity额外缓冲,超出时策略由onBufferOverflow决定
onBufferOverflow缓冲溢出策略:SUSPEND(挂起)/DROP_OLDEST(丢弃最旧)/DROP_LATEST(丢弃最新)

三、对比图

plain

时间线 ──────────────────────────────────────► 冷流 (Flow): 收集者1: [1]──[2]──[3]──[4]──[5] 收集者2: [1]──[2]──[3]──[4]──[5] ← 独立重新执行 热流 (StateFlow/SharedFlow): 数据源: [1]──[2]──[3]──[4]──[5] 收集者1: [1]──[2]──[3]──[4]──[5] 收集者2: [2]──[3]──[4]──[5] ← 从订阅时刻开始,共享数据源

四、实际应用场景

冷流场景:一次性数据获取

kotlin

// 网络请求、数据库查询 fun getUserProfile(userId: String): Flow<User> = flow { val user = api.fetchUser(userId) // 每次 collect 都会重新请求 emit(user) } // 使用 viewModelScope.launch { getUserProfile("123").collect { user -> updateUI(user) } }

热流场景:UI 状态 & 事件

kotlin

class NewsViewModel : ViewModel() { // StateFlow:UI 状态(始终有值) private val _newsState = MutableStateFlow<NewsUiState>(NewsUiState.Loading) val newsState: StateFlow<NewsUiState> = _newsState.asStateFlow() // SharedFlow:一次性事件(如 Toast、导航) private val _events = MutableSharedFlow<NewsEvent>() // replay=0 val events: SharedFlow<NewsEvent> = _events.asSharedFlow() fun loadNews() { viewModelScope.launch { _newsState.value = NewsUiState.Loading try { val news = repository.fetchNews() _newsState.value = NewsUiState.Success(news) } catch (e: Exception) { _newsState.value = NewsUiState.Error(e.message) _events.emit(NewsEvent.ShowToast("加载失败")) // 一次性事件 } } } }

Activity/Fragment 中收集:

kotlin

class NewsFragment : Fragment() { override fun onViewCreated(view: View, savedInstanceState: Bundle?) { // 收集状态:使用 repeatOnLifecycle 避免后台耗电 viewLifecycleOwner.lifecycleScope.launch { repeatOnLifecycle(Lifecycle.State.STARTED) { viewModel.newsState.collect { state -> when (state) { is NewsUiState.Loading -> showLoading() is NewsUiState.Success -> showNews(state.data) is NewsUiState.Error -> showError() } } } } // 收集事件:不需要 repeatOnLifecycle,事件不能丢 viewLifecycleOwner.lifecycleScope.launch { viewModel.events.collect { event -> when (event) { is NewsEvent.ShowToast -> Toast.makeText(context, event.msg, Toast.LENGTH_SHORT).show() is NewsEvent.Navigate -> findNavController().navigate(event.direction) } } } } }

五、冷流转热流:shareIn/stateIn

有时你需要把冷流转成热流,比如多个 UI 组件共享同一个数据流:

kotlin

class Repository @Inject constructor(private val api: Api) { // 冷流:每次 collect 都会触发网络请求 fun fetchData(): Flow<Data> = flow { emit(api.fetchData()) } // 转热流:在 ViewModel 中使用 val hotData: Flow<Data> = fetchData() .stateIn( scope = viewModelScope, started = SharingStarted.WhileSubscribed(5000), // 5秒内无订阅者则停止 initialValue = Data.Empty ) }

SharingStarted策略:

表格

策略行为
Eagerly立即开始,永不停止
Lazily第一个订阅者到来时开始,永不停止
WhileSubscribed(timeout)有订阅者时活跃,无订阅者后等待 timeout 停止(最省资源)

六、总结

表格

场景选择
网络请求、数据库查询冷流(flow { })
UI 状态(Loading/Content/Error)StateFlow
一次性事件(Toast、SnackBar、导航)SharedFlow(replay=0)
多个订阅者共享数据冷流 +shareIn/stateIn

核心记忆口诀:冷流按需重新生产,热流实时共享广播


Kotlin 冷流与热流 — 高频知识点


一、核心概念(必问)

Q1: 冷流和热流的本质区别是什么?

维度Cold FlowHot Flow
数据生产时机有订阅者(collect)才开始执行独立于订阅者,自行生产
多订阅者行为每个订阅者独立执行,数据重新生产所有订阅者共享同一数据源
数据完整性每个订阅者收到完整序列只能收到订阅之后的数据
类比音乐 App 按需播放广播电台实时广播
代表flow { },flowOf()StateFlow,SharedFlow

金句:冷流是"拉"模式(pull),热流是"推"模式(push)。


二、StateFlow 深度解析(超高频)

Q2: StateFlow 和 LiveData 的区别?

特性StateFlowLiveData
初始值必须提供初始值可以没有初始值
主线程安全需要手动确保(Dispatchers.Main自动在主线程观察
生命周期感知不感知,需配合repeatOnLifecycle自动感知生命周期
数据去重distinctUntilChanged()需手动调用自动去重(值不变不通知)
版本支持需要 Coroutines 依赖Android 原生支持
转换操作丰富的 Flow 操作符map,switchMap

kotlin

// StateFlow 去重需要手动处理 stateFlow .distinctUntilChanged() // 值不变时跳过 .collect { }

Q3: 为什么用 StateFlow 替代 LiveData?

  1. 一致性:Flow 操作符更丰富(debounce,flatMapLatest,combine等)

  2. 测试性:Flow 不依赖 Android 生命周期,单元测试更方便

  3. 组合能力:多个 Flow 可以用combine,zip,merge

  4. Kotlin 优先:与协程深度集成

Q4: StateFlow 的value赋值是线程安全的吗?

不是线程安全的!必须在单线程中更新(通常主线程):

kotlin

// ❌ 错误:可能在后台线程更新 viewModelScope.launch(Dispatchers.IO) { _state.value = newValue // 可能崩溃! } // ✅ 正确 viewModelScope.launch { _state.value = newValue // 默认 Dispatchers.Main }

如果需要在后台计算后更新,用update函数更安全:

kotlin

_state.update { it.copy(isLoading = true) }

三、SharedFlow 深度解析(高频)

Q5: SharedFlow 的replayextraBufferCapacityonBufferOverflow分别是什么?

kotlin

val sharedFlow = MutableSharedFlow<Int>( replay = 2, // 新订阅者能收到的历史值数量 extraBufferCapacity = 3, // 额外缓存容量 onBufferOverflow = BufferOverflow.DROP_OLDEST // 溢出策略 )

表格

参数作用默认值
replay缓存最近 N 个值给新订阅者0
extraBufferCapacity超出 replay 的额外缓存0
onBufferOverflow缓存满时的处理策略SUSPEND

溢出策略:

  • SUSPEND:挂起发送者(默认,可能阻塞)

  • DROP_OLDEST:丢弃最旧的数据

  • DROP_LATEST:丢弃最新的数据

Q6: SharedFlow 和 StateFlow 的关系?

kotlin

// StateFlow 是 SharedFlow 的特化版本 interface StateFlow<out T> : SharedFlow<T>

等价关系:

kotlin

// 以下两者等价 val stateFlow = MutableStateFlow(initialValue) val sharedFlow = MutableSharedFlow<T>( replay = 1, onBufferOverflow = BufferOverflow.DROP_OLDEST ).apply { tryEmit(initialValue) }

关键区别:

  • StateFlow必须有初始值,SharedFlow可以没有

  • StateFlowvalue属性可以直接读取当前值

  • SharedFlow更适合事件流(如 Toast、导航事件)


四、冷流转热流(高频)

Q7:shareInstateIn的区别?

kotlin

// shareIn:转为 SharedFlow val hotFlow = coldFlow.shareIn( scope = viewModelScope, started = SharingStarted.WhileSubscribed(5000), replay = 1 ) // stateIn:转为 StateFlow(必须有初始值) val stateFlow = coldFlow.stateIn( scope = viewModelScope, started = SharingStarted.WhileSubscribed(5000), initialValue = emptyList() )

Q8:SharingStarted三种策略的区别?

策略行为适用场景
Eagerly立即开始,永不停止应用全局数据
Lazily第一个订阅者来时开始,永不停止启动后持续需要的数据
WhileSubscribed(timeout)有订阅者时活跃,无订阅者 timeout 后停止最常用,省资源

WhileSubscribed(5000)中的 5000ms 是** grace period**:最后一个订阅者离开后,等待 5 秒再停止,避免配置变更(如旋转屏幕)时重复初始化。


五、生命周期与收集(高频)

Q9: 为什么收集 Flow 要用repeatOnLifecycle

kotlin

// ❌ 错误:后台持续收集,浪费资源,可能崩溃 lifecycleScope.launch { viewModel.state.collect { updateUI(it) } } // ✅ 正确:生命周期感知 lifecycleScope.launch { repeatOnLifecycle(Lifecycle.State.STARTED) { viewModel.state.collect { updateUI(it) } } }

问题背景:

  • lifecycleScope.launch在 Activity/Fragment 整个生命周期运行

  • 当页面进入后台(onStop),Flow 仍在收集,浪费资源

  • repeatOnLifecycleonStop时自动取消,在onStart时重新订阅

Q10:repeatOnLifecycleflowWithLifecycle的区别?

kotlin

// 方式1:repeatOnLifecycle(代码块级别) lifecycleScope.launch { repeatOnLifecycle(Lifecycle.State.STARTED) { flow1.collect { } flow2.collect { } // 顺序执行,flow1 不完成不会执行 flow2 } } // 方式2:flowWithLifecycle(流级别) lifecycleScope.launch { flow1.flowWithLifecycle(lifecycle, Lifecycle.State.STARTED) .collect { } flow2.flowWithLifecycle(lifecycle, Lifecycle.State.STARTED) .collect { } // 并行执行 }
repeatOnLifecycleflowWithLifecycle
作用域整个代码块单个 Flow
多个 Flow顺序执行(一个挂起,后面的不执行)可并行
使用场景单个 Flow 或需要顺序执行多个 Flow 并行收集

六、背压处理(中高频)

Q11: Flow 的背压是什么?如何处理?

背压:生产者速度 > 消费者速度,数据堆积。

kotlin

// 生产者每 100ms 发一个,消费者每 300ms 处理一个 flow { for (i in 1..100) { delay(100) emit(i) } }.collect { value -> delay(300) // 处理慢 println(value) }

解决方案:

操作符行为适用场景
buffer()缓冲数据,不阻塞生产者允许一定延迟
conflate()只保留最新值,丢弃中间值只关心最新状态
collectLatest { }有新值时取消旧值处理搜索输入等
flatMapLatest { }类似 collectLatest,但用于转换搜索请求

kotlin

// 示例:搜索框防抖 + 取消旧请求 searchQueryFlow .debounce(300) // 停止输入 300ms 后才触发 .flatMapLatest { query -> searchRepository.search(query) // 新搜索来时取消旧请求 } .collect { results -> updateUI(results) }

七、常见陷阱(面试加分项)

Q12: 以下代码有什么问题?

kotlin

// ❌ 问题代码 class MyViewModel : ViewModel() { val data = repository.fetchData() // 冷流 .stateIn(viewModelScope, SharingStarted.Lazily, emptyList()) }

问题fetchData()MyViewModel实例化时就被调用了(即使无人订阅)!

原因stateIn的参数是Flow,但fetchData()先执行返回 Flow,然后传给stateIn

修正

kotlin

// ✅ 正确:使用 lazy 或函数 class MyViewModel : ViewModel() { val data: StateFlow<List<Data>> = repository.fetchData() .stateIn(viewModelScope, SharingStarted.Lazily, emptyList()) } // 实际上上面的写法在 Kotlin 属性初始化时也会立即执行 fetchData() // 更好的方式: class MyViewModel : ViewModel() { val data by lazy { repository.fetchData() .stateIn(viewModelScope, SharingStarted.Lazily, emptyList()) } }

Q13:SharedFlow用于事件时,为什么可能丢失事件?

kotlin

// ❌ 问题:事件可能丢失 viewModelScope.launch { _events.emit(NavigateToDetail) // 如果此时无订阅者,事件丢失! }

原因SharedFlow(replay=0)不缓存历史事件,如果没有活跃的订阅者,事件直接丢失。

解决方案

  1. 使用Channel(推荐用于一次性事件):

kotlin

private val _events = Channel<Event>(Channel.BUFFERED) val events = _events.receiveAsFlow() // 转为 Flow fun sendEvent(event: Event) { viewModelScope.launch { _events.send(event) // 缓冲,不会丢失 } }
  1. 或使用SharedFlow增加 replay:

kotlin

private val _events = MutableSharedFlow<Event>(extraBufferCapacity = 1)

八、综合代码题

题目:实现一个带搜索、防抖、 loading 状态的 ViewModel

kotlin

class SearchViewModel( private val repository: SearchRepository ) : ViewModel() { private val _searchQuery = MutableStateFlow("") // 对外暴露只读 StateFlow val uiState: StateFlow<SearchUiState> = _searchQuery .debounce(300) // 防抖 300ms .filter { it.isNotBlank() } // 空内容不搜索 .flatMapLatest { query -> // 新搜索取消旧请求 flow { emit(SearchUiState.Loading) try { val results = repository.search(query) emit(SearchUiState.Success(results)) } catch (e: Exception) { emit(SearchUiState.Error(e.message)) } } } .stateIn( scope = viewModelScope, started = SharingStarted.WhileSubscribed(5000), initialValue = SearchUiState.Idle ) fun onSearchQueryChange(query: String) { _searchQuery.value = query } } sealed class SearchUiState { object Idle : SearchUiState() object Loading : SearchUiState() data class Success(val data: List<SearchResult>) : SearchUiState() data class Error(val message: String?) : SearchUiState() }

Activity 中收集:

kotlin

class SearchActivity : AppCompatActivity() { override fun onCreate(savedInstanceState: Bundle?) { // ... lifecycleScope.launch { repeatOnLifecycle(Lifecycle.State.STARTED) { viewModel.uiState.collect { state -> when (state) { is SearchUiState.Idle -> showIdle() is SearchUiState.Loading -> showLoading() is SearchUiState.Success -> showResults(state.data) is SearchUiState.Error -> showError(state.message) } } } } searchEditText.doAfterTextChanged { viewModel.onSearchQueryChange(it?.toString().orEmpty()) } } }

九、速记口诀

知识点口诀
冷流 vs 热流冷流按需重新跑,热流实时共享好
StateFlow有初始值、读当前值、UI状态管
SharedFlow无初始值、配缓存、事件通知用
stateIn/shareIn冷流转热流,WhileSubscribed 最省流
生命周期repeatOnLifecycle,STARTED 时收,STOPPED 时丢
背压buffer 做缓冲,conflate 留最新,collectLatest 取消旧
事件不丢失Channel 来缓冲,SharedFlow 配缓存
http://www.jsqmd.com/news/1310604/

相关文章:

  • 腾讯混元AngelSpec投机解码框架深度解析:MTP+块扩散双Draft策略与D-cut高并发吞吐优化
  • Unity权限问题深度解析:从UAC机制到项目路径规范
  • 小红书面经
  • SSDTTime终极指南:一键生成黑苹果完美SSDT补丁的完整教程
  • Java 微服务架构设计与 Spring Cloud 实战:基于 OpenFeign 与 Resilience4j 的容错底座
  • 2026年8月湖南省联通300M单宽带避坑指南!小白怎么选_ - 找卡家园
  • 如何让PS3手柄在Windows上完美重生:5种HID模式+智能震动+蓝牙连接全解析
  • 广州中小微企业主经济犯罪律师推荐:【法纳刑辩】成效斐然 - 秋山寄远
  • 2026年8月湖南省电信2000M融合宽带我的真实踩坑与实操 - 找卡家园
  • 计算机考研408高效备考攻略:从核心概念到实战策略
  • RS232转RS485/422转换器:原理、选型与工业通信组网实战指南
  • 猫抓浏览器资源嗅探扩展:高性能网络请求拦截与媒体资源捕获架构
  • 可持久化线段树(Persistent Segment Tree)详解
  • Unlimited-OCR 部署运行(11/13):环境变量块超限导致 spawn 子进程崩溃
  • 围棋AI智能教练KaTrain:免费开源工具助你快速提升棋力的终极指南
  • DSL崛起挑战Python生态:领域特定语言如何重塑开发者技能与薪酬
  • 2026下半年珠海古堡瓷砖实力厂商甄选:广东乐高尚品陶瓷以全案美学破局同质化 - 装修教育财税推荐2026
  • 如何用Zettelkasten打造你的第二大脑:免费开源知识管理终极指南
  • 从零打造64x64高密度全彩LED矩阵:HUB75接口驱动与ESP32实战
  • 变频无刷电机低噪音控制实战:从FOC原理到35分贝静音实现
  • 2026年8月湖南省电信2000M融合宽带实测对比宽带怎么选? - 找卡家园
  • RTL8019AS、555,DS3231等芯片的跳线方式和引脚定义
  • GD32F450Z移植LittleFS:构建掉电安全的SPI Flash存储方案
  • 2026年8月湖南省联通500M单宽带申请避坑实录 - 找卡家园
  • DSO Quad示波器固件构建:从STM32开发环境搭建到固件烧录全流程
  • [特殊字符]2026奇幻新片《星光继承者:暗黑仙境》 4K DVHDR超清画质,内封简中字幕 夸克网盘分享,支持在线观看!
  • Windows驱动管理的终极免费工具:Driver Store Explorer完整使用指南
  • 树莓派5寸HDMI显示屏选购配置全攻略:从参数解析到实战优化
  • 昆明理工大学信息工程与自动化2026届硕士就业流向全景解析
  • 2026年8月湖南省怀化市移动单宽带小白避坑指南 - 找卡家园