【问题标题】:Launch a number of coroutines and join them all with timeout (without cancelling)启动多个协程并在超时后将它们全部加入(不取消)
【发布时间】:2019-03-10 07:15:57
【问题描述】:

我需要启动一些将返回结果的作业。

在主代码中(不是协程),启动作业后,我需要等待他们全部完成任务 为给定的超时到期,以先到者为准。

如果我因为所有作业在超时前完成而退出等待,那很好,我会收集它们的结果。

但是,如果某些作业的超时时间更长,我的主要功能需要在超时到期后立即唤醒,检查哪些作业确实及时完成(如果有的话)以及哪些仍在运行,然后开始工作在那里,没有取消仍在运行的作业。

你会如何编写这种等待?

【问题讨论】:

    标签: kotlin kotlin-coroutines


    【解决方案1】:

    解决方案直接来自问题。首先,我们将为任务设计一个挂起函数。让我们看看我们的要求:

    如果某些作业花费的时间超过超时...而不取消仍在运行的作业。

    这意味着我们启动的作业必须是独立的(不是子级),因此我们将选择退出结构化并发并使用GlobalScope 来启动它们,手动收集所有作业。我们使用async coroutine builder 是因为我们计划稍后收集他们的某种类型的结果R

    val jobs: List<Deferred<R>> = List(numberOfJobs) { 
        GlobalScope.async { /* our code that produces R */ }
    }
    

    在启动作业后,我需要等待他们全部完成任务或给定超时到期,以先到者为准。

    让我们等待所有这些,然后等待超时:

    withTimeoutOrNull(timeoutMillis) { jobs.joinAll() }
    

    我们使用joinAll(而不是awaitAll)在其中一项作业失败时避免异常,并使用withTimeoutOrNull 来避免超时异常。

    我的主要功能需要在超时到期后立即唤醒,检查哪些作业已及时完成(如果有)以及哪些仍在运行

    jobs.map { deferred -> /* ... inspect results */ }
    

    在主代码中(不是协程)...

    由于我们的主代码不是协程,它必须以阻塞方式等待,因此我们使用runBlocking 桥接我们编写的代码。把它们放在一起:

    fun awaitResultsWithTimeoutBlocking(
        timeoutMillis: Long,
        numberOfJobs: Int
    ) = runBlocking {
        val jobs: List<Deferred<R>> = List(numberOfJobs) { 
            GlobalScope.async { /* our code that produces R */ }
        }    
        withTimeoutOrNull(timeoutMillis) { jobs.joinAll() }
        jobs.map { deferred -> /* ... inspect results */ }
    }
    

    附:我不建议在任何一种严肃的生产环境中部署这种解决方案,因为让你的后台作业在超时后运行(泄漏)以后总是会严重地咬你。仅当您完全了解这种方法的所有缺陷和风险时才这样做。

    【讨论】:

      【解决方案2】:

      您可以尝试使用whileSelectonTimeout 子句。但是你还是要克服你的主代码不是协程的问题。接下来的几行是whileSelect 语句的示例。该函数返回一个Deferred,其中包含在超时期间评估的结果列表和另一个未完成结果的Deferreds 列表。

      fun CoroutineScope.runWithTimeout(timeoutMs: Int): Deferred<Pair<List<Int>, List<Deferred<Int>>>> = async {
      
          val deferredList = (1..100).mapTo(mutableListOf()) {
              async {
                  val random = Random.nextInt(0, 100)
                  delay(random.toLong())
                  random
              }
          }
      
          val finished = mutableListOf<Int>()
          val endTime = System.currentTimeMillis() + timeoutMs
      
          whileSelect {
              var waitTime = endTime - System.currentTimeMillis()
              onTimeout(waitTime) {
                  false
              }
              deferredList.toList().forEach { deferred ->
                  deferred.onAwait { random ->
                      deferredList.remove(deferred)
                      finished.add(random)
                      true
                  }
              }
          }
      
          finished.toList() to deferredList.toList()
      }
      

      在您的主代码中,您可以使用不鼓励的方法runBlocking 来访问Deferrred

      fun main() = runBlocking<Unit> {
          val deferredResult = runWithTimeout(75)
          val (finished, pending) = deferredResult.await()
          println("Finished: ${finished.size} vs Pending: ${pending.size}")
      }
      

      【讨论】:

      • 谢谢。但是你的代码在极端情况下表现如何?例如,当deferredList 为空时,或者当所有作业在超时之前终止时?
      • 另外,为什么不鼓励runBlocking,你在哪里读到的?
      【解决方案3】:

      这是我想出的解决方案。将每个作业与状态(以及其他信息)配对:

      private enum class State { WAIT, DONE, ... }
      
      private data class MyJob(
          val job: Deferred<...>,
          var state: State = State.WAIT,
          ...
      )
      

      并编写一个显式循环:

      // wait until either all jobs complete, or a timeout is reached
      val waitJob = launch { delay(TIMEOUT_MS) }
      while (waitJob.isActive && myJobs.any { it.state == State.WAIT }) {
          select<Unit> {
              waitJob.onJoin {}
              myJobs.filter { it.state == State.WAIT }.forEach { 
                  it.job.onJoin {}
              }
          }
          // mark any finished jobs as DONE to exclude them from the next loop
          myJobs.filter { !it.job.isActive }.forEach { 
              it.state = State.DONE
          }
      }
      

      初始状态称为 WAIT(而不是 RUN),因为它并不一定意味着作业仍在运行,只是我的循环尚未考虑到它。

      我很想知道这是否足够地道,或者是否有更好的方法来编码这种行为。

      【讨论】:

        猜你喜欢
        • 1970-01-01
        • 1970-01-01
        • 2020-06-07
        • 1970-01-01
        • 2020-04-08
        • 1970-01-01
        • 2012-07-15
        • 2020-01-07
        • 1970-01-01
        相关资源
        最近更新 更多