使用UDF时触发Task序列化错误,IntelliJ IDEA运行Spark代码报错
解决Spark Task序列化错误问题
嘿,这个序列化错误在Spark开发里太常见了,我来帮你捋清楚问题出在哪,怎么解决。
首先先确认下你的场景:
你提供的DataFrame内容:
+------+------+ |nodeId| p_i| +------+------+ | 26|0.6914| | 29|0.6914| | 474| 0.0| | 65|0.4898| | 191|0.4445| | 418|0.4445| +------+------+
你的代码结构大概是这样的(我补全了缺失的UDF部分,应该和你的实际代码一致):
class MyUtils extends Serializable { def calculate(spark: SparkSession, df: DataFrame): DataFrame = { def myFunc(a: Double): String = { var result: String = "-" if (a > 1) { result = "A" } return result } val myUdf = udf(myFunc) df.withColumn("result", myUdf(col("p_i"))) } }
问题根源
虽然你已经让MyUtils实现了Serializable,但嵌套在calculate方法里的myFunc会隐式持有MyUtils实例的引用。当Spark要把这个UDF发送到Executor节点执行时,它需要序列化整个MyUtils实例——如果MyUtils里藏着任何不可序列化的成员(哪怕是Scala编译器自动生成的隐藏成员),或者序列化机制无法处理这种嵌套函数的引用关系,就会触发Task序列化错误。
解决办法
给你几个靠谱的方案,任选其一就行:
方案1:把函数移到伴生对象中
把myFunc放到MyUtils的伴生对象里,这样它就和实例无关了,不需要序列化整个类实例:
class MyUtils extends Serializable { def calculate(spark: SparkSession, df: DataFrame): DataFrame = { val myUdf = udf(MyUtils.myFunc) df.withColumn("result", myUdf(col("p_i"))) } } // 伴生对象里的方法是静态性质的,不会持有实例引用 object MyUtils { def myFunc(a: Double): String = { var result: String = "-" if (a > 1) { result = "A" } result // Scala会自动返回最后一行结果,可省略return } }
方案2:定义为顶级函数
如果这个函数不需要依赖MyUtils的任何成员,直接把它放在类外面作为顶级函数,彻底切断和类实例的关联:
// 顶级函数,不属于任何类实例 def myFunc(a: Double): String = { var result: String = "-" if (a > 1) { result = "A" } result } class MyUtils extends Serializable { def calculate(spark: SparkSession, df: DataFrame): DataFrame = { val myUdf = udf(myFunc) df.withColumn("result", myUdf(col("p_i"))) } }
方案3:优化类的序列化(不推荐,除非必须保留嵌套)
如果一定要把函数留在方法内部,那得确保MyUtils的所有成员都是可序列化的,并且用@transient标记那些不需要序列化的成员(比如如果有SparkSession的成员,一定要标记为@transient,因为SparkSession本身不可序列化)。不过这个方案容易踩坑,还是优先前两个方案。
小提醒
在Spark里写UDF的时候,尽量避免让函数依赖外部类的实例——静态方法/伴生对象方法、顶级函数都是更安全的选择,能少很多序列化的麻烦。
内容的提问来源于stack exchange,提问作者ScalaBoy
相关产品推荐
相关产品推荐

