乐于分享
好东西不私藏

【极速学习spark源码】DataFrame、Dataset、RDD 的转换关系

【极速学习spark源码】DataFrame、Dataset、RDD 的转换关系

基于 Spark 4.2,分支 branch-4.2

阅读时长约 10 分钟 · 入门到中级、

背景

Spark 里有三个数据抽象,互相可以转来转去:

val df = spark.read.json("...")       // DataFrameval ds = df.as[MyClass]               // Dataset[MyClass]val rdd = df.rdd                      // RDD[Row]val df2 = rdd.toDF()                  // 回到 DataFrame

但每次转换到底发生了什么?是零成本还是有代价?

这篇文章把三者的转换关系从源码层面拆开。


一、三者是什么

抽象类型核心表示特点
RDD[T]abstract class RDD分区 + 计算函数类型安全、无 schema、手动优化
DataFrametype DataFrame = Dataset[Row]LogicalPlan + Encoder无类型、有 schema、Catalyst 优化
Dataset[T]class Dataset[T]LogicalPlan + Encoder[T]类型安全、有 schema、Catalyst 优化

DataFrame 不是独立的类。package.scala:33

type DataFrame = Dataset[Row]

所以 DataFrame 和 Dataset 本质是同一个东西,区别只是 Row vs 具体类型 T


二、DataFrame → Dataset:as / toDF

toDF:Dataset[T] → DataFrame

API 层声明在 Dataset.scala:143

def toDF(): DataFrame

Classic 实现在 classic/Dataset.scala:518

def toDF(): DataFrame = withTypedPlan {  Dataset.ofRows(sparkSession, logicalPlan, queryExecution.tracker, ReuseTracker)}

关键:toDF 只是换了 Encoder。LogicalPlan 不变。它把 ExpressionEncoder[T] 换成 RowEncoder,让数据以 Row 形式暴露。

as:DataFrame → Dataset[T]

API 层声明在 Dataset.scala:164

def as[UEncoder]: Dataset[U]

Classic 实现在 classic/Dataset.scala:521

def as[U : Encoder]: Dataset[U] = withTypedPlan {  new Dataset(sparkSession, logicalPlan, encoderFor[U], ReuseTracker)}

同样:LogicalPlan 不变,只是换了 Encoder。

所以 df.as[MyClass].toDF() 几乎零成本——只是 Encoder 的替换。


三、Dataset → RDD:rdd

这是转换链中有实际代价的一步。

API 层声明在 Dataset.scala:3381

def rdd: RDD[T]

Classic 实现在 classic/Dataset.scala:1672-1675

def rdd: RDD[T] = withActive {  val rddObj = materializedRdd  SQLExecution.withSQLExecutionId(rddObj, sparkSession)}

materializedRdd 定义在 classic/Dataset.scala:1664-1668

private[sql] lazy val materializedRddRDD[T] = {  rddQueryExecution.toRdd.mapPartitions { rows =>    rows.map(encoder.deserialize)  }}

rddQueryExecution 定义在 classic/Dataset.scala:1659-1662

private[sql] lazy val rddQueryExecution: QueryExecution = {  val deserialized = CatalystSerde.deserialize[T](logicalPlan)  sparkSession.sessionState.executePlan(deserialized)}

所以 Dataset → RDD 实际做了:

  1. 在 LogicalPlan 上追加 CatalystSerde.deserialize(把 InternalRow 反序列化回 T

  2. 重新走一遍 QueryExecution:Analyze → Optimize → Plan → Prepare

  3. 调用 toRdd 得到 RDD[InternalRow]

  4. 再 mapPartitions 把 InternalRow 反序列化成 T

这个转换有双重代价

  • Catalyst 编译本身的耗时(如果还没编译过)

  • 反序列化:每个 InternalRow 要变成 Row 或自定义对象


四、QueryExecution.toRdd:从 LogicalPlan 到 RDD[InternalRow]

QueryExecution.scala:378-379

val lazyToRdd = LazyTry {  new SQLExecutionRDD(executedPlan.execute(), sparkSession.sessionState.conf)}

QueryExecution.scala:392-394

def toRdd: RDD[InternalRow] = withAbortTransactionOnFailure {  lazyToRdd.get}

SQLExecutionRDD.scala:33-34

class SQLExecutionRDD(    prev: RDD[InternalRow],    conf: SQLConf)

toRdd 产出的是 RDD[InternalRow]——Spark 内部的二进制行格式(Tungsten),不是用户可见的 Row

所以 Dataset → RDD 的完整链路:

Dataset[T]
  → CatalystSerde.deserialize(logicalPlan)    ← 追加反序列化算子
    → QueryExecution
      → executedPlan.execute()
        → RDD[InternalRow]
          → mapPartitions(encoder.deserialize)  ← InternalRow → T
            → RDD[T]

五、RDD → DataFrame / Dataset

spark.createDataFrame(rdd)

Classic 实现在 classic/SparkSession.scala:316-319

def createDataFrame[A <: Product : TypeTag](rdd: RDD[A]): DataFrame = {  val encoder = Encoders.product[A]  val logicalPlan = new ExternalRDD[A](rdd, self)(encoder)  Dataset.ofRows(self, logicalPlan)}

做了什么:

  1. 从 TypeTag 推导 ExpressionEncoder

  2. 把 RDD 包装成 ExternalRDD(一个 LogicalPlan 节点)

  3. Dataset.ofRows 创建 DataFrame

ExternalRDD 本质上是一个"桥接":它把 RDD 作为数据源接入 Catalyst 体系。后续优化和执行都走 Catalyst。

spark.createDataset(rdd)

Classic 实现在 classic/SparkSession.scala:397-399

def createDataset[T : Encoder](data: RDD[T]): Dataset[T] = {  val encoder = implicitly[Encoder[T]]  val logicalPlan = new ExternalRDD[T](data, self)(encoder)  new Dataset(self, logicalPlan, encoder)}

和 createDataFrame 几乎一样,只是 Encoder 不同。

rdd.toDF(隐式转换)

SQLImplicits.scala:61

implicit def rddToDatasetHolder[T : Encoder](rdd: RDD[T]): DatasetHolder[T]

Classic 实现在 classic/SQLImplicits.scala:34-35

implicit def rddToDatasetHolder[T : Encoder](rdd: RDD[T]): DatasetHolder[T] = {  DatasetHolder(session.createDataset(rdd))}

DatasetHolderDatasetHolder.scala:34)提供 toDF() 方法,最终调用底层 Dataset.toDF

所以 rdd.toDF 完整链路:

rdd.toDF
  → SQLImplicits.rddToDatasetHolder
    → spark.createDataset(rdd)
      → new ExternalRDD(rdd)
        → Dataset.ofRows / new Dataset

六、Encoder:三种抽象之间的"翻译器"

贯穿所有转换的核心是 Encoder

  • RowEncoderRowEncoder.scala:64):StructType ↔ InternalRow

  • ExpressionEncoderExpressionEncoder.scala:143):T ↔ InternalRow

Dataset.ofRows 用 RowEncoder

new Dataset[Row](qe, () => RowEncoder.encoderFor(qe.analyzed.schema)) // Dataset.scala:114

Dataset.rdd 用 ExpressionEncoder[T] 反序列化:

rows.map(encoder.deserialize) // Dataset.scala:1667

InternalRow 是 Catalyst 的"通用货币"。所有转换都经过 T → InternalRow → T' 这条路:

Dataset[T] ──encoder.serialize──→ InternalRow ──encoder.deserialize──→ Dataset[T']                ↑                                    ↓           QueryExecution                       mapPartitions           (Catalyst 优化)                    (反序列化)

七、转换代价图

DataFrame ←──→ Dataset[T]   as/toDF         as/toDF (零成本)         (零成本)DataFrame/Dataset ──────→ RDD       rdd   (有代价:追加反序列化算子    + Catalyst 编译 + InternalRow→T)RDD ──────→ DataFrame/Dataset  createDataFrame  createDataset / toDF  (有代价:ExternalRDD 桥接    + Catalyst 编译 + T→InternalRow)

记住

  • DataFrame ↔ Dataset:几乎零成本(只换 Encoder)

  • Dataset → RDD:有代价(Catalyst 编译 + 反序列化)

  • RDD → Dataset:有代价(Catalyst 编译 + 序列化)


八、实战:观察转换代价

    val df = spark.range(100).toDF("id")    // Dataset → RDD    val startTime = System.nanoTime()    val rdd = df.rdd    val rddTime = System.nanoTime() - startTime    println(s"df.rdd took: ${rddTime / 1e6} ms")    // RDD → DataFrame    val startTime2 = System.nanoTime()    val df2 = spark.createDataFrame(rdd, df.schema)    val dfTime = System.nanoTime() - startTime2    println(s"rdd.toDF took: ${dfTime / 1e6} ms")    // DataFrame → Dataset    val startTime3 = System.nanoTime()    val ds = df.as[Long]    val dsTime = System.nanoTime() - startTime3    println(s"df.as took: ${dsTime / 1e6} ms")

结果:

df.rdd took: 1254.79893 ms
rdd.toDF took: 9.690113 ms
df.as took: 28.541962 ms

九、总结

三个抽象之间的转换关系可以压缩成一个图:

     Encoder 替换(零成本)DataFrame ◄────────────────► Dataset[T]    │                            │    │ rdd (Catalyst + 反序列化)    │ rdd (Catalyst + 反序列化)    ▼                            ▼   RDD[Row]                  RDD[T]    ▲                            ▲    │ createDataFrame            │ createDataset / toDF    │ (ExternalRDD + 序列化)       │ (ExternalRDD + 序列化)    └────────────────────────────┘

最值得带走的认知

  1. DataFrame 和 Dataset 是同一个东西。DataFrame 是 Dataset[Row]as / toDF 只换 Encoder。

  2. Dataset → RDD 有实际代价。它在 LogicalPlan 上追加反序列化算子,Catalyst 重新编译,最后 InternalRow 反序列化成用户类型。

  3. InternalRow 是 Catalyst 的通用货币。无论怎么转,底层都是 T → InternalRow → T'


每天花费10分钟学习spark,让你技术之路走得更稳、更快。

喜欢的点个关注。