如何理解复杂的SharedFlow表达式及相关操作?
原代码:
sharedFlow.debounce(250) .onEach(::updateGroupingStrategy) .shareIn(viewModelScope, SharingStarted.WhileSubscribed(5000L), replay = 1)
debounce 解析
官方说明(中文翻译)
返回一个镜像原始流的流,但会过滤掉在给定超时时间内被新值跟随的旧值。最新的值总会被发射出来。
debounce 用于检测一段时间内没有新数据提交的状态,能让你在输入完成后再处理数据。
疑问解答
既不是延迟订阅SharedFlow,也不是每隔250毫秒订阅一次值,核心逻辑是防抖过滤:
- 上游
sharedFlow每发射一个值,就启动250毫秒的计时 - 若计时结束前上游又发新值,直接取消当前计时,重新开始算250毫秒
- 只有当250毫秒内没有新值产生,才会把这段时间里的最后一个值发射给下游
举个场景:上游在100ms、200ms、300ms连续发射A、B、C,debounce会在300ms+250ms=550ms时发射C,前两个值因为后续有新值覆盖,都被过滤掉了。
onEach 解析
官方说明(中文翻译)
返回一个流,该流会在上游流的每个值发射到下游之前,先调用给定的[action]操作。
疑问解答
onEach是Flow的中间操作,它不会触发流的订阅(Flow是冷流,只有调用collect这类终端操作才会启动订阅流程)。它的作用是给流添加一个副作用:当流最终被订阅、上游有值要向下传递时,会先执行updateGroupingStrategy这个函数,再把值传给下游的订阅者。
举个实际执行的例子:
val processedFlow = sharedFlow.debounce(250).onEach(::updateGroupingStrategy)... processedFlow.collect { value -> // 处理下游收到的值 }
当debounce后的流准备发射值时,会先调用updateGroupingStrategy(value)执行逻辑,之后才把value传给collect代码块处理。
shareIn 解析
官方说明(中文翻译)
当第一个订阅者出现时开始共享,默认情况下最后一个订阅者消失后立即停止共享,默认永久保留重放缓存。
它有以下可选参数:
stopTimeoutMillis — 配置最后一个订阅者消失到共享协程停止之间的延迟(毫秒),默认值为0(立即停止)。
replay — 向新订阅者重放的值的数量(不能为负,默认值为0)。
疑问解答
- 关于
WhileSubscribed(5000L):准确逻辑是最后一个订阅者取消订阅后,会等待5000毫秒;如果这段时间内没有新订阅者加入,就停止共享协程;若期间有新订阅者,就继续维持共享状态。 - 关于
replay = 1:当新订阅者订阅这个共享流时,会立即收到缓存中保存的最近1个值。比如之前共享流已经发射过A、B、C三个值,新订阅者一订阅就会先拿到C,之后再接收后续新发射的值。如果replay设为0,新订阅者只能拿到订阅之后上游发射的新值。
内容的提问来源于stack exchange,提问作者chrisChris

