Spark SQL 中 DataFrame 的创建方式及与 RDD 的区别

创建 DataFrame 的多种方式

1. 从结构化数据文件创建

JSON 文件
// 从 JSON 文件创建 DataFrame
val df = spark.read.json("path/to/file.json")
df.show()

// 指定 schema
import org.apache.spark.sql.types._
val customSchema = StructType(Array(
  StructField("name", StringType, true),
  StructField("age", IntegerType, true)
))
val dfWithSchema = spark.read.schema(customSchema).json("path/to/file.json")
CSV 文件
// 基本 CSV 读取
val csvDF = spark.read
  .option("header", "true")
  .option("inferSchema", "true")
  .csv("path/to/file.csv")

// 自定义选项
val advancedCsvDF = spark.read
  .option("header", "true")
  .option("delimiter", ";")
  .option("quote", "'")
  .option("escape", "\\")
  .csv("path/to/file.csv")
Parquet 文件
// Parquet 是推荐的列式存储格式
val parquetDF = spark.read.parquet("path/to/data.parquet")

// 写入 Parquet
df.write.mode("overwrite").parquet("output/path")

2. 从现有 RDD 创建

使用反射推断 Schema
// 定义 case class
case class Person(name: String, age: Int, city: String)

// 从 RDD 创建 DataFrame
val rdd = spark.sparkContext.textFile("people.txt")
val personRDD = rdd.map(_.split(",")).map(p => Person(p(0), p(1).toInt, p(2)))
val personDF = spark.createDataFrame(personRDD)
personDF.createOrReplaceTempView("people")
编程方式指定 Schema
import org.apache.spark.sql.types._

// 创建 RDD
val rdd = spark.spark(1, "Alice", 25), (2, "Bob", 30))

// 定义 schema
val schema = StructType(Array(
  StructField("id", IntegerType, false),
  StructField("name", StringType, false),
  StructField("age", IntegerType, false)
))

// 应用 schema 到 RDD
val rowRDD = rdd.map(attributes => Row(attributes._1, attributes._2, attributes._3))
val df = spark.createDataFrame(rowRDD, schema)

3. 从外部数据源创建

JDBC 数据库
val jdbcDF = spark.read
  .format("jdbc")
  .option("url", "jdbc:mysql://localhost:3306/mydb")
  .option("dbtable", "users")
  .option("user", "username")
  .option("password", "password")
  .load()
Hive 表
// 需要启用 Hive 支持
val hiveDF = spark.sql("SELECT * FROM hive_table")
// 或者
val hiveDF = spark.table("hive_table")

4. 通过编程方式创建

手动创建数据
import spark.implicits._

// 从序列创建
val data = Seq(
  ("Alice", 25),
  ("Bob", 30),
  ("Charlie", 35)
)
val df = data.toDF("name", "age")

// 从元组创建
val tupleDF = Seq((1, "a"), (2, "b")).toDF("id", "value")

DataFrame 与 RDD 的主要区别

1. 抽象层次对比

RDD - 弹性分布式数据集
低级抽象
记录级别的操作
DataFrame - 分布式数据集合
高级抽象
表结构的操作

2. 详细比较表格

特性 RDD DataFrame
抽象级别 低级API,记录级别操作 高级API,表结构操作
类型安全 编译时检查 运行时检查
优化 无自动优化 Catalyst优化器优化
内存使用 Java对象序列化开销大 Tungsten引擎高效编码
易用性 需要更多代码 类似SQL的表达式
性能 相对较低 更高(可达2倍以上)

3. 代码示例对比

RDD 方式
// RDD 实现年龄筛选和计数
val rdd = spark.sparkContext.parallelize(Seq(("Alice", 25), ("Bob", 30), ("Charlie", 35)))
val filteredRDD = rdd.filter(_._2 > 25)
val count = filteredRDD.count()
println(s"Age > 25的人数: $count")
DataFrame 方式
// DataFrame 实现相同功能
import spark.implicits._
val df = Seq(("Alice", 25), ("Bob", 30), ("Charlie", 35)).toDF("name", "age")
val count = df.filter($"age" > 25).count()
println(s"Age > 25的人数: $count")

// 或使用 SQL
df.createOrReplaceTempView("people")
val result = spark.sql("SELECT COUNT(*) FROM people WHERE age > 25")
result.show()

4. 性能差异原因

Catalyst 优化器工作流程
Logical Plan
Analyzer
Logical Optimization
Physical Plan Generation
Cost-based Optimization
Optimized Physical Plan
Tungsten 内存管理优势
  • 二进制内存布局减少垃圾回收压力
  • 缓存感知算法提高内存访问效率
  • 堆外内存使用避免 JVM GC 开销

5. 适用场景选择

使用 RDD 的情况:
  • 需要精确控制数据分布和分区
  • 复杂的函数式变换需求
  • 非结构化数据处理
  • 需要在不同类型数据间进行复杂转换
使用 DataFrame/Dataset 的情况:
  • 结构化或半结构化数据分析
  • 需要高性能查询执行
  • 希望利用自动优化功能
  • 团队成员熟悉 SQL 或关系代数

6. 互操作性

两者之间可以相互转换:

// DataFrame 转 RDD
val rddFromDF = df.rdd

// RDD 转 DataFrame(如前所示)
val dfFromRDD = spark.createDataFrame(rdd, schema)

总的来说,DataFrame 提供了比 RDD 更高层次的抽象,在大多数情况下能够提供更好的性能和更简洁的 API,而 RDD 则提供了更多的控制灵活性。在实际应用中,应该根据具体需求选择合适的抽象层级。

Logo

码道开发者社区,聚焦华为云码道 CodeArts 代码智能体,沉淀 Agent、Skill、鸿蒙开发实战内容,供开发者查阅资料、交流技术、分享工程实践

更多推荐