RDD序列化

Spark 序列化:什么时候需要序列化?Java vs Kryo 怎么选?NotSerializableException 怎么破?

序列化在 Spark 里发生在三个地方:

  1. Driver → Executor 传数据:广播变量、闭包(函数)要序列化后才能发给 Executor
  2. Shuffle:Map 输出的数据要序列化后通过网络传给 Reduce
  3. 缓存: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 {
  @transient 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 的自定义类
  • 闭包里的外部变量

排查方法:

  1. 看报错栈顶,找出哪个类没序列化
  2. 看看这个类在闭包里哪里被引用了
  3. 用上面三种方案之一解决