TT Lab
开始
学习 学习路径 课程

Apache Spark — 慢作业的答案在执行计划和事件日志里

不要把模式交给推断,损坏的行不要丢弃而要分拣出来

在 TT Lab 中继续学习

一句话总结

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 选项决定。文档定义的三种模式如下。

损坏记录列的名字可以用 columnNameOfCorruptRecord 选项修改,不指定则使用配置 spark.sql.columnNameOfCorruptRecord 的默认值 _corrupt_record。

把同一个 CSV 的四行用三种模式读取的示意图。对于数量字段里是 two 的第二行,PERMISSIVE 只把数量置为 null,并把原文放进损坏记录列,四行全部输出;DROPMALFORMED 悄悄去掉这一行,输出三行;FAILFAST 在这一行以 MALFORMED_RECORD_IN_PARSING 停止。下面写着一个陷阱:只做 count 的话会跳过解析,三种模式都会得到四行

用 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 的表头行。如果合作方调换了列的顺序,就不是按名字而是按位置放入,悄悄变成错误的值。文档也建议,为了避免错误的结果,把这个选项关掉。

实际工作中真正重要的事

下一项实验要做什么

用 DDL schema 读取合作方发来的 6,000 行订单 CSV。先确认交给推断时 qty 和 order_ts 被推断为字符串,再用 PERMISSIVE 和损坏记录列把损坏的行筛出来,数一数有几行。用 DROPMALFORMED 读取同一个文件,比较只计数得到的值和写出后再计数得到的值有什么不同,并从用 FAILFAST 写入时失败的异常里拿到解析错误条件名。按逗号拆分后,按字段数来计数,比较损坏的行中字段数不对的有多少;最后把解析失败,以及空数量、负数数量这类业务规则违反,连同原因列一起单独隔离,得到干净的结果。