Azure Databricks读取畸形数据CSV报错,如何定位并从源修复?
定位并修复CSV格式异常问题
一、先找出问题记录和异常字符
既然不想用dropmalformed或permissive绕过,第一步得先定位到具体出错的行和异常点。临时修改读取代码,开启坏记录捕获:
df2 = (spark.read.format('csv') .option("delimiter", ",") .option("quote", '"') .option("escape", '"') .option("header", "false") .option("path", '/mnt/d365/'+absolute+'/'+table_name+"/*.csv") .option("mode", "permissive") # 临时切宽容模式,不中断读取流程 .option("lineSep", "\r\n") .option("multiLine", "true") .option("columnNameOfCorruptRecord", "_corrupt_record") # 把坏记录原内容存到这个字段 .schema(schema) .load() ) # 筛选出所有有问题的记录 bad_records = df2.filter(df2._corrupt_record.isNotNull()) display(bad_records)
查看_corrupt_record字段的原始内容,重点排查这几个常见问题:
- 未闭合的双引号:比如某字段里单独出现一个
",没有按规则转义成"" - 未被引号包裹的逗号:字符串字段里直接包含逗号,导致Spark误将其识别为列分隔符
- 异常换行符:虽然开启了
multiLine,但如果存在单独的\n而非\r\n,可能导致行分割错误 - 不可见控制字符:比如
\0、\t这类特殊字符,可用正则提取排查:
from pyspark.sql.functions import regexp_extract # 提取行中的控制字符 bad_records.withColumn("bad_chars", regexp_extract("_corrupt_record", r"[\x00-\x1F\x7F]", 0)).display()
二、定位问题来源
- 找到具体出错的文件:
如果读取的是多份CSV,添加输入文件路径字段,就能精准定位到有问题的文件:
from pyspark.sql.functions import input_file_name df2_with_path = df2.withColumn("source_file", input_file_name()) bad_records_with_path = df2_with_path.filter(df2_with_path._corrupt_record.isNotNull()) display(bad_records_with_path)
- 回溯数据生成环节:
- 检查上游系统:比如D365导出的话,确认导出配置是否开启了引号包裹、特殊字符转义
- 检查ETL脚本:如果是自定义脚本生成CSV,查看是否正确处理了引号和逗号(比如是否启用了全字段引号、双引号转义规则)
- 检查数据源头:比如用户输入内容里是否包含未规范的特殊字符(比如表单中直接输入
"测试,内容")
三、从源头解决问题
- 上游系统/脚本修复:
- 如果是D365导出:调整导出配置,确保所有字符串字段用双引号包裹,字段内的双引号自动转义为两个双引号(
"") - 如果是Python生成CSV:用标准库
csv模块,指定正确参数规避格式问题:
import csv with open('output.csv', 'w', newline='') as f: # QUOTE_ALL 强制所有字段加引号,doublequote=True 用""转义字段内的" writer = csv.writer(f, delimiter=',', quotechar='"', quoting=csv.QUOTE_ALL, doublequote=True) writer.writerows(data)
- 临时预处理修复(上游无法快速修改时使用):
可以在读取前先清洗文件:
- 用shell命令批量替换未转义的引号(根据实际情况调整正则):
for file in /dbfs/mnt/d365/xxx/*.csv; do sed -i 's/"\([^",]*\)"/""\1"""/g' $file done
- 或者用Spark先读成文本行,清洗后再转成DataFrame:
text_df = spark.read.text('/mnt/d365/'+absolute+'/'+table_name+"/*.csv") # 把单独的"转成"",修复未闭合问题 cleaned_df = text_df.withColumn("cleaned_line", regexp_replace("value", '([^"])"([^"])', '$1""$2')) # 拆分列并映射到目标schema from pyspark.sql.functions import split final_df = cleaned_df.select(split("cleaned_line", ",").alias("cols")).selectExpr( *[f"cols[{i}] as {schema.fields[i].name}" for i in range(len(schema.fields))] )
内容的提问来源于stack exchange,提问作者Greencolor
相关产品推荐
相关产品推荐

