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

PySpark读取损坏记录并存储:代码未达预期的问题排查

问题分析与解决方案

问题原因

  1. 参数拼写错误:你使用的option("badrecords")不是PySpark CSV数据源的有效参数,正确参数名为badRecordsPath(仅Spark 3.0及以上版本支持)。
  2. 模式特性导致记录全部保留:PERMISSIVE是PySpark CSV读取的默认模式,它不会丢弃损坏记录,而是将格式错误的完整记录内容存入_corrupt_record字段,因此你会看到所有记录被加载。
  3. 损坏记录判定:ID为3、4的记录因为字段数量超过schema定义的6个(多了India字段),被识别为损坏记录;ID为5的记录仅address字段为空,符合schema中该字段允许为null的设置,属于正常记录。

修正方案

方案一:手动过滤并分离正常/损坏记录(兼容所有Spark版本)

先加载所有记录,再通过_corrupt_record字段筛选分离,分别存储:

from pyspark.sql.types import StructType, StructField, IntegerType, StringType

# 定义schema(保留_corrupt_record字段用于识别坏记录)
emp_schema = StructType(
    [
        StructField("id", IntegerType(), True),
        StructField("name", StringType(), True),
        StructField("age", IntegerType(), True),
        StructField("salary", IntegerType(), True),
        StructField("address", StringType(), True),
        StructField("nominee", StringType(), True),
        StructField("_corrupt_record", StringType(), True),
    ]
)

# 读取CSV文件
df_after = spark.read.format("csv")
            .option("header", "true")
            .schema(emp_schema)
            .option("mode", "PERMISSIVE")
            .load("/FileStore/tables/corrupt-2.csv")

# 筛选正常记录(_corrupt_record为null)
normal_df = df_after.filter("_corrupt_record IS NULL")
# 筛选损坏记录(_corrupt_record不为null)
bad_df = df_after.filter("_corrupt_record IS NOT NULL")

# 存储正常记录(可根据需求调整格式,比如csv/parquet)
normal_df.write.mode("overwrite").csv("/FileStore/tables/normal_records")
# 存储损坏记录
bad_df.write.mode("overwrite").csv("/FileStore/tables/badrecords")

方案二:使用badRecordsPath自动存储坏记录(Spark 3.0+)

利用Spark 3.0新增的badRecordsPath参数,配合DROPMALFORMED模式直接只加载正常记录,同时自动将坏记录写入指定路径:

from pyspark.sql.types import StructType, StructField, IntegerType, StringType

emp_schema = StructType(
    [
        StructField("id", IntegerType(), True),
        StructField("name", StringType(), True),
        StructField("age", IntegerType(), True),
        StructField("salary", IntegerType(), True),
        StructField("address", StringType(), True),
        StructField("nominee", StringType(), True),
    ]
)

# 读取时自动分离正常/坏记录
df_after = spark.read.format("csv")
            .option("header", "true")
            .schema(emp_schema)
            .option("mode", "DROPMALFORMED")
            .option("badRecordsPath", "/FileStore/tables/badrecords")
            .load("/FileStore/tables/corrupt-2.csv")

# 此时df_after仅包含ID为1、2、5的正常记录,坏记录已自动写入指定路径

内容的提问来源于stack exchange,提问作者user_program

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.23 13:12:32