【问题标题】:When to use Scala Futures?何时使用 Scala 期货?
【发布时间】:2020-01-20 17:54:42
【问题描述】:

我是 spark Scala 程序员。我有一个火花工作,它有完成整个工作的子任务。我想用期货来并行完成子任务。完成整个工作后,我必须返回整个工作响应。

我听说 scala Futures 是,一旦主线程执行并停止,其余线程将被杀死,并且您将得到空响应。

我必须使用 Await.result 来收集结果。但是所有的博客都告诉你应该避免 Await.result,这是一种不好的做法。

在 mycase 中使用 Await.result 是否正确?

def computeParallel(): Future[String] = {
  val f1 = Future {  "ss" }
  val f2 = Future { "sss" }
  val f3 = Future { "ssss" }

  for {
    r1 <- f1
    r2 <- f2
    r3 <- f3
  } yield (r1 + r2 + r3)
} 

computeParallel().map(result => ???)



根据我的理解,我们必须在 Webservice 类型的应用程序中使用 Futures,它有一个始终在运行且不会退出的进程。但就我而言,一旦逻辑执行(scala 程序)完成,它将退出。

我可以使用 Futures 来解决我的问题吗?

提前致谢

【问题讨论】:

标签: scala apache-spark apache-spark-sql future futuretask


【解决方案1】:

在 Spark 中使用 future 可能是不可取的,除非在特殊情况下,并且简单地并行计算不是其中之一(为阻塞 I/O 提供非阻塞包装器(例如向外部服务发出请求)是很可能的唯一的特殊情况)。

请注意,Future 不保证并行性(它们是否以及如何并行执行取决于它们在其中运行的ExecutionContext),只是异步。此外,如果您在 Spark 转换中生成计算性能期货(即在执行程序上,而不是驱动程序上),则很可能不会有任何性能改进,因为 Spark 往往做得很好让 executor 上的核心保持忙碌,所有产生这些 future 所做的就是与 Spark 竞争核心。

总的来说,在组合 Spark RDDs/DStreams/Dataframes、actors 和 futures 等并行抽象时要非常小心:有很多潜在的雷区,这些组合可能会违反各种组件中的保证和/或约定。

还值得注意的是,Spark 对中间值的可序列化有要求,并且期货通常不可序列化,因此 Spark 阶段不能导致未来;这意味着你基本上别无选择,只能在一个阶段产生的期货上Await。

如果您仍想在 Spark 阶段生成期货(例如,将它们发布到 Web 服务),最好使用 Future.sequence 将期货折叠成一个,然后在上面使用 Await(请注意,我有没有测试过这个想法:我假设有一个隐含的CanBuildFrom[Iterator[Future[String]], String, Future[String]] 可用):

def postString(s: String): Future[Unit] = ???

def postStringRDD(rdd: RDD[String]): RDD[String] = {
  rdd.mapPartitions { strings =>
    // since this is only get used for combining the futures in the Await, it's probably OK to use the implicit global execution context here
    implicit val ectx = ???
    Await.result(strings.map(postString))
  }
  rdd  // Pass through the original RDD
}

【讨论】:

    猜你喜欢
    • 2013-07-20
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2020-11-05
    • 2020-11-10
    • 2019-03-23
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多