基于 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、手动优化 |
| DataFrame | type 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(): DataFrameClassic 实现在 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[U: Encoder]: 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 = materializedRddSQLExecution.withSQLExecutionId(rddObj, sparkSession)}
materializedRdd 定义在 classic/Dataset.scala:1664-1668:
private[sql] lazy val materializedRdd: RDD[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 实际做了:
在 LogicalPlan 上追加
CatalystSerde.deserialize(把InternalRow反序列化回T)重新走一遍
QueryExecution:Analyze → Optimize → Plan → Prepare调用
toRdd得到RDD[InternalRow]再
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)}
做了什么:
从
TypeTag推导ExpressionEncoder把 RDD 包装成
ExternalRDD(一个 LogicalPlan 节点)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))}
DatasetHolder(DatasetHolder.scala:34)提供 toDF() 方法,最终调用底层 Dataset.toDF。
所以 rdd.toDF 完整链路:
rdd.toDF
→ SQLImplicits.rddToDatasetHolder
→ spark.createDataset(rdd)
→ new ExternalRDD(rdd)
→ Dataset.ofRows / new Dataset
六、Encoder:三种抽象之间的"翻译器"
贯穿所有转换的核心是 Encoder:
RowEncoder(RowEncoder.scala:64):StructType ↔ InternalRowExpressionEncoder(ExpressionEncoder.scala:143):T ↔ InternalRow
Dataset.ofRows 用 RowEncoder:
new Dataset[Row](qe, () => RowEncoder.encoderFor(qe.analyzed.schema)) // Dataset.scala:114Dataset.rdd 用 ExpressionEncoder[T] 反序列化:
rows.map(encoder.deserialize) // Dataset.scala:1667InternalRow 是 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 ──────→ RDDrdd(有代价:追加反序列化算子+ Catalyst 编译 + InternalRow→T)RDD ──────→ DataFrame/DatasetcreateDataFramecreateDataset / toDF(有代价:ExternalRDD 桥接+ Catalyst 编译 + T→InternalRow)
记住:
DataFrame ↔ Dataset:几乎零成本(只换 Encoder)
Dataset → RDD:有代价(Catalyst 编译 + 反序列化)
RDD → Dataset:有代价(Catalyst 编译 + 序列化)
八、实战:观察转换代价
val df = spark.range(100).toDF("id")// Dataset → RDDval startTime = System.nanoTime()val rdd = df.rddval rddTime = System.nanoTime() - startTimeprintln(s"df.rdd took: ${rddTime / 1e6} ms")// RDD → DataFrameval startTime2 = System.nanoTime()val df2 = spark.createDataFrame(rdd, df.schema)val dfTime = System.nanoTime() - startTime2println(s"rdd.toDF took: ${dfTime / 1e6} ms")// DataFrame → Datasetval startTime3 = System.nanoTime()val ds = df.as[Long]val dsTime = System.nanoTime() - startTime3println(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 + 序列化)└────────────────────────────┘
最值得带走的认知:
DataFrame 和 Dataset 是同一个东西。DataFrame 是
Dataset[Row],as/toDF只换 Encoder。Dataset → RDD 有实际代价。它在 LogicalPlan 上追加反序列化算子,Catalyst 重新编译,最后
InternalRow反序列化成用户类型。InternalRow 是 Catalyst 的通用货币。无论怎么转,底层都是
T → InternalRow → T'。
每天花费10分钟学习spark,让你技术之路走得更稳、更快。
喜欢的点个关注。
夜雨聆风