Apache Spark — 慢作业的答案在执行计划和事件日志里
不要把模式交给推断,损坏的行不要丢弃而要分拣出来
一句话总结
CSV 里没有类型,所以必须有人来确定 schema。推断要把数据多读一遍,而且一个值就会让整列的类型摇摆。基本功是:自己提供 schema,对损坏的行用 PERMISSIVE 模式和损坏记录列筛出来而不是丢掉,并用数字留下有几行、为什么损坏。
为什么 schema 会成为问题
像 Parquet 这样的格式,文件内部带有 schema。CSV 只是用逗号分隔的文字。2 是整数,还是 2026-01-03 10:00:00 是时间戳,文件本身不会告诉你。所以 Spark 读取 CSV 时必须二选一:由人来提供 schema,或者由 Spark 看着数据去猜。
猜是有代价的。CSV 数据源文档写道,inferSchema 的默认值是 false,开启后“需要再扫描一遍数据”。推断使用的比例 samplingRatio 默认值是 1.0——也就是全部。如果是合作方每天发来的 10GB 文件,就等于每天要多读 10GB。
更大的问题是结果会随数据而摇摆。推断会选出能容纳该列所有值的类型。如果哪天数量列里混进了一行文字 two,那一列就不会被推断为整数,而是字符串。本实验中合作方的文件恰恰如此。交给推断的话,qty 和 order_ts 都会变成 string。昨天还正常运行的 sum("qty") 今天给出了奇怪的值,而代码一个字都没有改。
工作原理
PySpark 的 csv 文档建议,如果不想把全部数据多读一遍,就关掉推断或者自己提供 schema,而且 schema 除了 StructType,也可以用 DDL 字符串来传入。
ddl = ("order_id STRING, customer_id STRING, product_id STRING, qty INT, "
"order_ts TIMESTAMP, status STRING, channel STRING, _corrupt_record STRING")
df = (spark.read
.option("header", True)
.option("mode", "PERMISSIVE") # 기본값이지만 적어 둔다
.schema(ddl)
.csv("/data/partner/orders.csv"))
bad = df.where(F.col("_corrupt_record").isNotNull())
提供 schema 之后,Spark 会按该 schema 去转换值,遇到无法转换的行。怎么处理那一行,由 mode 选项决定。文档定义的三种模式如下。
- PERMISSIVE(默认):遇到损坏的行时,把原文字符串放进损坏记录列,把无法转换的字段置为 null。想保留原文,必须在用户 schema 中加入该名字的字符串列。
- DROPMALFORMED:忽略整个损坏的行。
- FAILFAST:遇到损坏的行就抛出异常。
损坏记录列的名字可以用 columnNameOfCorruptRecord 选项修改,不指定则使用配置 spark.sql.columnNameOfCorruptRecord 的默认值 _corrupt_record。
用 FAILFAST 失败时,外层异常是表示读取文件失败的 FAILED_READ_FILE,其原因是 MALFORMED_RECORD_IN_PARSING。这个错误的说明文字提示:如果想把损坏的记录当作 null 处理,就把 mode 设为 PERMISSIVE。
count 为什么会说谎
这是本模块中最容易出错的地方。同一份文档的 mode 说明里附有一条简短的警告:在列裁剪之下,CSV 只会尝试解析所需要的列,所以哪些行会被判定为损坏,取决于所要求的列集合。这个行为由 spark.sql.csv.parser.columnPruning.enabled 控制,默认是开启的。
count() 不需要任何列。所以解析器连把 qty 转换成整数都不试,只数行数。用 DROPMALFORMED 读取再计数,得到的是没有去掉损坏行的数字;用 FAILFAST 读取再计数,则平安无事地结束。在实验的 6,000 行文件里,两种模式都得到 6,000。只有进行需要所有列的写入,才会得到去掉了损坏行之后的 5,560 行。“用 FAILFAST 跑过了,所以文件是干净的”这个结论,离开了“当时要求了什么”就不成立。
文档与实际行为不一致的地方
文档中对 PERMISSIVE 的说明写道:字段数比 schema 少或多的行,在 CSV 中不被视为损坏记录。少了就用 null 补齐剩下的字段,多了就丢弃多出来的 token。然而在这个 Pod 的 Spark 4.2 上,用带有损坏记录列的 schema 去读,字段数不一致的行,其原文也会被放进损坏记录列。实验文件中这样被筛出来的行一共有 440 行。版本一变,这类细节就会变动。文档是出发点,判定要靠亲手数出来的数字。
在现场相遇的样子
悄悄丢弃的模式最危险。DROPMALFORMED 看上去很方便,但它不会留下任何记录说明丢掉了几行、为什么丢掉。合作方改了格式,哪怕一半的行都损坏了,管道仍然是绿灯。在现场,是用 PERMISSIVE 读取,把损坏记录列有值的行单独写进隔离表,并把它们的数量留作指标。
解析通过了,业务规则还是另一回事。数量 -3 能被顺利转换成整数。解析器只看格式,不懂含义。像负数的数量、未来的日期这样业务上不可能的值,要在解析之后用单独的条件筛出来,送进同一张隔离表。干净的结果,是通过两道关卡的行。
表头行默认会被忽略。文档中 enforceSchema 的默认值是 true,此时会强制套用指定的或推断出的 schema,同时忽略 CSV 的表头行。如果合作方调换了列的顺序,就不是按名字而是按位置放入,悄悄变成错误的值。文档也建议,为了避免错误的结果,把这个选项关掉。
实际工作中真正重要的事
- schema 写在代码里。推断要把数据多读一遍,而且一个值就会改变列的类型。
- 用 PERMISSIVE 和损坏记录列把损坏的行筛出来隔离。不是丢弃,而是计数并留下。
- 不要用 count 来确认损坏的行。因为列裁剪,如果不需要任何列,就根本不会解析。
- 格式检查和业务规则检查是不同的关卡。通过了解析的负数数量也要筛掉。
- 关掉 enforceSchema,并验证表头行。防止列顺序被调换的文件悄悄变成错误的值。
下一项实验要做什么
用 DDL schema 读取合作方发来的 6,000 行订单 CSV。先确认交给推断时 qty 和 order_ts 被推断为字符串,再用 PERMISSIVE 和损坏记录列把损坏的行筛出来,数一数有几行。用 DROPMALFORMED 读取同一个文件,比较只计数得到的值和写出后再计数得到的值有什么不同,并从用 FAILFAST 写入时失败的异常里拿到解析错误条件名。按逗号拆分后,按字段数来计数,比较损坏的行中字段数不对的有多少;最后把解析失败,以及空数量、负数数量这类业务规则违反,连同原因列一起单独隔离,得到干净的结果。