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

求满足初始值、重放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不会合并值,但存在三个缺陷:

  1. 没有同步的value属性(仅提供挂起函数last());
  2. 允许发布重复值;
  3. 无默认初始值。
    虽然可以手动补充这些特性,但实现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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.20 08:21:58