【问题标题】:Play/Scala: Making unknown number of I/O calls in parallell, watining for the resultsPlay/Scala:并行进行未知数量的 I/O 调用,等待结果
【发布时间】:2015-03-25 04:05:15
【问题描述】:

所以,我阅读了关于并行理解的文章here。他给出了以下代码示例:

// Make 3 parallel async calls
val fooFuture = WS.url("http://foo.com").get()
val barFuture = WS.url("http://bar.com").get()
val bazFuture = WS.url("http://baz.com").get()

for {
  foo <- fooFuture
  bar <- barFuture
  baz <- bazFuture
} yield {
  // Build a Result using foo, bar, and baz
  Ok(...)
} 

到目前为止一切都很好,但是,我不知道我需要始终执行多少个 WS.get(),我希望它是动态的。比如:

val checks = Seq(callOne(param), callTwo(param))

电话在哪里:

def callOne(param: String): Future[Boolean] = {
  // do something and return the Future with a true/false value
  Future(true)
}

def callTwo(param: String): Future[Boolean] = {
  // do something and return the Future with a true/false value
  Future(false)
}

所以,我的问题是,我应该如何对我的序列结果与 WS 调用(或与此相关的数据库查询)做出反应?

我给出了两个调用示例,但我希望相同的代码能够并行处理 1 到多个调用,并在 for-yield 中收集结果以最终继续执行其他操作。

重要提示:所有调用都应该并行执行,最快的将在慢的之前完成,而不考虑它们被触发的顺序。

【问题讨论】:

    标签: scala concurrency playframework-2.0


    【解决方案1】:

    Future.sequence 可能是您想要的。

    示例用法:

    val futures = List(WS.url("http://foo.com").get(), WS.url("http://bar.com").get())
    Future.sequence(futures) # => Transforms a Seq[Future[_]] to Future[Seq[_]]
    

    Future.sequence 的未来返回只有在输入序列中的所有期货都完成后才会完成。

    奖励:

    如果您的期货是异构类型的,并且您需要保留该类型,则可以使用 Hlist。我编写了以下 sn-p ,它将获取一个期货 Hlist,并将其转换为一个包含已解析值 Hlist 的 Future:

    import shapeless._
    import scala.concurrent.{ExecutionContext,Future}
    
    object FutureHelpers {
      object FutureReducer extends Poly2 {
        import scala.concurrent.ExecutionContext.Implicits.global
        implicit def f[A, B <: HList] = at[Future[A], Future[B]] { (f, resultFuture) =>
          for {
            result <- resultFuture
            value <- f
          } yield value :: result
        }
      }
    
      // Like Future.sequence, but for HList
      // hsequence(Future { 1 } :: Future { "string" } :: HNil)
      // => Future { 1 :: "string" :: HNil }
      def hsequence[T <: HList](hlist: T)(implicit
        executor: ExecutionContext,
        folder: RightFolder[T, Future[HNil], FutureReducer.type]) = {
        hlist.foldRight(Future.successful[HNil](HNil))(FutureReducer)
      }
    }
    

    【讨论】:

    • 看起来这就是我要找的东西,但是,我如何在等待所有人完成之前动态填充 futures
    • 如果你有一个 urls: Seq[String|, urls map (WS.url).
    • 好的,所以我对 Future.sequence 进行了一些测试,但它们是按顺序执行的,1、2、3 ......并等待所有这些都结束。
    • 您正在寻找的正是 Future.sequence 所做的。
    • 换句话说,您的期货将在您实例化它们后立即安排。因此,如果您实例化一个 Futures 列表,它们将被同时安排(只要您的 ExecutionContext 允许;您可能需要根据您的设备调整执行程序以获得更高的并行度,因为它默认为机器上核心数的 2 倍) . Future.sequence 将等待列表中的每个未来,然后为您提供结果列表(与提供的期货列表的顺序相同)。
    猜你喜欢
    • 2019-10-21
    • 1970-01-01
    • 1970-01-01
    • 2021-06-28
    • 2013-02-21
    • 2023-03-18
    • 1970-01-01
    • 2016-07-13
    • 1970-01-01
    相关资源
    最近更新 更多