PySpark读取损坏记录并存储:代码未达预期的问题排查
问题分析与解决方案
问题原因
- 参数拼写错误:你使用的
option("badrecords")不是PySpark CSV数据源的有效参数,正确参数名为badRecordsPath(仅Spark 3.0及以上版本支持)。 - 模式特性导致记录全部保留:
PERMISSIVE是PySpark CSV读取的默认模式,它不会丢弃损坏记录,而是将格式错误的完整记录内容存入_corrupt_record字段,因此你会看到所有记录被加载。 - 损坏记录判定: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
相关产品推荐
相关产品推荐

