Spark badRecordsPath未按预期写入脏数据至指定路径问题排查
解决方案
badRecordsPath确实是Databricks Runtime专属的扩展特性,并非Apache Spark原生支持的配置项,这就是本地Spark环境无法生效的核心原因。在Apache Spark原生环境中,可通过以下方式实现类似的脏数据处理效果:
方式1:直接丢弃脏记录
将读取模式设置为DROPMALFORMED,Spark会自动丢弃整条不符合Schema的记录:
val badDF = spark.read .format("csv") .option("header", true) .schema(schema_string) .mode("DROPMALFORMED") // 丢弃不符合Schema的记录 .option("path", "D:/spark_practice/data/bad_emp.csv") .load
方式2:捕获脏记录到单独字段
通过spark.sql.columnNameOfCorruptRecord配置,将不符合Schema的整条记录捕获到指定字段,后续可过滤并保存这些脏记录:
// 配置存储脏记录的字段名 spark.conf.set("spark.sql.columnNameOfCorruptRecord", "_corrupt_record") // 读取数据,脏记录会被存入_corrupt_record字段 val badDF = spark.read .format("csv") .option("header", true) .schema(schema_string) .option("path", "D:/spark_practice/data/bad_emp.csv") .load // 过滤脏记录并保存到指定路径 badDF.filter("_corrupt_record is not null") .write .mode("overwrite") .text("D:/spark_practice/bad_dir") // 处理正常记录 badDF.filter("_corrupt_record is null") .drop("_corrupt_record") .show()
方式3:字段级脏数据处理
如果仅需处理单个字段的脏数据,可自定义UDF转换非法值,同时收集脏记录:
import org.apache.spark.sql.functions._ // 自定义UDF:尝试将字符串转为Int,失败则返回Null val safeToInt = udf((s: String) => { try { Some(s.toInt) } catch { case _: Exception => None } }) // 先读取原始数据 val rawDF = spark.read .format("csv") .option("header", true) .option("path", "D:/spark_practice/data/bad_emp.csv") .load // 处理salary字段 val processedDF = rawDF.withColumn("salary", safeToInt(col("salary"))) // 保存脏记录 processedDF.filter(col("salary").isNull) .write .mode("overwrite") .csv("D:/spark_practice/bad_dir") // 展示正常记录 processedDF.filter(col("salary").isNotNull) .show()
内容的提问来源于stack exchange,提问作者Vivek Puurkayastha
相关产品推荐
相关产品推荐

