Spark任务序列化异常:排查AssetFileNameConventionCheck不可序列化问题
Spark任务序列化异常分析与修复
问题定位
从异常栈可直接定位核心问题:
Caused by: java.io.NotSerializableException: dna.dat.dam.datazones.AssetFileNameConventionCheck
你的AssetFileNameConventionCheck类未实现Serializable接口,且UDF内部引用了类成员变量regex(继承自父类)和dtf,导致Spark序列化UDF时会尝试捕获整个AssetFileNameConventionCheck实例,触发序列化失败。
修复方案
1. 让涉及类实现Serializable
Spark分布式计算要求跨节点传输的对象必须可序列化,因此所有相关类需要继承Serializable接口。
2. 优化UDF引用,避免捕获整个类实例
将依赖的正则、日期格式化器改为类内独立成员或静态成员,减少对类实例的依赖,降低序列化开销。
修改后的完整代码
import java.time.LocalDateTime import java.time.format.DateTimeFormatter import org.apache.spark.sql.functions._ import scala.util.matching.Regex abstract class StructuralCheck(val checkName: String) extends Serializable { val regex: Regex def check(colName: String, inputDf: DataFrame): DataFrame } class FileNameConventionCheck(override val regex: Regex) extends StructuralCheck("isFileNameValid") with Serializable { override def check(colName: String, inputDf: DataFrame): DataFrame = inputDf.withColumn("isFileNameValid", when(col("name").rlike(regex.regex), true).otherwise(false)) } class AssetFileNameConventionCheck() extends FileNameConventionCheck("""asset-created-(\d{14})\.json""".r) with Serializable { // 独立定义成员,避免UDF捕获整个类实例 private val dtf = DateTimeFormatter.ofPattern("yyyyMMddHHmmss") private val assetRegex = """asset-created-(\d{14})\.json""".r private val isValidFileFormatUdf = udf((name: String) => { assetRegex.findFirstMatchIn(name) match { case Some(matchedPattern) => val timestampString = matchedPattern.group(1) try { LocalDateTime.parse(timestampString, dtf) true } catch { case _: java.time.format.DateTimeParseException => false } case None => false } }) override def check(colName: String, inputDf: DataFrame): DataFrame = { inputDf.withColumn("isFileNameValidDateTime", isValidFileFormatUdf(col("name"))) } }
关键修改说明
- 所有类(
StructuralCheck、FileNameConventionCheck、AssetFileNameConventionCheck)添加Serializable继承,确保实例可序列化。 - 在
AssetFileNameConventionCheck中重新定义assetRegex,避免直接依赖父类的regex(父类成员仍会关联类实例)。 - 删除原代码中未使用的
SimpleDateFormat冗余代码。
额外优化建议
- 尽量避免在UDF中引用外部类成员,优先在UDF内部定义常量或使用静态成员,减少序列化对象体积。
- 可使用Spark内置函数替代自定义UDF,彻底规避序列化问题且性能更优:
inputDf.withColumn("timestamp_str", regexp_extract(col("name"), """asset-created-(\d{14})\.json""", 1)) .withColumn("isFileNameValidDateTime", when(to_timestamp(col("timestamp_str"), "yyyyMMddHHmmss").isNotNull, true).otherwise(false))
内容的提问来源于stack exchange,提问作者Aravind Yarram
相关产品推荐
相关产品推荐

