在 Spark SQL 中,如何创建 DataFrame?DataFrame 与 RDD 有什么区别?
·
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. 抽象层次对比
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 优化器工作流程
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 则提供了更多的控制灵活性。在实际应用中,应该根据具体需求选择合适的抽象层级。
更多推荐


所有评论(0)