news 2026/8/2 2:26:30

Kotlin 冷流与热流详解

作者头像

张小明

前端开发工程师

1.2k 24
文章封面图
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 配缓存
版权声明: 本文来自互联网用户投稿,该文观点仅代表作者本人,不代表本站立场。本站仅提供信息存储空间服务,不拥有所有权,不承担相关法律责任。如若内容造成侵权/违法违规/事实不符,请联系邮箱:809451989@qq.com进行投诉反馈,一经查实,立即删除!
网站建设 2026/8/2 2:17:45

Fate/Grand Automata:解放双手的FGO自动化战斗助手终极指南

Fate/Grand Automata&#xff1a;解放双手的FGO自动化战斗助手终极指南 【免费下载链接】FGA Auto-battle app for F/GO Android 项目地址: https://gitcode.com/gh_mirrors/fg/FGA 你是否厌倦了在Fate/Grand Order中重复刷取素材的枯燥操作&#xff1f;是否希望将宝贵的…

作者头像 李华
网站建设 2026/8/2 2:16:04

终极Cherry MX键帽3D模型库:免费开源个性化键盘改造完全指南

终极Cherry MX键帽3D模型库&#xff1a;免费开源个性化键盘改造完全指南 【免费下载链接】cherry-mx-keycaps 3D models of Chery MX keycaps 项目地址: https://gitcode.com/gh_mirrors/ch/cherry-mx-keycaps 厌倦了千篇一律的机械键盘外观&#xff1f;想要打造独一无二…

作者头像 李华
网站建设 2026/8/2 2:15:10

忍者龙剑传4豪华版免费下载

下载链接 游戏介绍 终极忍者动作冒险游戏《忍者龙剑传4(NINJA GAIDEN 4)》继承了系列的丰富传统&#xff0c;同时引入了一些创新元素&#xff0c;游戏背景设定在一个近未来的东京&#xff0c;玩家将体验到城市被古老宿敌威胁的紧张情节。年轻的天才忍者八云肩负重任&#xff…

作者头像 李华
网站建设 2026/8/2 2:13:24

【AI大模型进阶】流式输出(Streaming):让AI逐字回复,体验“打字机”的快感

【AI大模型进阶】流式输出(Streaming):让AI逐字回复,体验“打字机”的快感 这是【AI大模型进阶】系列第七十七课,聚焦大模型工程落地与用户体验优化的核心实战技能。 在前几节课程中,我们掌握了OpenAI新版API、Function Calling工具调用、模型超参数调优、损失函数与优…

作者头像 李华
网站建设 2026/8/2 2:11:42

LangChain 实战指南:从上线前检查开始讲

如果你正准备往大模型方向转&#xff0c;《会用LangChain只是起点&#xff0c;能解释失败才算真正入门》这类问题别只看热度。更重要的是判断自己该补哪块能力&#xff0c;以及怎么证明你真的会。 摘要 摘要&#xff1a;很多人学完LangChain&#xff0c;自己跑个RAG Demo挺丝…

作者头像 李华
网站建设 2026/8/2 2:11:40

全网资源宝藏库-夸父资源社

在信息海洋中寻宝&#xff0c;你是否感到疲惫&#xff1f;从热播影视到冷门古籍&#xff0c;从精品课程到稀缺素材&#xff0c;每一次搜索都像一场漫长的远征。现在&#xff0c;这场远征有了终点——欢迎来到「夸父资源社」&#xff0c;一个致力于追逐并汇集全网光芒的资源宝藏…

作者头像 李华