【问题标题】:Why would one use DataFrame.select over DataFrame.rdd.map (or vice versa)?为什么要使用 DataFrame.select 而不是 DataFrame.rdd.map(反之亦然)?
【发布时间】:2017-04-09 19:52:27
【问题描述】:

DataFrame 上使用select 来获取我们需要的信息和为同一目的映射底层RDD 的每一行之间有什么“机械”区别吗?

我所说的“机械”是指执行操作的机制。换句话说,实现细节。

提供的两个中哪一个更好/性能更高?

df = # create dataframe ...
df.select("col1", "col2", ...)

df = # create dataframe ...
df.rdd.map(lambda row: (row[0], row[1], ...))

我正在进行性能测试,因此我将找出哪个更快,但我想知道实现差异和优缺点。

【问题讨论】:

  • 总之,dataframe 的操作总是比 rdd 的操作快。
  • @mtoto 我不会这么肯定。数据集必须对数据进行序列化和反序列化,以便在类型化操作中处理 + Scala 代码,而 UDF 可能会使优化消失(这对于 RDD 来说不会像“危险”那样)。
  • @JacekLaskowski Datasets 中的坏代码可能比 RDDs 中的好代码慢 ;) 但是更多时候 Dataset 将比普通 RDD 快。请注意 Spark、ML 和 Streaming 中的新 API 将以数据集为中心

标签: performance apache-spark dataframe apache-spark-sql rdd


【解决方案1】:

在这个过于简单的例子中,DataFrame.selectDataFrame.rdd.map 我认为差异几乎可以忽略不计。

毕竟您已经加载了数据集并且只进行了投影。最终,两者都必须对 Spark 的 InternalRow 列格式的数据进行反序列化,以计算操作的结果。

您可以通过explain(extended = true) 查看DataFrame.select 发生的情况,您将了解物理计划(以及物理计划)。

scala> spark.version
res4: String = 2.1.0-SNAPSHOT

scala> spark.range(5).select('id).explain(extended = true)
== Parsed Logical Plan ==
'Project [unresolvedalias('id, None)]
+- Range (0, 5, step=1, splits=Some(4))

== Analyzed Logical Plan ==
id: bigint
Project [id#17L]
+- Range (0, 5, step=1, splits=Some(4))

== Optimized Logical Plan ==
Range (0, 5, step=1, splits=Some(4))

== Physical Plan ==
*Range (0, 5, step=1, splits=Some(4))

将物理计划(即SparkPlan)与您正在使用的rdd.maptoDebugString)进行比较,您就会知道什么可能“更好”。

scala> spark.range(5).rdd.toDebugString
res5: String =
(4) MapPartitionsRDD[8] at rdd at <console>:24 []
 |  MapPartitionsRDD[7] at rdd at <console>:24 []
 |  MapPartitionsRDD[6] at rdd at <console>:24 []
 |  MapPartitionsRDD[5] at rdd at <console>:24 []
 |  ParallelCollectionRDD[4] at rdd at <console>:24 []

(在这个人为的例子中,我认为没有赢家——两者都尽可能高效)。

请注意DataFrame实际上是Dataset[Row],它使用RowEncoder将数据编码(即序列化)为InternalRow列二进制格式。如果您要在管道中执行更多运算符,那么坚持使用Dataset 可以获得比RDD 更好的性能,这仅仅是因为低级别的幕后逻辑查询计划优化和列式二进制格式。

有很多优化,试图击败它们通常会浪费您的时间。您必须熟记 Spark 内部结构才能获得更好的性能(而且价格肯定是可读性)。

其中有很多内容,我强烈建议您观看 Herman van Hovell 的演讲 A Deep Dive into the Catalyst Optimizer,以了解和欣赏所有优化。

我的看法是......“除非你知道自己在做什么,否则远离 RDD”

【讨论】:

    【解决方案2】:

    RDD 只是转换和动作的图谱。

    DataFrame 有一个逻辑计划,该计划在执行操作之前由 Catalyst 逻辑查询优化器进行内部优化。

    在你的情况下是什么意思?

    如果你有 DataFrame,那么你应该使用select - 任何额外的工作,如过滤、加入等,都将得到优化。优化后的 DataFrame 可以比普通 RDD 快 10 倍。也就是说,在执行select之前,Spark 会尝试让查询更快。使用dataFrame.rdd.map()时不会这样做

    还有一个:rdd 的值是懒惰计算的:

    lazy val rdd: RDD[T] = {
        val objectType = exprEnc.deserializer.dataType
        val deserialized = CatalystSerde.deserialize[T](logicalPlan)
        sparkSession.sessionState.executePlan(deserialized).toRdd.mapPartitions { rows =>
          rows.map(_.get(0, objectType).asInstanceOf[T])
        }
      }
    

    所以 Spark 将使用它的 RDD、地图和投射内容。两个版本的 DAG 在查询中几乎相同,因此性能相似。然而,在更高级的情况下,使用 Datasets 的好处将非常明显,正如 Spark PMC 在 Databricks 博客上所写的那样,在 Catalyst 优化后,Datasets 甚至可以快 100 倍

    请注意,DataFrame=Dataset[Row] 并且它在后台使用 RDD - 但 RDD 的图形是在优化后创建的

    注意:Spark 是统一的 API。 Spark ML 现在以 DataFrame 为中心,不应使用旧的 API。流媒体正在转向结构化流媒体。因此,即使您的情况不会有太大的性能提升,也请考虑使用 DataFrames。这对未来的发展会是更好的决定,当然会比使用普通的 RDD 更快

    【讨论】:

      猜你喜欢
      • 2015-03-01
      • 2013-10-25
      • 1970-01-01
      • 2018-06-09
      • 2010-11-14
      • 2011-12-07
      • 2014-02-24
      • 1970-01-01
      • 2010-10-23
      相关资源
      最近更新 更多