【问题标题】:How to use callbackFlow within a flow?如何在流中使用回调流?
【发布时间】:2023-03-26 22:35:01
【问题描述】:

我正在尝试将 callbackFlow 包装在外部 flow 中 - 我想从外部流中发出一些项目,但我有一个旧的回调接口,我想适应 Kotlin 流程。我查看了几个examples of usage of callbackFlow,但我不知道如何在另一个流程中正确触发它。

这是一个例子:

class Processor {
    fun start(processProgress: ProcessProgressListener) {
        processProgress.onFinished() //finishes as soon as it starts!
    }
}

interface ProcessProgressListener {
    fun onFinished()
}

//main method here:
 fun startProcess(processor: Processor): Flow<String> {
     val mainFlow = flow {
         emit("STARTED")
         emit("IN_PROGRESS")
     }

     return merge(processProgressFlow(processor), mainFlow)
 }

 fun processProgressFlow(processor: Processor) = callbackFlow {
     val listener = object : ProcessProgressListener {
         override fun onFinished() {
             trySend("FINISHED")
         }
     }

     processor.start(listener)
 }

Processor 带有一个侦听器,该侦听器在进程完成时触发。发生这种情况时,我想发出最后一个项目FINISHED

我调用整个流程的方式如下:

     runBlocking {
         startProcess(Processor()).collect {
             print(it)
         }
     }

但是,我没有得到任何输出。但是,如果我不使用megre 并且只返回mainFlow,我会得到STARTEDIN_PROGRESS 项目。

我做错了什么?

【问题讨论】:

    标签: kotlin kotlin-coroutines kotlin-flow


    【解决方案1】:

    您忘记在callbackFlow 块末尾调用awaitClose

    fun processProgressFlow(processor: Processor) = callbackFlow<String> {
        val listener = object : ProcessProgressListener {
            override fun onFinished() {
                trySend("FINISHED")
                channel.close()
            }
        }
    
        processor.start(listener)
    
        /*
         * Suspends until 'channel.close() or cancel()' is invoked
         * or flow collector is cancelled (e.g. by 'take(1)' or because a collector's coroutine was cancelled).
         * In both cases, callback will be properly unregistered.
         */
        awaitClose { /* unregister listener here */ }
    }
    

    awaitClose {} 应该用在callbackFlow 块的末尾。 否则,在外部取消的情况下,回调/侦听器可能会泄漏。

    根据callbackFlow docs

    awaitClose 应该用于保持流程运行,否则在阻塞完成时通道将立即关闭。 awaitClose 参数在流消费者取消流收集或基于回调的 API 手动调用 SendChannel.close 时调用,通常用于在完成后清理资源,例如取消注册回调。使用awaitClose 强制是为了防止流收集被取消时内存泄漏,否则即使流收集器已经完成,回调也可能继续运行。为避免此类泄漏,如果块返回,但通道尚未关闭,此方法将抛出 IllegalStateException

    【讨论】:

    • 我明白了,所以我实际上遇到了内存泄漏,这就是为什么没有打印出来的原因?顺便说一句,这完成了工作,谢谢!还请注意 - 这当然意味着我必须更改 Processor 类的设计 - 它必须能够设置/取消设置侦听器,而不是将侦听器作为 start 的参数。
    • IllegalStateException 被抛出,所以什么也没打印出来。
    猜你喜欢
    • 1970-01-01
    • 2017-01-10
    • 2020-09-30
    • 2019-11-16
    • 1970-01-01
    • 2018-11-30
    • 1970-01-01
    • 1970-01-01
    • 2019-06-09
    相关资源
    最近更新 更多