You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.05 04:15:30