求满足初始值、重放1条、无重复无合并的Kotlin Flow多平台实现
问题描述
我需要实现一个值流,能将生产者的值传递给一个或多个消费者,具体需求如下:
- 流具备默认初始值;
- 消费者订阅时接收最近发布的值(重放最近值);
- 可通过
value属性同步获取最近发布的值; - 连续发布的值不会重复;
- 所有发布的值必须被消费者接收(值不可被合并);
- 消费基于
FlowCollector进行; - 支持多平台(Android、iOS、JS)。
现有方案分析
MutableStateFlow
MutableStateFlow几乎满足所有需求,但不满足第5条:它的值更新始终会被合并,慢收集器会跳过快速更新,仅能收集到最新值。
示例代码:
val flow = MutableStateFlow(0) // 在协程1中执行 flow.collect { println("Collect: $it") delay(100) } // 在协程2中执行 repeat(10) { flow.emit(it + 1) }
该代码只会打印0和10,而非中间的所有数值,因为慢收集器导致中间值被合并。
MutableSharedFlow
MutableSharedFlow不会合并值,但存在三个缺陷:
- 没有同步的
value属性(仅提供挂起函数last()); - 允许发布重复值;
- 无默认初始值。
虽然可以手动补充这些特性,但实现value属性和重复值检查的逻辑并不简单。
Reaktive的BehaviorSubject
通过BehaviorSubject配合distinctUntilChanged()可以完美满足所有需求:
val subject = BehaviorSubject(0) val flow = subject.asFlow().distinctUntilChanged()
但该方案需要引入第三方库Reaktive。
问题
是否存在无需手动补充特性的Kotlin Flow原生方案,能满足上述所有需求?
回答
目前Kotlin协程原生的Flow组件(MutableStateFlow/MutableSharedFlow)没有现成的实现能直接满足所有需求。不过你可以基于MutableSharedFlow封装一个轻量的自定义流,完美匹配你的需求:
import kotlinx.coroutines.flow.MutableSharedFlow import kotlinx.coroutines.flow.distinctUntilChanged import kotlinx.coroutines.sync.Mutex import kotlinx.coroutines.sync.withLock class MutableBehaviorFlow<T>(initialValue: T) { private val mutex = Mutex() private var currentValue = initialValue private val sharedFlow = MutableSharedFlow<T>(replay = 1, extraBufferCapacity = Int.MAX_VALUE) init { sharedFlow.tryEmit(initialValue) } val value: T get() = currentValue val flow = sharedFlow.distinctUntilChanged() suspend fun emit(value: T) { mutex.withLock { if (value != currentValue) { currentValue = value sharedFlow.emit(value) } } } fun tryEmit(value: T): Boolean { return mutex.withLock { if (value != currentValue) { currentValue = value sharedFlow.tryEmit(value) } else { true } } } }
该实现的特性说明
- 支持初始值,新订阅者会收到最近一次发布的值(通过
replay=1的MutableSharedFlow实现); - 提供同步
value属性,可直接获取当前最新值; - 通过锁保护的重复值检查+
distinctUntilChanged(),确保连续发布的重复值被过滤; - 设置
extraBufferCapacity = Int.MAX_VALUE(可根据业务场景调整为合理值),避免值被合并,保证所有发布的非重复值都能被消费者接收; - 基于原生Flow构建,支持Android、iOS、JS多平台;
- 消费端可直接通过
flow属性使用FlowCollector进行收集。
如果不想自行封装,使用Reaktive库的BehaviorSubject方案是最直接的选择,它原生支持所有你需要的特性。
内容的提问来源于stack exchange,提问作者Emanuel Moecklin
相关产品推荐
相关产品推荐

