SparkSQL 之序列化问题深度剖析

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闭包引用的外部对象不可序列化。// ❌ 经典错误valconnDriverManager.getConnection(url)rdd.map(rowconn.execute(sINSERT ...$row))// NotSerializableException!// ✅ 正确: mapPartitions 内创建rdd.mapPartitions{itervalconnDriverManager.getConnection(url)iter.map(rowconn.execute(...))}二、三层序列化体系Layer路径序列化器① Task ClosureDriver→Executorspark.serializer (Kryo)② Shuffle WireExecutor↔ExecutorRDD: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 调优清单valconfnewSparkConf().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.cnGitHubstarzy1990.github.io专注 AI Agent · LangGraph · RAG · 大数据架构 · 数据工程实践