在并行计算中正确使用期货

syn*_*pse 7 scala

我正在实现一个可以轻松并行化的算法,但无法弄清楚如何创建适当数量的期货以及如何提前中止.目前,代码的轮廓沿着这些方向

def solve: Boolean = {
  var result = false
  while(!result && i < iterations) {
    val futures = (1 to threads) map { _ => solveIter(geInitialValues()) }
    val loopResult = Future.fold(futures)(false)((acc, r) => acc || r )
    result = Await.result(loopResult, Duration.Inf)
    i+=1
  }
}

def solveIter(initialValues: Values): Future[Boolean] = Future {
  /* Takes a lot of time */
}
Run Code Online (Sandbox Code Playgroud)

显而易见的问题是明确设置并行级别,该级别可能适合或可能不适合当前执行上下文.如果所有的期货都立即创建,那么如何Future.fold提前中止?

j-k*_*eck 2

您无法取消 Future,因为 Future 是只读的。但是你可以使用 Promise,它是“来自 Future 的写入部分”。

示例代码:

  • 此代码在 5 秒后超时,因为 future 尚未完成(solveIter 永远不会完成)
  • 要完成承诺,请删除将已完成的承诺添加到“承诺”的注释 - 其他期货将被取消
  • 删除 'promises.foreach(_.trySuccess(false))' 并且您再次超时,因为其他 future 没有被取消

import scala.concurrent.ExecutionContext.Implicits.global
import scala.concurrent.duration._
import scala.concurrent.{Await, Future, Promise}
import scala.util.Try


// create a bunch of promises
val promises = ((1 to 10) map { _ =>
  val p = Promise[Boolean]()
  p.completeWith(solveIter())
  p
}) // :+ Promise().success(true)
// ^^ REMOVE THIS COMMENT TO ADD A PROMISE WHICH COMPLETES

// get the futures from the promises
val futures = promises.map(_.future)

// loop over all futures
futures.foreach(oneFuture =>
  // register callback when future is done
  oneFuture.foreach{
    case true =>
      println("future with 'true' result found")

      // stop others
      promises.foreach(_.trySuccess(false))

    case _ => // future completes with false
  })



// wait at most 5 seconds till all futures are done
Try(Await.ready(Future.sequence(futures), 5.seconds)).recover { case _ =>
  println("TIMEOUT")
}


def solveIter(): Future[Boolean] = Future {
  /* Takes a VERY VERY VERY .... lot of time */
  Try(Await.ready(Promise().future, Duration.Inf))
  false
}
Run Code Online (Sandbox Code Playgroud)