【问题标题】:How to pass result as it comes using coroutines?如何使用协程传递结果?
【发布时间】:2020-04-22 08:39:01
【问题描述】:

假设我有 repos 列表。我想遍历所有这些。当每个 repo 都返回结果时,我想把它传递下去。

val repos = listOf(repo1, repo2, repo3)
val deferredItems = mutableListOf<Deferred<List<result>>>()

repos.forEach { repo ->
    deferredItems.add(async { getResult(repo) })
}

val results = mutableListOf<Any>()
deferredItems.forEach { deferredItem ->
    results.add(deferredItem.await())
}

println("results :: $results")

在上述情况下,它等待每个 repo 返回结果。它按顺序填充resultsrepo1 的结果后跟repo2 的结果。如果repo1repo2 花费更多的时间来返回结果,我们将等待repo1 的结果,即使我们有repo2 的结果。

有什么方法可以在我们得到结果后立即传递repo2 的结果?

【问题讨论】:

    标签: asynchronous kotlin coroutine kotlin-coroutines


    【解决方案1】:

    你应该使用Channels

    suspend fun loadReposConcurrent() = coroutineScope {
        val repos = listOf(repo1, repo2, repo3)
        val channel = Channel<List<YourResultType>>()
    
        for (repo in repos) {
            launch {
                val result = getResult(repo)
                channel.send(result)
            }
        }
    
        var allResults = emptyList<YourResultType>()
        repeat(repos.size) {
            val result = channel.receive()
            allResults = allResults + result
    
            println("results :: $result")
            //updateUi(allResults)
        }
    }
    

    在上面for (repo in repos) {...} 的代码中,使用launch 循环在单独的协程中计算的所有请求,一旦它们的结果准备好,就会发送到channel

    repeat(repos.size) {...} 中,channel.receive() 等待来自所有协程的新值并使用它们。

    【讨论】:

      【解决方案2】:

      这就是频道的用途:

      val repos = listOf("repo1", "repo2", "repo3")
      val results = Channel<Result>()
      repos.forEach { repo ->
          launch {
              val res = getResult(repo)
              results.send(res)
          }
      }
      
      for (r in results) {
          println(r)
      }
      

      这个例子是不完整的,因为我没有关闭通道,所以生成的代码将永远挂起。确保在收到所有结果后在您的真实代码中关闭通道:

      val count = AtomicInteger()
      
      for (r in results) {
          println(r)
          if (count.incrementAndGet() == repos.size) {
              results.close()
          }
      }
      

      【讨论】:

        【解决方案3】:

        Flow API 几乎直接支持这一点:

        repos.asFlow()
                .flatMapMerge { flow { emit(getResult(it)) } }
                .collect { println(it) }
        

        flatMapMerge 首先收集来自您传递给它的 lambda 的所有 Flows,然后 同时 收集这些并在其中任何一个完成后立即将它们发送到下游。

        【讨论】:

          猜你喜欢
          • 2020-07-16
          • 2019-12-23
          • 1970-01-01
          • 1970-01-01
          • 1970-01-01
          • 2014-09-14
          • 1970-01-01
          • 1970-01-01
          • 1970-01-01
          相关资源
          最近更新 更多