ARTICLE DETAIL

资讯详情

深耕编程入门与网站建设的一线实战洞察。

Kotlin SharedFlowStateFlow 热流到底有多热?

Kotlin SharedFlowStateFlow 热流到底有多热? 前言协程系列文章一个小故事讲明白进程、线程、Kotlin 协程到底啥关系少年你可知 Kotlin 协程最初的样子讲真Kotlin 协程的挂起/恢复没那么神秘(故事篇)讲真Kotlin 协程的挂起/恢复没那么神秘(原理篇)Kotlin 协程调度切换线程是时候解开真相了Kotlin 协程之线程池探索之旅(与Java线程池PK)Kotlin 协程之取消与异常处理探索之旅(上)Kotlin 协程之取消与异常处理探索之旅(下)来跟我一起撸Kotlin runBlocking/launch/join/async/delay 原理使用继续来同我一起撸Kotlin Channel 深水区Kotlin 协程 Select看我如何多路复用Kotlin Sequence 是时候派上用场了Kotlin Flow啊你将流向何方Kotlin Flow 背压和线程切换竟然如此相似Kotlin SharedFlowStateFlow 热流到底有多热前面分析的都是冷流冷热是对应的有冷就有热本篇将重点分析热流SharedFlowStateFlow的使用及其原理探究其热度。通过本篇文章你将了解到冷流与热流区别SharedFlow 使用方式与应用场景SharedFlow 原理不一样的角度分析StateFlow 使用方式与应用场景StateFlow 原理一看就会StateFlow/SharedFlow/LiveData 区别与应用1. 冷流与热流区别2. SharedFlow 使用方式与应用场景使用方式流的两端分别是消费者(观察者/订阅者)生产者(被观察者/被订阅者)因此只需要关注两端的行为即可。1. 生产者先发送数据fun test1() { runBlocking { //构造热流 val flow MutableSharedFlowString() //发送数据(生产者) flow.emit(hello world) //开启协程 GlobalScope.launch { //接收数据(消费者) flow.collect { println(collect: $it) } } } }Q先猜测一下结果A没有任何打印我们猜测生产者先发送了数据因为此时消费者还没来得及接收因此数据被丢弃了。2. 生产者延后发送数据我们很容易想到变换一下时机让消费者先注册等待fun test2() { runBlocking { //构造热流 val flow MutableSharedFlowString() //开启协程 GlobalScope.launch { //接收数据(消费者) flow.collect { println(collect: $it) } } //发送数据(生产者) delay(200)//保证消费者已经注册上 flow.emit(hello world) } }这个时候消费者成功打印数据。3. 历史数据的保留(重放)虽然2的方式连通了生产者和消费者但是你对1的失败耿耿于怀觉得SharedFlow有点弱啊限制有点狠LiveData每次新的观察者到来都能收到当前的数据而SharedFlow不行。实际上SharedFlow对于历史数据的重放比LiveData更强大LiveData始终只有个值也就是每次只重放1个值而SharedFlow可配置重放任意值当然不能超过Int的范围。换一下使用姿势fun test3() { runBlocking { //构造热流 val flow MutableSharedFlowString(1) //发送数据(生产者) flow.emit(hello world) //开启协程 GlobalScope.launch { //接收数据(消费者) flow.collect { println(collect: $it) } } } }此时达成的效果与2一致MutableSharedFlow(1)表示设定生产者保留1个值当有新的消费者来了之后将会获取到这个保留的值。当然也可以保留更多的值fun test3() { runBlocking { //构造热流 val flow MutableSharedFlowString(4) //发送数据(生产者) flow.emit(hello world1) flow.emit(hello world2) flow.emit(hello world3) flow.emit(hello world4) //开启协程 GlobalScope.launch { //接收数据(消费者) flow.collect { println(collect: $it) } } } }此时消费者将打印出hell world1~hello world4此时也说明了不管有没有消费者生产者都生产了数据由此说明SharedFlow 是热流4. collect是挂起函数在2里我们开启了协程去执行消费者逻辑flow.collect不单独开启协程执行会怎样fun test4() { runBlocking { //构造热流 val flow MutableSharedFlowString() //接收数据(消费者) flow.collect { println(collect: $it) } println(start emit)//① flow.emit(hello world) } }最后发现①没打印出来因为collect是挂起函数此时由于生产者还没来得及生产数据消费者调用collect时发现没数据后便挂起协程。因此生产者和消费者要处在不同的协程里5. emit是挂起函数消费者要等待生产者生产数据所以collect设计为挂起函数反过来生产者是否要等待消费者消费完数据才进行下一次emit呢fun test5() { runBlocking { //构造热流 val flow MutableSharedFlowString() //开启协程 GlobalScope.launch { //接收数据(消费者) flow.collect { delay(2000) println(collect: $it) } } //发送数据(生产者) delay(200)//保证消费者先执行 println(emit 1 ${System.currentTimeMillis()}) flow.emit(hello world1) println(emit 2 ${System.currentTimeMillis()}) flow.emit(hello world2) println(emit 3 ${System.currentTimeMillis()}) flow.emit(hello world3) println(emit 4 ${System.currentTimeMillis()}) flow.emit(hello world4) } }从打印可以看出生产者每次emit都需要等待消费者消费完成之后才能进行下次emit。6. 缓存的设定在之前分析Flow的时候有说过Flow的背压问题以及使用Buffer来解决它同样的在SharedFlow里也有缓存的概念。fun test6() { runBlocking { //构造热流 val flow MutableSharedFlowString(0, 10) //开启协程 GlobalScope.launch { //接收数据(消费者) flow.collect { delay(2000) println(collect: $it) } } //发送数据(生产者) delay(200)//保证消费者先执行 println(emit 1 ${System.currentTimeMillis()}) flow.emit(hello world1) println(emit 2 ${System.currentTimeMillis()}) flow.emit(hello world2) println(emit 3 ${System.currentTimeMillis()}) flow.emit(hello world3) println(emit 4 ${System.currentTimeMillis()}) flow.emit(hello world4) } }MutableSharedFlow(0, 10) 第2个参数10表示额外的缓存大小为10生产者通过emit先将数据放到缓存里此时它并没有被消费者的速度拖累。7. 重放与额外缓存个数public fun T MutableSharedFlow( replay: Int 0,//重放个数 extraBufferCapacity: Int 0,//额外的缓存个数 onBufferOverflow: BufferOverflow BufferOverflow.SUSPEND ):重放主要用来给新进的消费者重放特定个数的历史数据而额外的缓存个数是为了应付背压问题总的缓存个数重放个数额外的缓存个数。应用场景如有以下需求可用SharedFlow需要重放历史数据可以配置缓存需要重复发射/接收相同的值3. SharedFlow 原理不一样的角度分析带着问题找答案重点关注的无非是emit和collect函数它俩都是挂起函数而是否挂起取决于是否满足条件。同时生产者和消费出现的时机也会影响这个条件因此列举生产者、消费者出现的时机即可。只有生产者当只有生产者没有消费者此时生产者调用emit会挂起协程吗如果不是那么什么情况会挂起从emit函数源码入手override suspend fun emit(value: T) { //如果发射成功则直接退出函数 if (tryEmit(value)) return // fast-path //否则挂起协程 emitSuspend(value) }先看tryEmit(xx)override fun tryEmit(value: T): Boolean { var resumes: ArrayContinuationUnit? EMPTY_RESUMES val emitted kotlinx.coroutines.internal.synchronized(this) { //尝试emit if (tryEmitLocked(value)) { //遍历所有消费者找到需要唤醒的消费者协程 resumes findSlotsToResumeLocked(resumes) true } else { false } } //恢复消费者协程 for (cont in resumes) cont?.resume(Unit) //emittedtrue表示发射成功 return emitted } private fun tryEmitLocked(value: T): Boolean { //nCollectors 表示消费者个数若是没有消费者则无论如何都会发射成功 if (nCollectors 0) return tryEmitNoCollectorsLocked(value) // always returns true if (bufferSize bufferCapacity minCollectorIndex replayIndex) { //如果缓存已经满并且有消费者没有消费最旧的数据(replayIndex)则进入此处 when (onBufferOverflow) { //挂起生产者 BufferOverflow.SUSPEND - return false // will suspend //直接丢弃最新数据认为发射成功 BufferOverflow.DROP_LATEST - return true // just drop incoming //丢弃最旧的数据 BufferOverflow.DROP_OLDEST - {} // force enqueue drop oldest instead } } //将数据加入到缓存队列里 enqueueLocked(value) //缓存数据队列长度 bufferSize // value was added to buffer // drop oldest from the buffer if it became more than bufferCapacity if (bufferSize bufferCapacity) dropOldestLocked() // keep replaySize not larger that needed if (replaySize replay) { // increment replayIndex by one updateBufferLocked(replayIndex 1, minCollectorIndex, bufferEndIndex, queueEndIndex) } return true } private fun tryEmitNoCollectorsLocked(value: T): Boolean { kotlinx.coroutines.assert { nCollectors 0 } //没有设置重放则直接退出丢弃发射的值 if (replay 0) return true // no need to replay, just forget it now //加入到缓存里 enqueueLocked(value) // enqueue to replayCache bufferSize // value was added to buffer // drop oldest from the buffer if it became more than replay //若是超出了重放个数则丢弃最旧的值 if (bufferSize replay) dropOldestLocked() minCollectorIndex head bufferSize // a default value (max allowed) //发射成功 return true }再看emitSuspend(value)private suspend fun emitSuspend(value: T) suspendCancellableCoroutineUnit sc{ cont - var resumes: ArrayContinuationUnit? EMPTY_RESUMES val emitter kotlinx.coroutines.internal.synchronized(this) lock{ ... //构造为Emitter加入到buffer里 SharedFlowImpl.Emitter(this, head totalSize, value, cont).also { enqueueLocked(it) //单独记录挂起的emit queueSize // added to queue of waiting emitters // synchronous shared flow might rendezvous with waiting emitter if (bufferCapacity 0) resumes findSlotsToResumeLocked(resumes) } } }用图表示整个emit流程现在可以回到上面的问题了。如果没有消费者生产者调用emit函数永远不会挂起有消费者注册了并且缓存容量已满并且最旧的数据没有被消费则生产者emit函数有机会被挂起如果设定了挂起模式则会被挂起最旧的数据下面会分析只有消费者当只有消费者时消费者调用collect会被挂起吗从collect函数源码入手。override suspend fun collect(collector: FlowCollectorT) { //分配slot val slot allocateSlot()//① try { if (collector is SubscribedFlowCollector) collector.onSubscription() val collectorJob currentCoroutineContext()[Job] while (true) { //死循环 var newValue: Any? while (true) { //尝试获取值 ② newValue tryTakeValue(slot) // attempt no-suspend fast path first if (newValue ! NO_VALUE) break//拿到值退出内层循环 //没拿到值挂起等待 ③ awaitValue(slot) // await signal that the new value is available } collectorJob?.ensureActive() //拿到值消费数据 collector.emit(newValue as T) } } finally { freeSlot(slot) } }重点看三点① allocateSlot()先看Slot数据结构private class SharedFlowSlot : AbstractSharedFlowSlotSharedFlowImpl*() { //消费者当前应该消费的数据在生产者缓存里的索引 var index -1L // current to-be-emitted index, -1 means the slot is free now //挂起的消费者协程体 var cont: ContinuationUnit? null // collector waiting for new value }每此调用collect都会为其生成一个AbstractSharedFlowSlot对象该对象存储在AbstractSharedFlowSlot对象数组slots里allocateSlot() 有两个作用给slots数组扩容往slots数组里存放AbstractSharedFlowSlot对象② tryTakeValue(slot)创建了slot之后就可以去取值了private fun tryTakeValue(slot: SharedFlowSlot): Any? { var resumes: ArrayContinuationUnit? EMPTY_RESUMES val value kotlinx.coroutines.internal.synchronized(this) { //找到slot对应的buffer里的数据索引 val index tryPeekLocked(slot) if (index 0) { //没找到 NO_VALUE } else { //找到 val oldIndex slot.index //根据索引从buffer里获取值 val newValue getPeekedValueLockedAt(index) //slot索引增加指向buffer里的下个数据 slot.index index 1 // points to the next index after peeked one //更新游标等信息并返回挂起的生产者协程 resumes updateCollectorIndexLocked(oldIndex) newValue } } //如果可以则唤起生产者协程 for (resume in resumes) resume?.resume(Unit) return value }该函数有可能取到值也可能取不到。③ awaitValueprivate suspend fun awaitValue(slot: kotlinx.coroutines.flow.SharedFlowSlot): Unit suspendCancellableCoroutine { cont - kotlinx.coroutines.internal.synchronized(this) lock{ //再次尝试获取 val index tryPeekLocked(slot) // recheck under this lock if (index 0) { //说明没数据可取此时记录当前协程后续恢复时才能找到 slot.cont cont // Ok -- suspending } else { //有数据了则唤醒 cont.resume(Unit) // has value, no need to suspend returnlock } slot.cont cont // suspend, waiting } }对比生产者emit和消费者collect流程显然collect流程比emit流程简单多了。现在可以回到上面的问题了。无论是否有生产者只要没拿到数据collect都会被挂起slot与buffer以上分别分析了emit和collect流程我们知道了emit可能被挂起被挂起后可以通过collect唤醒同样的collect也可能被挂起挂起后通过emit唤醒。重点在于两者是如何交换数据的也就是slot对象和buffer是怎么关联的如上图简介其流程SharedFlow设定重放个数为4额外容量为3总容量为437生产者将数据堆到buffer里此时消费者还没开始collect消费者开始collect因为设置了重放个数因此构造Slot对象时slot.index0根据index找到buffer下标为0的元素即为可以消费的元素拿到0号数据后slot.index1找到buffer下标为1的元素index重复4的步骤因为collect消费了数据因此emit可以继续放新的数据此时又有新的collect加入进来新加入的消费者collect时构造Slot对象因为此时的buffer最旧的值为buffer下标为2因此Slot初始化Slot.index 2取第2个数据同样的继续往后取值此时有了2个消费者假设消费者2消费速度很慢它停留在了index3而消费者1消费速度快变成了如下图消费者1在取index4的值可以继续往后消费数据消费者2在取index3的值生产者此时已经填充满buffer了buffer里最旧的值为index4为了保证消费者2能够获取到index4的值此时它不能再emit新的数据了于是生产者被挂起等到消费者2消费了index4的值就会唤醒正在挂起的生产者继续生产数据由此得出一个结论SharedFlow的emit可能会被最慢的collect拖累从而挂起该现象用代码查看打印比较直观fun test7() { runBlocking { //构造热流 val flow MutableSharedFlowString(4, 3) //开启协程 GlobalScope.launch { //接收数据(消费者1) flow.collect { println(collect1: $it) } } GlobalScope.launch { //接收数据(消费者2) flow.collect { //模拟消费慢 delay(10000) println(collect2: $it) } } //发送数据(生产者) delay(200)//保证消费者先执行 var count 0 while (true) { flow.emit(emit:${count}) } } }4. StateFlow 使用方式与应用场景使用方式1. 重放功能上面花了很大篇幅分析SharedFlow而StateFlow是SharedFlow的特例先来看其简单使用。fun test8() { runBlocking { //构造热流 val flow MutableStateFlow() flow.emit(hello world) flow.collect { //消费者 println(it) } } }我们发现并没有给Flow设置重放此时消费者依然能够消费到数据说明StateFlow默认支持历史数据重放。2. 重放个数具体能重放几个值呢fun test10() { runBlocking { //构造热流 val flow MutableStateFlow() flow.emit(hello world) flow.emit(hello world1) flow.emit(hello world2) flow.emit(hello world3) flow.emit(hello world4) flow.collect { //消费者 println(it) } } }最后发现消费者只有1次打印说明StateFlow只重放1次并且是最新的值。3. 防抖fun test9() { runBlocking { //构造热流 val flow MutableStateFlow() flow.emit(hello world) GlobalScope.launch { flow.collect { //消费者 println(it) } } //再发送 delay(1000) flow.emit(hello world) // flow.emit(hello world) } }生产者发送了两次数据猜猜此时消费者有几次打印答案是只有1次因为StateFlow设计了防抖当emit时会检测当前的值和上一次的值是否一致若一致则直接抛弃当前数据不做任何处理collect当然就收不到值了。若是我们将注释放开则会有2次打印。应用场景StateFlow 和LiveData很像都是只维护一个值旧的值过来就会将新值覆盖。适用于通知状态变化的场景如下载进度。适用于只关注最新的值的变化。如果你熟悉LiveData就可以理解为StateFlow基本可以做到替换LiveData功能。5. StateFlow 原理一看就会如果你看懂了SharedFlow原理那么对StateFlow原理的理解就不在话下了。emit 过程override suspend fun emit(value: T) { //value 为StateFlow维护的值每次emit都会修改它 this.value value } public override var value: T get() NULL.unbox(_state.value)//从state取出 set(value) { updateState(null, value ?: NULL) } private fun updateState(expectedState: Any?, newState: Any): Boolean { var curSequence 0 var curSlots: ArrayStateFlowSlot?? this.slots // benign race, we will not use it kotlinx.coroutines.internal.synchronized(this) { val oldState _state.value if (expectedState ! null oldState ! expectedState) return false // CAS support //新旧值一致则无需更新 if (oldState newState) return true // Dont do anything if value is not changing, but CAS - true //更新到state里 _state.value newState curSequence sequence //... curSlots slots // read current reference to collectors under lock } while (true) { curSlots?.forEach { //遍历消费者修改状态或是将挂起的消费者唤醒 it?.makePending() } ... } }emit过程就是修改value值的过程无论是否修改成功emit函数都会退出它不会被挂起。collect 过程override suspend fun collect(collector: FlowCollectorT) { //分配slot val slot allocateSlot() try { if (collector is SubscribedFlowCollector) collector.onSubscription() val collectorJob currentCoroutineContext()[Job] var oldState: Any? null // previously emitted T!! | NULL (null -- nothing emitted yet) while (true) { val newState _state.value collectorJob?.ensureActive() //值不相同才调用collect闭包 if (oldState null || oldState ! newState) { collector.emit(NULL.unbox(newState)) oldState newState } if (!slot.takePending()) { // try fast-path without suspending first //挂起协程 slot.awaitPending() // only suspend for new values when needed } } } finally { freeSlot(slot) } }StateFlow 也有slot叫做StateFlowSlot它比SharedFlowSlot简单多了因为始终只需要维护一个值所以不需要index。里面有个成员变量_state该值既可以是消费者协程当前的状态也可以表示协程体。当表示为协程体时说明此时消费者被挂起了等到生产者通过emit唤醒该协程。6. StateFlow/SharedFlow/LiveData 区别与应用StateFlow 是SharedFlow特例SharedFlow 多用于事件通知StateFlow/LiveData多用于状态变化StateFlow 有默认值LiveData没有StateFlow.collect闭包可在子线程执行LiveData.observe需要在主线程监听StateFlow没有关联生命周期LiveData关联了生命周期StateFlow防抖LiveData不防抖等等。随着本篇的完结Kotlin协程系列也告一段落了接下来将重点放在协程工程架构实践上敬请期待。以上为Flow背压和线程切换的全部内容下篇将分析Flow的热流。本文基于Kotlin 1.5.3文中完整Demo请点击您若喜欢请点赞、关注、收藏您的鼓励是我前进的动力持续更新中和我一起步步为营系统、深入学习Android/Kotlin1、Android各种Context的前世今生2、Android DecorView 必知必会3、Window/WindowManager 不可不知之事4、View Measure/Layout/Draw 真明白了5、Android事件分发全套服务6、Android invalidate/postInvalidate/requestLayout 彻底厘清7、Android Window 如何确定大小/onMeasure()多次执行原因8、Android事件驱动Handler-Message-Looper解析9、Android 键盘一招搞定10、Android 各种坐标彻底明了11、Android Activity/Window/View 的background12、Android Activity创建到View的显示过13、Android IPC 系列14、Android 存储系列15、Java 并发系列不再疑惑16、Java 线程池系列17、Android Jetpack 前置基础系列18、Android Jetpack 易学易懂系列19、Kotlin 轻松入门系列20、Kotlin 协程系列全面解读作者小鱼人爱编程链接https://juejin.cn/post/7195569817940164668最后如果想要成为架构师或想突破20~30K薪资范畴那就不要局限在编码业务要会选型、扩展提升编程思维。此外良好的职业规划也很重要学习的习惯很重要但是最重要的还是要能持之以恒任何不能坚持落实的计划都是空谈。如果你没有方向这里给大家分享一套由阿里高级架构师编写的《Android八大模块进阶笔记》帮大家将杂乱、零散、碎片化的知识进行体系化的整理让大家系统而高效地掌握Android开发的各个知识点。相对于我们平时看的碎片化内容这份笔记的知识点更系统化更容易理解和记忆是严格按照知识体系编排的。全套视频资料一、面试合集二、源码解析合集三、开源框架合集欢迎大家一键三连支持若需要文中资料直接点击文末CSDN官方认证微信卡片免费领取↓↓↓
返回列表