【问题标题】:How to wait for several Futures?如何等待几个Future?
【发布时间】:2013-04-21 19:10:31
【问题描述】:

假设我有几个未来,需要等到其中一个其中任何一个失败全部成功。

例如:假设有 3 个期货:f1f2f3

  • 如果f1 成功并且f2 失败,我不会等待f3(并将失败返回给客户端)。

  • 如果f2 失败而f1f3 仍在运行,我不会等待它们(并返回失败

  • 如果f1 成功,然后f2 成功,我继续等待f3

你会如何实现它?

【问题讨论】:

标签: scala concurrency future


【解决方案1】:

你可以用这个:

val l = List(1, 6, 8)

val f = l.map{
  i => future {
    println("future " +i)
    Thread.sleep(i* 1000)
    if (i == 12)
      throw new Exception("6 is not legal.")
    i
  }
}

val f1 = Future.sequence(f)

f1 onSuccess{
  case l => {
    logInfo("onSuccess")
    l.foreach(i => {

      logInfo("h : " + i)

    })
  }
}

f1 onFailure{
  case l => {
    logInfo("onFailure")
  }

【讨论】:

    【解决方案2】:

    您可能想查看 Twitter 的 Future API。尤其是 Future.collect 方法。它完全符合您的要求:https://twitter.github.io/scala_school/finagle.html

    Future.scala 源代码可在此处获得: https://github.com/twitter/util/blob/master/util-core/src/main/scala/com/twitter/util/Future.scala

    【讨论】:

      【解决方案3】:

      这个问题已经得到解答,但我发布了我的价值类解决方案(价值类是在 2.10 中添加的),因为这里没有。欢迎批评指正。

        implicit class Sugar_PimpMyFuture[T](val self: Future[T]) extends AnyVal {
          def concurrently = ConcurrentFuture(self)
        }
        case class ConcurrentFuture[A](future: Future[A]) extends AnyVal {
          def map[B](f: Future[A] => Future[B]) : ConcurrentFuture[B] = ConcurrentFuture(f(future))
          def flatMap[B](f: Future[A] => ConcurrentFuture[B]) : ConcurrentFuture[B] = concurrentFutureFlatMap(this, f) // work around no nested class in value class
        }
        def concurrentFutureFlatMap[A,B](outer: ConcurrentFuture[A], f: Future[A] => ConcurrentFuture[B]) : ConcurrentFuture[B] = {
          val p = Promise[B]()
          val inner = f(outer.future)
          inner.future onFailure { case t => p.tryFailure(t) }
          outer.future onFailure { case t => p.tryFailure(t) }
          inner.future onSuccess { case b => p.trySuccess(b) }
          ConcurrentFuture(p.future)
        }
      

      ConcurrentFuture 是一个无开销的 Future 包装器,它将默认的 Future map/flatMap 从 do-this-then-that 更改为 combine-all-and-fail-if-any-fail。用法:

      def func1 : Future[Int] = Future { println("f1!");throw new RuntimeException; 1 }
      def func2 : Future[String] = Future { Thread.sleep(2000);println("f2!");"f2" }
      def func3 : Future[Double] = Future { Thread.sleep(2000);println("f3!");42.0 }
      
      val f : Future[(Int,String,Double)] = {
        for {
          f1 <- func1.concurrently
          f2 <- func2.concurrently
          f3 <- func3.concurrently
        } yield for {
         v1 <- f1
         v2 <- f2
         v3 <- f3
        } yield (v1,v2,v3)
      }.future
      f.onFailure { case t => println("future failed $t") }
      

      在上面的示例中,f1、f2 和 f3 将同时运行,如果任何顺序失败,则元组的未来将立即失败。

      【讨论】:

      • 太棒了!有没有提供这种实用功能的库?
      • 是的,我已经创建了一个广泛的 Future 实用程序库:github.com/S-Mach/s_mach.concurrent 请参阅示例代码中的 async.par。
      【解决方案4】:

      你可以使用如下的理解:

      val fut1 = Future{...}
      val fut2 = Future{...}
      val fut3 = Future{...}
      
      val aggFut = for{
        f1Result <- fut1
        f2Result <- fut2
        f3Result <- fut3
      } yield (f1Result, f2Result, f3Result)
      

      在此示例中,期货 1、2 和 3 并行启动。然后,在 for comprehension 中,我们等到结果 1,然后是 2,然后是 3。如果 1 或 2 失败,我们将不再等待 3。如果所有 3 个都成功,则 aggFut val 将保存一个具有 3 个插槽的元组,对应于 3 个期货的结果。

      现在,如果您需要在说 fut2 先失败时停止等待的行为,事情会变得有点棘手。在上面的示例中,您必须等待 fut1 完成才能意识到 fut2 失败。为了解决这个问题,你可以尝试这样的事情:

        val fut1 = Future{Thread.sleep(3000);1}
        val fut2 = Promise.failed(new RuntimeException("boo")).future
        val fut3 = Future{Thread.sleep(1000);3}
      
        def processFutures(futures:Map[Int,Future[Int]], values:List[Any], prom:Promise[List[Any]]):Future[List[Any]] = {
          val fut = if (futures.size == 1) futures.head._2
          else Future.firstCompletedOf(futures.values)
      
          fut onComplete{
            case Success(value) if (futures.size == 1)=> 
              prom.success(value :: values)
      
            case Success(value) =>
              processFutures(futures - value, value :: values, prom)
      
            case Failure(ex) => prom.failure(ex)
          }
          prom.future
        }
      
        val aggFut = processFutures(Map(1 -> fut1, 2 -> fut2, 3 -> fut3), List(), Promise[List[Any]]())
        aggFut onComplete{
          case value => println(value)
        }
      

      现在这可以正常工作,但问题在于知道在成功完成后要从Map 中删除哪个Future。只要您有某种方法可以将结果与产生该结果的 Future 正确关联,那么这样的事情就可以工作。它只是递归地不断从 Map 中删除已完成的 Futures,然后在剩余的 Futures 上调用 Future.firstCompletedOf 直到没有剩余,沿途收集结果。这并不漂亮,但如果你真的需要你正在谈论的行为,那么这个或类似的东西可以工作。

      【讨论】:

      • 谢谢。如果 fut2fut1 之前失败会发生什么?在那种情况下,我们还会等待fut1 吗?如果我们愿意,那不是我想要的。
      • 但是如果3先失败,我们还是等1和2能早点回来。有什么方法可以在不需要对期货进行排序的情况下做到这一点?
      • 您可以为fut2 安装onFailure 处理程序以快速失败,并在aggFut 上安装onSuccess 以处理成功。 aggFut 上的成功意味着 fut2 已成功完成,因此您只调用了一个处理程序。
      • 我在答案中添加了一些内容,以显示如果任何期货失败的快速失败的可能解决方案。
      • 在您的第一个示例中,1 2 和 3 不并行运行,而是串行运行。用 printlines 试试看
      【解决方案5】:

      为此,我将使用 Akka 演员。与 for-comprehension 不同,只要任何一个 future 失败,它就会失败,所以从这个意义上说它更有效率。

      class ResultCombiner(futs: Future[_]*) extends Actor {
      
        var origSender: ActorRef = null
        var futsRemaining: Set[Future[_]] = futs.toSet
      
        override def receive = {
          case () =>
            origSender = sender
            for(f <- futs)
              f.onComplete(result => self ! if(result.isSuccess) f else false)
          case false =>
            origSender ! SomethingFailed
          case f: Future[_] =>
            futsRemaining -= f
            if(futsRemaining.isEmpty) origSender ! EverythingSucceeded
        }
      
      }
      
      sealed trait Result
      case object SomethingFailed extends Result
      case object EverythingSucceeded extends Result
      

      然后,创建actor,向它发送一条消息(以便它知道将其回复发送到哪里)并等待回复。

      val actor = actorSystem.actorOf(Props(new ResultCombiner(f1, f2, f3)))
      try {
        val f4: Future[Result] = actor ? ()
        implicit val timeout = new Timeout(30 seconds) // or whatever
        Await.result(f4, timeout.duration).asInstanceOf[Result] match {
          case SomethingFailed => println("Oh noes!")
          case EverythingSucceeded => println("It all worked!")
        }
      } finally {
        // Avoid memory leaks: destroy the actor
        actor ! PoisonPill
      }
      

      【讨论】:

      • 对于这么简单的任务来说看起来有点太复杂了。我真的需要一个演员来等待未来吗?还是谢谢。
      • 我在 API 中找不到任何合适的方法可以完全满足您的要求,但也许我错过了一些东西。
      【解决方案6】:

      这是一个不使用演员的解决方案。

      import scala.util._
      import scala.concurrent._
      import java.util.concurrent.atomic.AtomicInteger
      
      // Nondeterministic.
      // If any failure, return it immediately, else return the final success.
      def allSucceed[T](fs: Future[T]*): Future[T] = {
        val remaining = new AtomicInteger(fs.length)
      
        val p = promise[T]
      
        fs foreach {
          _ onComplete {
            case s @ Success(_) => {
              if (remaining.decrementAndGet() == 0) {
                // Arbitrarily return the final success
                p tryComplete s
              }
            }
            case f @ Failure(_) => {
              p tryComplete f
            }
          }
        }
      
        p.future
      }
      

      【讨论】:

        【解决方案7】:

        您可以使用承诺,并将第一次失败或最终完成的聚合成功发送给它:

        def sequenceOrBailOut[A, M[_] <: TraversableOnce[_]](in: M[Future[A]] with TraversableOnce[Future[A]])(implicit cbf: CanBuildFrom[M[Future[A]], A, M[A]], executor: ExecutionContext): Future[M[A]] = {
          val p = Promise[M[A]]()
        
          // the first Future to fail completes the promise
          in.foreach(_.onFailure{case i => p.tryFailure(i)})
        
          // if the whole sequence succeeds (i.e. no failures)
          // then the promise is completed with the aggregated success
          Future.sequence(in).foreach(p trySuccess _)
        
          p.future
        }
        

        如果你想阻止,你可以AwaitFuture,或者只是map它变成别的东西。

        与 for comprehension 的区别在于,在这里您会得到第一个失败的错误,而使用 for comprehension 您会在输入集合的遍历顺序中得到第一个错误(即使另一个首先失败)。例如:

        val f1 = Future { Thread.sleep(1000) ; 5 / 0 }
        val f2 = Future { 5 }
        val f3 = Future { None.get }
        
        Future.sequence(List(f1,f2,f3)).onFailure{case i => println(i)}
        // this waits one second, then prints "java.lang.ArithmeticException: / by zero"
        // the first to fail in traversal order
        

        还有:

        val f1 = Future { Thread.sleep(1000) ; 5 / 0 }
        val f2 = Future { 5 }
        val f3 = Future { None.get }
        
        sequenceOrBailOut(List(f1,f2,f3)).onFailure{case i => println(i)}
        // this immediately prints "java.util.NoSuchElementException: None.get"
        // the 'actual' first to fail (usually...)
        // and it returns early (it does not wait 1 sec)
        

        【讨论】:

          【解决方案8】:

          您可以单独使用期货来做到这一点。这是一个实现。请注意,它不会提前终止执行!在这种情况下,您需要做一些更复杂的事情(并且可能自己实现中断)。但是,如果您只是不想继续等待无法正常工作的事情,关键是继续等待第一件事完成,并在没有任何内容或遇到异常时停止:

          import scala.annotation.tailrec
          import scala.util.{Try, Success, Failure}
          import scala.concurrent._
          import scala.concurrent.duration.Duration
          import ExecutionContext.Implicits.global
          
          @tailrec def awaitSuccess[A](fs: Seq[Future[A]], done: Seq[A] = Seq()): 
          Either[Throwable, Seq[A]] = {
            val first = Future.firstCompletedOf(fs)
            Await.ready(first, Duration.Inf).value match {
              case None => awaitSuccess(fs, done)  // Shouldn't happen!
              case Some(Failure(e)) => Left(e)
              case Some(Success(_)) =>
                val (complete, running) = fs.partition(_.isCompleted)
                val answers = complete.flatMap(_.value)
                answers.find(_.isFailure) match {
                  case Some(Failure(e)) => Left(e)
                  case _ =>
                    if (running.length > 0) awaitSuccess(running, answers.map(_.get) ++: done)
                    else Right( answers.map(_.get) ++: done )
                }
            }
          }
          

          下面是一个在一切正常时的示例:

          scala> awaitSuccess(Seq(Future{ println("Hi!") }, 
            Future{ Thread.sleep(1000); println("Fancy meeting you here!") },
            Future{ Thread.sleep(2000); println("Bye!") }
          ))
          Hi!
          Fancy meeting you here!
          Bye!
          res1: Either[Throwable,Seq[Unit]] = Right(List((), (), ()))
          

          但是当出现问题时:

          scala> awaitSuccess(Seq(Future{ println("Hi!") }, 
            Future{ Thread.sleep(1000); throw new Exception("boo"); () }, 
            Future{ Thread.sleep(2000); println("Bye!") }
          ))
          Hi!
          res2: Either[Throwable,Seq[Unit]] = Left(java.lang.Exception: boo)
          
          scala> Bye!
          

          【讨论】:

          • 很好的实现。但请注意,如果您将一个空的期货序列传递给 awaitSuccess,它将永远等待......
          猜你喜欢
          • 1970-01-01
          • 2018-09-18
          • 1970-01-01
          • 1970-01-01
          • 2021-02-24
          • 2019-03-02
          • 1970-01-01
          • 1970-01-01
          • 1970-01-01
          相关资源
          最近更新 更多