Scala中Spark DataFrame浮点数列过滤与序列化异常排查
搞定Spark DataFrame浮点列过滤的Task Not Serializable异常
嘿,我来帮你解决这个头疼的序列化问题!首先得搞明白为啥会触发这个异常:Spark在分布式跑任务的时候,得把闭包里的所有对象序列化后发到Worker节点执行。如果你的过滤逻辑里引用了不可序列化的对象(比如没实现Serializable接口的类,或者不小心在闭包里用了SQLContext这种本身不能序列化的组件),就会炸出这个错误。
下面给你两个靠谱的解决方案,既能实现浮点数值过滤,又能彻底避开序列化坑:
方案一:用Spark内置正则函数(首推!)
Spark SQL自带的rlike函数天生就是分布式安全、序列化友好的,根本不用自己写自定义UDF,直接就能用,完美避开序列化问题。
示例Scala代码
假设你要过滤的列存在columnNames列表里,咱们这么写:
import org.apache.spark.sql.functions._ // 定义匹配浮点数的正则(可以根据你的需求调整,比如要不要允许正负号、整数部分为空的情况) val floatRegex = "^[-+]?([0-9]*\\.[0-9]+|[0-9]+\\.[0-9]*)$" // 先读入原始CSV数据 val rawDf = sqlContext.read .schema(dfSchema) .csv(filename) // 循环过滤指定列,只留符合浮点格式的行 val filteredDf = columnNames.foldLeft(rawDf) { (df, colName) => df.filter(col(colName).rlike(floatRegex)) } // 最后写入CSV(记得按需设置header、分隔符这些参数) filteredDf.write .option("header", "true") .csv(outputFilename)
正则小说明
^[-+]?:允许数值开头带正负号[0-9]*\.[0-9]+:匹配.123、0.123这种格式[0-9]+\.[0-9]*:匹配123.、123.456这种格式- 如果只想要严格的“整数+小数点+小数”格式(比如不接受
.123或123.),把正则改成^[-+]?[0-9]+\.[0-9]+$就行
方案二:序列化安全的自定义UDF
要是你非得用自定义UDF,那得确保UDF里用的对象都是可序列化的。关键是把正则Pattern提前编译好,而且要放在UDF外面——Java的Pattern本身是实现了Serializable的,这样就没问题了。
示例Scala代码
import org.apache.spark.sql.functions.udf import java.util.regex.Pattern // 提前编译正则Pattern(一定要放在UDF外面,只编译一次,而且是可序列化的) val floatPattern = Pattern.compile("^[-+]?([0-9]*\\.[0-9]+|[0-9]+\\.[0-9]*)$") // 定义序列化安全的UDF,判断字符串是不是浮点数格式 val isFloat = udf((value: String) => { if (value == null) false else floatPattern.matcher(value).matches() }) // 读入原始数据 val rawDf = sqlContext.read .schema(dfSchema) .csv(filename) // 过滤指定列 val filteredDf = columnNames.foldLeft(rawDf) { (df, colName) => df.filter(isFloat(col(colName))) } // 写入CSV filteredDf.write .option("header", "true") .csv(outputFilename)
为啥之前的代码会踩序列化坑?
大概率是这几个原因之一:
- 在闭包里面(比如UDF内部)编译正则
Pattern,不仅性能差,还可能不小心引用了外部不可序列化的对象 - 过滤逻辑里直接用了
SQLContext实例(SQLContext本身不能序列化,绝对不能放在闭包里) - 用了自己写的、没实现
Serializable接口的类对象参与过滤
额外提醒
- 如果你的列不是字符串类型,记得先转成字符串再匹配,比如用
col(colName).cast(StringType) - 写CSV的时候,根据数据情况设置
nullValue、delimiter这些参数,避免导出的数据乱掉 - 测试的时候可以先用
sample(0.1)取10%的数据验证过滤效果,省得跑全量数据白忙活
内容的提问来源于stack exchange,提问作者Cassie
相关产品推荐
相关产品推荐

