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 多了类型检查和编码转换。但差距很小,复杂查询中可忽略。