【问题标题】:Kotlin Flow: How to unsubscribe/stopKotlin Flow:如何取消订阅/停止
【发布时间】:2019-11-27 00:43:11
【问题描述】:

更新协程 1.3.0-RC

工作版本:

@FlowPreview
suspend fun streamTest(): Flow<String> = channelFlow {
    listener.onSomeResult { result ->
        if (!isClosedForSend) {
            offer(result)
        }
    }

    awaitClose {
        listener.unsubscribe()
    }
}

还可以查看 Roman Elizarov 的这篇 Medium 文章:Callbacks and Kotlin Flows

原始问题

我有一个发送多个字符串的流:

@FlowPreview
suspend fun streamTest(): Flow<String> = flowViaChannel { channel ->
    listener.onSomeResult { result ->
            if (!channel.isClosedForSend) {
                channel.sendBlocking(result)
            }
    }
}

一段时间后,我想取消订阅该流。目前我做以下事情:

viewModelScope.launch {
    beaconService.streamTest().collect {
        Timber.i("stream value $it")
        if(it == "someString")
            // Here the coroutine gets canceled, but streamTest is still executed
            this.cancel() 
    }
}

如果协程被取消,流仍然会被执行。只是没有订阅者在听新的价值观。如何取消订阅和停止stream 功能?

【问题讨论】:

  • 我觉得这个问题和stackoverflow.com/questions/59680533/…是一样的
  • 感谢您的评论。这不完全一样。我的问题是我的流发射器如何检测流是否不再需要并且它可以取消订阅侦听器。
  • 有一些扩展函数可以让你从你的启动 {} 块中取消范围。如果您使用 Kotlin 1.3.7+,您现在应该可以安全地调用 cancel(),因为您在示例代码中拥有它。在这里查看答案:stackoverflow.com/a/65121663/1738090

标签: kotlin kotlin-coroutines


【解决方案1】:

使用当前版本的协程/Flows (1.2.x) 我现在不是一个好的解决方案。使用onCompletion,您将在流停止时收到通知,但您随后将处于streamTest 功能之外,并且很难停止侦听新事件。

beaconService.streamTest().onCompletion {

}.collect {
    ...
}

使用协程的下一个版本(1.3.x),这将非常容易。函数flowViaChannel 已被弃用,取而代之的是channelFlow。此功能允许您等待流程关闭并在此时执行某些操作,例如。移除监听器:

channelFlow<String> {
    println("Subscribe to listener")

    awaitClose {
        println("Unsubscribe from listener")
    }
}

【讨论】:

  • 我已将 flowViaChannel 更改为 channelFlow,现在 isClosedForSend 始终返回 true。我需要指定频道吗?
【解决方案2】:

解决方案不是取消流程,而是取消流程的启动范围。

val job = scope.launch { flow.cancellable().collect { } }
job.cancel()

注意:如果您希望在取消Job 时停止收集器,则应在collect 之前调用cancellable()

【讨论】:

  • @Farid 正确的方法是什么?
  • @barryalan2633 我建议进行编辑以在流程后添加cancellable() 标志并获得批准。现在答案是正确的。我正在删除我之前的回复,以免造成任何进一步的混乱
【解决方案3】:

当流在 cououtin 范围内运行时,您可以从中获取作业以控制停止订阅。

// Make member variable if you want.
var jobForCancel : Job? = null

// Begin collecting
jobForCancel = viewModelScope.launch {
    beaconService.streamTest().collect {
        Timber.i("stream value $it")
        if(it == "someString")
            // Here the coroutine gets canceled, but streamTest is still executed
            // this.cancel() // Don't
    }
}

// Call whenever to canceled
jobForCancel?.cancel()

【讨论】:

    【解决方案4】:

    您可以在 Flow 上使用 takeWhile 运算符。

    flow.takeWhile { it != "someString" }.collect { emittedValue ->
             //Do stuff until predicate is false  
           }
    

    【讨论】:

    • 这个解决方案帮助了我!文档不是最明显的,但是当 takeWhile{} 中的谓词返回 false 时,集合就完成了。
    • @n 你的意思是即使执行 collect 时条件发生变化?
    【解决方案5】:

    对于那些愿意在 Coroutine 范围内取消订阅 Flow 的人,这种方法对我有用:

    viewModelScope.launch {
    
            beaconService.streamTest().collect {
    
                //Do something then
                this.coroutineContext.job.cancel()
    
            }
    }
    

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2022-01-19
      • 2012-03-14
      • 2021-02-27
      • 2019-12-05
      • 1970-01-01
      相关资源
      最近更新 更多