Spark任务序列化失败:VarScoreData引发NotSerializableException问题
Spark NotSerializableException 排查与解决(VarScoreData引发)
异常原因分析
从堆栈跟踪和代码来看,核心问题是自定义Case Class VarScoreData未被注册到Kryo序列化器中:
- 虽然Spark配置指定了Kryo序列化器,并使用了Hoodie的
HoodieSparkKryoRegistrar,但该注册器仅负责注册Hoodie框架相关类,不会自动注册业务自定义类VarScoreData。 - 当Spark需要将
VarScoreData实例在Driver与Executor间传输时,Kryo无法识别该类,直接抛出NotSerializableException。
另外需注意代码中VarScoreData(tel.substring(0, 2), payment, tel, varArray, scoreArr)里的payment变量:如果它是Driver端的外部引用对象且未实现序列化,也可能间接触发序列化问题,但本次异常直接指向VarScoreData,优先处理类的Kryo注册即可。
解决方案
方案1:自定义Kryo注册器,显式注册VarScoreData
创建自定义Kryo注册类,继承Hoodie的注册器并添加自定义类注册逻辑:
import com.esotericsoftware.kryo.Kryo import org.apache.spark.serializer.KryoRegistrator import org.apache.spark.HoodieSparkKryoRegistrar class CustomKryoRegistrar extends HoodieSparkKryoRegistrar { override def registerClasses(kryo: Kryo): Unit = { super.registerClasses(kryo) // 先注册Hoodie框架类 kryo.register(classOf[VarScoreData]) // 注册自定义业务类 } }
修改Spark配置,替换为自定义注册器:
val spark = SparkSession.builder() .config("spark.serializer", "org.apache.spark.serializer.KryoSerializer") .config("sql.catalog.spark_catalog", "org.apache.spark.sql.hudi.catalog.HoodieCatalog") .config("spark.sql.extensions", "org.apache.spark.sql.hudi.HoodieSparkSessionExtension") .config("spark.kryo.registrator", "com.your.package.CustomKryoRegistrar") // 替换为实际类路径 .getOrCreate()
方案2:显式让VarScoreData实现Serializable接口
虽然Scala Case Class默认会自动混入Serializable trait,但显式声明可避免潜在的序列化识别问题:
case class VarScoreData(part: String, day: String, tel: String, var_array: Array[Double], score_array: Array[Double]) extends Serializable
方案3:检查外部引用变量的序列化性
确认payment变量的类型是否实现Serializable:
- 如果是自定义对象,需让其实现
Serializabletrait; - 如果是大对象,可转为广播变量减少传输开销:
val paymentBro = spark.sparkContext.broadcast(payment) // 在map算子中使用:paymentBro.value
内容的提问来源于stack exchange,提问作者Ahian Liu
相关产品推荐
相关产品推荐

