You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

能否扩展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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.10.02 23:45:01