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: 状态2StateFlow 特点:
必须有一个初始值
新订阅者立即收到当前最新值
适合 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: 5SharedFlow 配置参数:
表格
| 参数 | 说明 | 示例 |
|---|---|---|
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 Flow | Hot Flow |
|---|---|---|
| 数据生产时机 | 有订阅者(collect)才开始执行 | 独立于订阅者,自行生产 |
| 多订阅者行为 | 每个订阅者独立执行,数据重新生产 | 所有订阅者共享同一数据源 |
| 数据完整性 | 每个订阅者收到完整序列 | 只能收到订阅之后的数据 |
| 类比 | 音乐 App 按需播放 | 广播电台实时广播 |
| 代表 | flow { },flowOf() | StateFlow,SharedFlow |
金句:冷流是"拉"模式(pull),热流是"推"模式(push)。
二、StateFlow 深度解析(超高频)
Q2: StateFlow 和 LiveData 的区别?
| 特性 | StateFlow | LiveData |
|---|---|---|
| 初始值 | 必须提供初始值 | 可以没有初始值 |
| 主线程安全 | 需要手动确保(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分别是什么?
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可以没有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 要用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 仍在收集,浪费资源repeatOnLifecycle在onStop时自动取消,在onStart时重新订阅
Q10:repeatOnLifecycle和flowWithLifecycle的区别?
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 { } // 并行执行 }repeatOnLifecycle | flowWithLifecycle | |
|---|---|---|
| 作用域 | 整个代码块 | 单个 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)不缓存历史事件,如果没有活跃的订阅者,事件直接丢失。
解决方案:
使用
Channel(推荐用于一次性事件):
kotlin
private val _events = Channel<Event>(Channel.BUFFERED) val events = _events.receiveAsFlow() // 转为 Flow fun sendEvent(event: Event) { viewModelScope.launch { _events.send(event) // 缓冲,不会丢失 } }或使用
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 配缓存 |