DataSet编程

DataSet:给 DataFrame 加上类型安全,编译时就能发现错误

DataFrame 好用,但有一个问题:弱类型。

// DataFrame:写错列名,编译不报错,运行才报错
df.filter("ageee > 18")   // 列名写错了,编译能过,运行挂掉

DataSet 就是解决这个问题的:强类型 DataFrame。

DataFrame DataSet
类型 弱类型(Row) 强类型(Person)
列名检查 运行时 编译时
访问字段 row.getInt(0) 或 $"age" person.age
Python 支持 支持 不支持
适用场景 快速分析、Python 开发 需要类型安全的复杂业务逻辑

一句话:DataSet 是给 Scala/Java 开发者用的”类型安全版 DataFrame”。

DataSet 长什么样?

定义样例类,然后用 toDS() 创建:

// 1. 定义样例类(字段名即列名)
case class Person(name: String, age: Int, city: String)

// 2. 创建 DataSet
import spark.implicits._
val ds = Seq(
  Person("Alice", 25, "北京"),
  Person("Bob", 30, "上海")
).toDS()

// 3. 操作:用对象字段访问
ds.filter(_.age > 18).show()
ds.select($"name", $"age").show()

关键区别:

// DataFrame:用字符串列名
df.filter("age > 18")

// DataSet:用对象字段(编译时检查)
ds.filter(_.age > 18)   // 写错成 _.ageee 编译报错

DataSet 的创建方式

2.1 从集合创建

case class User(name: String, age: Long)
val ds = Seq(User("Alice", 25), User("Bob", 30)).toDS()

2.2 从 RDD 转换

val rdd = sc.parallelize(Seq(User("Alice", 25), User("Bob", 30)))
val ds = rdd.toDS()

2.3 从 DataFrame 转换(加类型)

val df = spark.read.json("people.json")
val ds = df.as[Person]   // DataFrame → DataSet[Person]

前提: DataFrame 的列名和类型必须跟 Person 的字段名和类型完全匹配。

DataSet 支持的操作

3.1 DataFrame 风格操作(用 $ 语法)

ds.select($"name", $"age").show()
ds.filter($"age" > 18).show()
ds.groupBy("city").count().show()

3.2 类型安全操作(用 . 访问字段)

// map、filter 里用对象字段
ds.filter(_.age > 18).map(_.name).show()

// 复杂转换
ds.map(person => (person.name, person.age * 2)).show()

3.3 Join

case class Order(id: Long, userName: String, amount: Double)
val orders = Seq(Order(1, "Alice", 100.0)).toDS()

ds.join(orders, ds("name") === orders("userName")).show()

3.4 写数据

ds.write.parquet("output.parquet")
ds.write.json("output.json")

DataSet ↔ DataFrame ↔ RDD 互转

// DataSet → DataFrame
val df = ds.toDF()

// DataFrame → DataSet
val ds = df.as[Person]

// DataSet → RDD
val rdd = ds.rdd   // RDD[Person]

// RDD → DataSet
val ds = rdd.toDS()

注意事项:

  • DataFrame 转 DataSet 时,列名和类型必须跟样例类匹配
  • toDS() 和 as[T] 需要 import spark.implicits._

DataSet 的局限性

1. Python 不支持

Python 是动态类型语言,没有 DataSet。Python 开发者只能用 DataFrame。

2. 编码器限制

DataSet 依赖编码器(Encoder)把对象转成二进制。Spark 对样例类支持很好,但复杂嵌套类型可能需要自定义编码器。

3. 性能略低于 DataFrame

类型安全有代价,DataSet 多了类型检查和编码转换。但差距很小,复杂查询中可忽略。