能否扩展Kotlin SharedFlow实现由消费者指定事件回放长度?
实现方案
完全可以基于官方MutableSharedFlow封装实现需求,同时解决你提到的三个问题,具体实现如下:
import kotlinx.coroutines.flow.MutableSharedFlow import kotlinx.coroutines.flow.SharedFlow import kotlinx.coroutines.flow.asFlow import kotlinx.coroutines.flow.collect import kotlinx.coroutines.flow.emitAll import kotlinx.coroutines.flow.flow import kotlinx.coroutines.flow.replay import kotlin.contracts.ExperimentalContracts import kotlin.contracts.InvocationKind import kotlin.contracts.contract class ReplayConfigurableSharedFlow<T>( // 全局允许的最大回放数量上限,初始化时指定 val maxReplay: Int ) { private val _source = MutableSharedFlow<T>(replay = maxReplay) // 可按需对外暴露只读SharedFlow val readOnlyFlow: SharedFlow<T> = _source suspend fun emit(event: T) = _source.emit(event) @OptIn(ExperimentalContracts::class) suspend inline fun collectWithReplay( count: Int, crossinline action: suspend (T) -> Unit ) { contract { callsInPlace(action, InvocationKind.AT_LEAST_ONCE) } require(count in 0..maxReplay) { "回放数量$count 超出允许范围[0, $maxReplay]" } flow { // 先发射指定数量的历史回放事件 emitAll(_source.replayCache.takeLast(count).asFlow()) // 订阅后续新事件,丢弃底层默认的maxReplay个回放避免重复消费 emitAll(_source.replay(count = 0)) }.collect { action(it) } } }
方案优势
- 完全复用官方
SharedFlow内置的线程安全回放缓存机制,无需自行维护缓存变量,避免并发读写问题 - 不存在事件丢失:
flow构建块内部的缓存读取与后续订阅是连续执行的,replay(0)会从订阅瞬间开始接收所有新事件,不会漏掉中间发出的事件 - 保留inline性能优化:通过
inline+crossinline修饰符实现和官方collect方法完全一致的性能表现,无额外函数调用开销
使用示例
// 初始化实例,设置最大允许回放2个事件 val eventBus = ReplayConfigurableSharedFlow<String>(maxReplay = 2) // 消费者1:不需要回放,仅接收订阅后的新事件 eventBus.collectWithReplay(count = 0) { event -> println("消费者1收到:$event") } // 消费者2:回放最近1个历史事件+接收后续新事件 eventBus.collectWithReplay(count = 1) { event -> println("消费者2收到:$event") } // 消费者3:回放全部2个历史事件+接收后续新事件 eventBus.collectWithReplay(count = 2) { event -> println("消费者3收到:$event") }
如果你不想封装独立类,也可以直接给SharedFlow加扩展函数,前提是你初始化SharedFlow的时候已经设置了对应的replay上限:
@OptIn(ExperimentalContracts::class) suspend inline fun <T> SharedFlow<T>.collectWithReplay( count: Int, crossinline action: suspend (T) -> Unit ) { contract { callsInPlace(action, InvocationKind.AT_LEAST_ONCE) } require(count in 0..replayCache.size) { "回放数量$count 超出当前缓存可用数量[0, ${replayCache.size}]" } flow { emitAll(replayCache.takeLast(count).asFlow()) emitAll(replay(count = 0)) }.collect { action(it) } }
内容的提问来源于stack exchange,提问作者Nilzor
相关产品推荐
相关产品推荐

