【问题标题】:In Apache Spark, how to make an RDD/DataFrame operation lazy?在 Apache Spark 中,如何使 RDD/DataFrame 操作变得惰性?
【发布时间】:2016-05-28 00:32:02
【问题描述】:

假设我想写一个函数 foo 来转换一个 DataFrame:

object Foo {
def foo(source: DataFrame): DataFrame = {
...complex iterative algorithm with a stopping condition...
}
}

由于 foo 的实现有很多“Action”(collect、reduce 等),调用 foo 会立即触发代价高昂的执行。

这不是一个大问题,但是由于 foo 只将一个 DataFrame 转换为另一个,按照惯例,最好允许延迟执行:只有当结果 DataFrame 或其派生项是正在驱动程序上使用(通过另一个“动作”)。

到目前为止,唯一可靠地实现这一点的方法是将所有实现写入 SparkPlan,并将其叠加到 DataFrame 的 SparkExecution 中,这非常容易出错并且涉及大量样板代码。推荐的方法是什么?

【问题讨论】:

    标签: scala apache-spark apache-spark-sql rdd lazy-evaluation


    【解决方案1】:

    我并不完全清楚你试图实现什么,但 Scala 本身至少提供了一些你可能会觉得有用的工具:

    • 惰性值:

      val rdd = sc.range(0, 10000)
      
      lazy val count = rdd.count  // Nothing is executed here
      // count: Long = <lazy>
      
      count  // count is evaluated only when it is actually used 
      // Long = 10000   
      
    • 按名称调用(在函数定义中由=&gt; 表示):

      def  foo(first: => Long, second: => Long, takeFirst: Boolean): Long =
        if (takeFirst) first else second
      
      val rdd1 = sc.range(0, 10000)
      val rdd2 = sc.range(0, 10000)
      
      foo(
        { println("first"); rdd1.count },
        { println("second"); rdd2.count },
        true  // Only first will be evaluated
      )
      // first
      // Long = 10000
      

      注意:在实践中,您应该创建本地惰性绑定,以确保不会在每次访问时评估参数。

    • 无限的惰性集合,例如Stream

      import org.apache.spark.mllib.random.RandomRDDs._
      
      val initial = normalRDD(sc, 1000000L, 10)
      
      // Infinite stream of RDDs and actions and nothing blows :)
      val stream: Stream[RDD[Double]] = Stream(initial).append(
        stream.map {
          case rdd if !rdd.isEmpty => 
            val mu = rdd.mean
            rdd.filter(_ > mu)
          case _ => sc.emptyRDD[Double]
        }
      )
      

    其中一些子集应该足以实现复杂的惰性计算。

    【讨论】:

    • 您也可以使用() =&gt; Futurescalaz.Task
    • 非常感谢您的回答,但在我的情况下,foo(...) 的输出应该像 df.select(...) 产生的任何数据帧一样使用,将其设置为lazy val 或 call-by-name 无法自动启用此功能。
    • 说实话,我看不出有什么理由不这样做。到目前为止,您在问题和评论中描述的内容完全描述了惰性值。
    • 找了半天相信你是对的,这可能是最简单的解决方案了。
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 2016-04-03
    • 2020-09-26
    • 1970-01-01
    • 2019-11-15
    • 1970-01-01
    • 1970-01-01
    • 2014-05-13
    相关资源
    最近更新 更多