SparkSQL 之序列化问题深度剖析
·
摘要:本文从 Spark 三层序列化体系(Task Closure / Shuffle Wire / Dataset Internal)、四种序列化路径差异、Encoder vs Kryo vs Java Serializer 性能对比、NotSerializableException 五种解决方案、以及 Kryo 调优清单五个维度,彻底解答 Spark 序列化的一切疑问。
关键词:Spark 序列化, Kryo, Encoder, NotSerializableException, Tungsten, InternalRow, Shuffle
一、开篇
几乎每个 Spark 开发者都遇到过 Task not serializable。本质原因:Driver 端创建的闭包需要序列化后发送到 Executor,闭包引用的外部对象不可序列化。
// ❌ 经典错误
val conn = DriverManager.getConnection(url)
rdd.map(row => conn.execute(s"INSERT ... $row")) // NotSerializableException!
// ✅ 正确: mapPartitions 内创建
rdd.mapPartitions { iter =>
val conn = DriverManager.getConnection(url)
iter.map(row => conn.execute(...))
}
二、三层序列化体系

| Layer | 路径 | 序列化器 |
|---|---|---|
| ① Task Closure | Driver→Executor | spark.serializer (Kryo) |
| ② Shuffle Wire | Executor↔Executor | RDD:Kryo / Dataset:Encoder |
| ③ Dataset Internal | 算子执行 | Encoder(Tungsten)旁路Kryo |
三、Encoder vs Kryo vs Java 对比

Java(默认) Kryo(推荐) Encoder(Tungsten)
───────────────────────────────────────────────────────
体积 1x (基准) 0.1x 0.01x (100x!)
速度 1x (基准) 10x 100x
适用范围 所有对象 RDD+闭包 Dataset API
配置 无需 需注册类 import implicits
Shuffle 是 是 是(绕过Kryo)
四、NotSerializableException 五种解法
① @transient lazy val — 延迟初始化不可序列化字段
② 广播变量 — 大对象只序列化一次
③ mapPartitions — 每个分区内创建一次
④ 实现 Serializable — 让类可序列化
⑤ 局部变量捕获 — 只捕获需要的值而非整个对象
五、Kryo 调优清单
val conf = new SparkConf()
.set("spark.serializer", "org.apache.spark.serializer.KryoSerializer")
conf.registerKryoClasses(Array(classOf[MyKey], classOf[MyData]))
conf.set("spark.kryo.registrationRequired", "true")
conf.set("spark.kryoserializer.buffer.max", "128m")
六、总结
-
三层序列化:Task Closure(Kryo) + Shuffle(RDD用Kryo/Dataset用Encoder) + Dataset Internal(Encoder旁路)。
-
性能排名:Encoder(100x) > Kryo(10x) > Java(1x)。Dataset Shuffle 自动走 Encoder。
-
避坑五法:@transient lazy / 广播变量 / mapPartitions / Serializable / 局部变量捕获。
作者:starzy
博客:blog.starzy.cn
GitHub:starzy1990.github.io
专注 AI Agent · LangGraph · RAG · 大数据架构 · 数据工程实践
更多推荐



所有评论(0)