运行 kotlin 流一次,但下游收到两次

mar*_*rs8 4 android kotlin kotlin-flow kotlin-sharedflow

设想

EventHandler.sharedFlow单击按钮时会发出热流。

Repository该流由在 中执行某些操作接收OnEach{}。

EventCollectorA然后,两个事件收集器和接收存储库流EventCollectorB。

然后,事件收集器流被组合并收集在 中MyViewModel。

问题

这两个事件收集器会导致onEach{...}每次单击时运行两次。但是我只想运行onEach{...}一次并在两个事件收集器中接收它。我怎样才能实现这个目标?

注意:我使用 Hilt 只拥有一个Repository,EventCollectorA实例EventCollectorB

流程图

代码

@AndroidEntryPoint
class MainActivity : AppCompatActivity() {
    override fun onCreate(savedInstanceState: Bundle?) {
        super.onCreate(savedInstanceState)
        val binding = ActivityMainBinding.inflate(layoutInflater)
        setContentView(binding.root)

        val viewModel = ViewModelProvider(this).get(MyViewModel::class.java)
        binding.buttonB.setOnClickListener {
            viewModel.userClickEvent("Click Event")
        }
    }
}
Run Code Online (Sandbox Code Playgroud)
@HiltViewModel
class MyViewModel @Inject constructor(
    private val eventHandler: EventHandler,
    private val eventCollectorA: EventCollectorA,
    private val eventCollectorB: EventCollectorB,
) : ViewModel() {
    fun userClickEvent(event: String) = viewModelScope.launch {
        eventHandler.userClick(event)
    }

    init {
        viewModelScope.launch {
            combine(
                eventCollectorA.sharedFlow,
                eventCollectorB.sharedFlow
            ) { a, b ->
                {/*do something*/}
            }.collect()
        }
    }
}
Run Code Online (Sandbox Code Playgroud)
class EventHandler  {
    private val _sharedFlow = MutableSharedFlow<String>()
    val sharedFlow = _sharedFlow.asSharedFlow()

    suspend fun userClick(event: String) {
        _sharedFlow.emit(event)
    }
}
Run Code Online (Sandbox Code Playgroud)
class Repository constructor(
    eventHandler: EventHandler,
) {
    val sharedFlow = eventHandler.sharedFlow
            .filter { it == "Click Event" }
            .onEach {/*do something*/} /*onEach is called twice on click event. I only want it called once*/ 
            .onStart { emit("Begin") }
}
Run Code Online (Sandbox Code Playgroud)
class EventCollectorA constructor(repository: Repository) {
    val sharedFlow = repository.sharedFlow.map {
        it
    }
}

class EventCollectorB constructor(repository: Repository) {
    val sharedFlow = repository.sharedFlow.map {
        it
    }
}
Run Code Online (Sandbox Code Playgroud)

bro*_*oot 6

这里的问题是 whileeventHandler.sharedFlow是 a SharedFlow,在对其应用任何运算符后,我们得到一个常规的、非共享的流。filter(),onEach()并且onStart()针对每个新集合单独运行。如果您想在集合之间共享它们,则需要在应用它们之后构建另一个共享流:

    val sharedFlow = eventHandler.sharedFlow
            .filter { it == "Click Event" }
            .onEach {/*do something*/}
            .onStart { emit("Begin") }
            .shareIn(...)   
Run Code Online (Sandbox Code Playgroud)

进一步解释

我们需要意识到,常规的冷流与实时数据流不同。它更像是此类流的来源,并且对于每个新集合,我们都会启动全新的数据流。例如,如果我们使用flow { }构建器创建一个流,那么我们只有一个流对象,但如果我们collect {}多次调用它,那么对于每个集合,lambda 将被一次又一次地调用。同样,我们用来构造新流的每个运算符也会为每个集合单独调用。

您可以认为shareIn()创建一个服务来观察其上游流并将其数据复制到每个下游流。无论我们收集多少次共享流量,上游流量都只会被收集一次。上面的运算符shareIn()将被调用一次,而下面的运算符shareIn()将为每个集合单独调用。