Spark中RDD与DataFrame序列化行为差异的原因探究
为什么Spark RDD用外部变量会触发序列化问题,而DataFrame不会?
这是个非常典型的Spark执行模型差异问题,核心在于RDD和DataFrame/Dataset的底层执行逻辑完全不同,尤其是对外部变量的处理方式:
1. RDD的闭包序列化问题
RDD是基于JVM对象的分布式集合,当你调用map这类转换操作时,传入的函数是一个闭包——它会捕获外部作用域的变量(比如你代码里的num)。Spark需要把这个闭包和它引用的所有对象一起序列化,发送到各个Executor节点执行。
在你的EX1代码里:
object Example { val r = 1 to 1000000 toList val rdd = sc.parallelize(r,3) val num = 1 val rdd2 = rdd.map(_ + num) rdd2.collect }
这里的num是Example这个单例object的成员变量。Scala的object默认不会自动实现Serializable接口,所以当Spark尝试序列化闭包时,会连带尝试序列化整个Example对象,最终触发序列化失败的异常。
2. DataFrame的Catalyst优化器帮你规避了问题
DataFrame(以及强类型的Dataset)是基于Spark的Catalyst优化器和Tungsten执行引擎工作的,它的处理流程和RDD完全不同:
- 你写的
$"b" + n这类表达式,首先会被转换成逻辑执行计划,而不是直接生成闭包函数。 - Catalyst会在Driver端对这个逻辑计划进行优化,其中就包括常量折叠——它会直接把
n的值(1)嵌入到执行计划里,最终生成的物理执行计划是b + 1,根本不需要把n或者Example对象序列化发送到Executor。
看你的EX2代码:
object Example { import spark.implicits._ import org.apache.spark.sql.functions._ val n = 1 val df = sc.parallelize(Seq( ("r1", 1, 1), ("r2", 6, 4), ("r3", 4, 1), ("r4", 1, 2) )).toDF("ID", "a", "b") df.repartition(3).withColumn("plus1", $"b" + n).show(false) }
这里的n在Driver端就被解析成了常量,Executor拿到的是已经优化好的执行计划,不需要处理任何来自Driver的闭包或外部对象,自然不会有序列化问题。
3. 例外情况:当DataFrame也会触发序列化问题
并不是所有DataFrame场景都不会有序列化问题——如果你用**UDF(用户自定义函数)**引用外部非序列化对象,就会回到和RDD一样的逻辑:
// 自定义一个不实现Serializable的类 class NonSerializableClass(val value: Int) object Example { import spark.implicits._ val badObj = new NonSerializableClass(1) // 不可序列化的对象 val df = sc.parallelize(Seq(("r1", 1))).toDF("ID", "b") // UDF引用了badObj,需要序列化闭包和badObj val addBadObj = udf((x: Int) => x + badObj.value) df.withColumn("plus", addBadObj($"b")).show() // 这里会触发序列化异常 }
因为UDF本质上还是一个闭包函数,Spark需要把它和引用的badObj一起序列化到Executor,所以如果badObj不可序列化,就会报错。
总结
- RDD:依赖JVM序列化传递闭包和外部变量,要求所有被引用的对象都必须实现
Serializable接口。 - DataFrame/Dataset:通过Catalyst优化器提前处理表达式,把外部常量嵌入执行计划,避免了闭包序列化的问题;但使用UDF时仍需注意序列化规则。
内容的提问来源于stack exchange,提问作者Ged
相关产品推荐
相关产品推荐

