RDD序列化
Spark 序列化:什么时候需要序列化?Java vs Kryo 怎么选?NotSerializableException 怎么破?
序列化在 Spark 里发生在三个地方:
- Driver → Executor 传数据:广播变量、闭包(函数)要序列化后才能发给 Executor
- Shuffle:Map 输出的数据要序列化后通过网络传给 Reduce
- 缓存:RDD 存到内存/磁盘时,如果用了序列化存储级别,需要序列化
说白了:任何跨节点或跨存储的数据传递,都需要序列化。
Java 序列化 vs Kryo 序列化
Spark 默认用 Java 序列化。它不用配,啥都能序列化,但慢、体积大。
Kryo 序列化是备选方案:快(10倍)、体积小(1/3~1/5),但要手动配置。
| 对比项 | Java 序列化 | Kryo 序列化 |
|---|---|---|
| 配置 | 不用配 | 要配 |
| 速度 | 慢 | 快 10 倍 |
| 体积 | 大 | 小 3-5 倍 |
| 兼容性 | 高(任何 Serializable) | 中(需注册类) |
什么时候用 Kryo? 任何时候。Spark 官方文档也建议生产环境用 Kryo。
Kryo 配置三步走
第一步:在 SparkConf 里启用 Kryo
val conf = new SparkConf()
.setAppName("MyApp")
.set("spark.serializer", "org.apache.spark.serializer.KryoSerializer")
第二步:注册自定义类
conf.registerKryoClasses(Array(
classOf[Person],
classOf[Order],
classOf[Array[String]]
))
注册的目的是让 Kryo 提前知道这些类,不用运行时动态发现,省时间。
第三步(推荐):强制注册
conf.set("spark.kryo.registrationRequired", "true")
开启后,如果遇到没注册的类,直接报错,而不是偷偷用 Java 序列化兜底。这样你能发现所有需要注册的类,保证性能。
Kryo 调参
| 参数 | 默认值 | 什么时候调 |
|---|---|---|
spark.kryoserializer.buffer |
64KB | 对象太大,报 “buffer overflow” 时调大 |
spark.kryoserializer.buffer.max |
64MB | 单个超大对象(比如 100MB 的 HashMap)时调大 |
spark.kryo.registrationRequired |
false | 生产环境建议开 true |
spark.kryo.referenceTracking |
true | 没有循环引用可以关掉,省点开销 |
conf.set("spark.kryoserializer.buffer.max", "512m")
经典问题:NotSerializableException
这是 Spark 新手最常见的报错。闭包里的对象没实现序列化。
错误示例:
class DatabaseConnection { // 没实现 Serializable
def query(id: Int): String = ???
}
val conn = new DatabaseConnection()
// 报错:闭包捕获了 conn,但 conn 不可序列化
val rdd = sc.parallelize(1 to 10).map(id => conn.query(id))
原因: map 里的函数会在 Executor 上执行,函数里用到的 conn 必须从 Driver 序列化后传过去。conn 不可序列化,就报错了。
解决方案:
方案1:用 mapPartitions,在 Executor 端创建
val rdd = sc.parallelize(1 to 10).mapPartitions { iter =>
// 在 Executor 端创建,不需要从 Driver 传过去
val conn = new DatabaseConnection()
iter.map(id => conn.query(id))
}
方案2:标记 @transient,在 Executor 端重新初始化
class DatabaseConnection { def query(id: Int): String = ??? }
class MyClass extends Serializable {
lazy val conn = new DatabaseConnection() // transient:不序列化
def process(id: Int) = conn.query(id)
}
val obj = new MyClass()
rdd.map(id => obj.process(id)) // MyClass 序列化传过去,conn 在 Executor 端重建
方案3:让对象实现 Serializable(如果可以的话)
class DatabaseConnection extends Serializable { ... }
什么时候需要显式关注序列化?
看到 NotSerializableException 就要排查闭包里有没有不可序列化的对象。
常见不可序列化的东西:
- 数据库连接
- 文件流
- 非
Serializable的自定义类 - 闭包里的外部变量
排查方法:
- 看报错栈顶,找出哪个类没序列化
- 看看这个类在闭包里哪里被引用了
- 用上面三种方案之一解决