Kotlin 冷流与热流详解

Kotlin 冷流与热流详解
Kotlin 冷流与热流详解核心区别特性冷流 (Cold Flow)热流 (Hot Flow)数据生产时机有订阅者才开始生产独立于订阅者自行生产订阅者接收数据每个订阅者收到完整序列订阅后才开始接收多订阅者行为各自独立数据重新生产共享同一数据源典型代表flow { }StateFlow,SharedFlow类比音乐 App 按需播放广播电台实时广播一、冷流 (Cold Flow)冷流是按需生产的。只有收集器collector开始收集时Flow 内部的代码才会执行。1. 基本示例kotlinimport 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. 常见冷流操作符kotlinflow { emit(1) } // 基础构建 flowOf(1, 2, 3) // 固定值 listOf(1,2,3).asFlow() // 集合转 Flow (1..10).asFlow()二、热流 (Hot Flow)热流是独立存在的数据生产不依赖于订阅者。新订阅者只能收到订阅之后的数据。1. StateFlowStateFlow是一个有状态的热流始终持有一个最新值。kotlinimport 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: StateFlowString _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: 状态2StateFlow 特点必须有一个初始值新订阅者立即收到当前最新值适合 UI 状态管理如 Loading/Success/Error2. SharedFlowSharedFlow是一个无状态的热流更灵活可配置缓存。kotlinimport kotlinx.coroutines.flow.MutableSharedFlow import kotlinx.coroutines.launch import kotlinx.coroutines.runBlocking fun main() runBlocking { // replay2新订阅者会收到最近2个值 val sharedFlow MutableSharedFlowInt(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 ← replay2收到最近2个2, 3 订阅者1: 3 订阅者1: 4 --- 订阅者2加入 --- 订阅者2: 3 ← replay2收到最近2个3, 4 订阅者2: 4 订阅者1: 5 订阅者2: 5SharedFlow 配置参数表格参数说明示例replay新订阅者能收到的历史值数量replay0不缓存replay1类似 StateFlowextraBufferCapacity额外缓冲超出时策略由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): FlowUser flow { val user api.fetchUser(userId) // 每次 collect 都会重新请求 emit(user) } // 使用 viewModelScope.launch { getUserProfile(123).collect { user - updateUI(user) } }热流场景UI 状态 事件kotlinclass NewsViewModel : ViewModel() { // StateFlowUI 状态始终有值 private val _newsState MutableStateFlowNewsUiState(NewsUiState.Loading) val newsState: StateFlowNewsUiState _newsState.asStateFlow() // SharedFlow一次性事件如 Toast、导航 private val _events MutableSharedFlowNewsEvent() // replay0 val events: SharedFlowNewsEvent _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 中收集kotlinclass 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 组件共享同一个数据流kotlinclass Repository Inject constructor(private val api: Api) { // 冷流每次 collect 都会触发网络请求 fun fetchData(): FlowData flow { emit(api.fetchData()) } // 转热流在 ViewModel 中使用 val hotData: FlowData fetchData() .stateIn( scope viewModelScope, started SharingStarted.WhileSubscribed(5000), // 5秒内无订阅者则停止 initialValue Data.Empty ) }SharingStarted策略表格策略行为Eagerly立即开始永不停止Lazily第一个订阅者到来时开始永不停止WhileSubscribed(timeout)有订阅者时活跃无订阅者后等待 timeout 停止最省资源六、总结表格场景选择网络请求、数据库查询冷流(flow { })UI 状态Loading/Content/ErrorStateFlow一次性事件Toast、SnackBar、导航SharedFlow(replay0)多个订阅者共享数据冷流 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一致性Flow 操作符更丰富debounce,flatMapLatest,combine等测试性Flow 不依赖 Android 生命周期单元测试更方便组合能力多个 Flow 可以用combine,zip,mergeKotlin 优先与协程深度集成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 的replay、extraBufferCapacity、onBufferOverflow分别是什么kotlinval sharedFlow MutableSharedFlowInt( replay 2, // 新订阅者能收到的历史值数量 extraBufferCapacity 3, // 额外缓存容量 onBufferOverflow BufferOverflow.DROP_OLDEST // 溢出策略 )表格参数作用默认值replay缓存最近 N 个值给新订阅者0extraBufferCapacity超出 replay 的额外缓存0onBufferOverflow缓存满时的处理策略SUSPEND溢出策略SUSPEND挂起发送者默认可能阻塞DROP_OLDEST丢弃最旧的数据DROP_LATEST丢弃最新的数据Q6: SharedFlow 和 StateFlow 的关系kotlin// StateFlow 是 SharedFlow 的特化版本 interface StateFlowout T : SharedFlowT等价关系kotlin// 以下两者等价 val stateFlow MutableStateFlow(initialValue) val sharedFlow MutableSharedFlowT( replay 1, onBufferOverflow BufferOverflow.DROP_OLDEST ).apply { tryEmit(initialValue) }关键区别StateFlow必须有初始值SharedFlow可以没有StateFlow的value属性可以直接读取当前值SharedFlow更适合事件流如 Toast、导航事件四、冷流转热流高频Q7:shareIn和stateIn的区别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 要用repeatOnLifecyclekotlin// ❌ 错误后台持续收集浪费资源可能崩溃 lifecycleScope.launch { viewModel.state.collect { updateUI(it) } } // ✅ 正确生命周期感知 lifecycleScope.launch { repeatOnLifecycle(Lifecycle.State.STARTED) { viewModel.state.collect { updateUI(it) } } }问题背景lifecycleScope.launch在 Activity/Fragment 整个生命周期运行当页面进入后台onStopFlow 仍在收集浪费资源repeatOnLifecycle在onStop时自动取消在onStart时重新订阅Q10:repeatOnLifecycle和flowWithLifecycle的区别kotlin// 方式1repeatOnLifecycle代码块级别 lifecycleScope.launch { repeatOnLifecycle(Lifecycle.State.STARTED) { flow1.collect { } flow2.collect { } // 顺序执行flow1 不完成不会执行 flow2 } } // 方式2flowWithLifecycle流级别 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: StateFlowListData 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(replay0)不缓存历史事件如果没有活跃的订阅者事件直接丢失。解决方案使用Channel推荐用于一次性事件kotlinprivate val _events ChannelEvent(Channel.BUFFERED) val events _events.receiveAsFlow() // 转为 Flow fun sendEvent(event: Event) { viewModelScope.launch { _events.send(event) // 缓冲不会丢失 } }或使用SharedFlow增加 replaykotlinprivate val _events MutableSharedFlowEvent(extraBufferCapacity 1)八、综合代码题题目实现一个带搜索、防抖、 loading 状态的 ViewModelkotlinclass SearchViewModel( private val repository: SearchRepository ) : ViewModel() { private val _searchQuery MutableStateFlow() // 对外暴露只读 StateFlow val uiState: StateFlowSearchUiState _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: ListSearchResult) : SearchUiState() data class Error(val message: String?) : SearchUiState() }Activity 中收集kotlinclass 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 最省流生命周期repeatOnLifecycleSTARTED 时收STOPPED 时丢背压buffer 做缓冲conflate 留最新collectLatest 取消旧事件不丢失Channel 来缓冲SharedFlow 配缓存