Apache Spark — 慢作业的答案在执行计划和事件日志里
用三种模式读取合作方 CSV 中损坏的行
目标
确定 schema,读取合作方发来的 6,000 行订单文件,用数字确认 PERMISSIVE、DROPMALFORMED、FAILFAST 三种模式分别如何处理损坏的行,然后再加上业务规则,把数据分成干净的表和隔离表。
为什么重要
CSV 里没有类型。只要某一列里混进一次“two”,inferSchema 就会把这整列退回成字符串,而且这件事表现出来的不是错误,而是悄悄发生的类型变化。昨天还是整数的列今天变成字符串,后面的合计就全都坏了。所以管道不去推断 schema,而是确定它再读取。
确定了 schema,接下来就要决定不符合 schema 的行怎么办。Spark 的 CSV 读取有三种模式。PERMISSIVE 保留损坏的行并把原文放进损坏记录列,DROPMALFORMED 丢弃,FAILFAST 停止。哪个是对的,由业务决定——如果是要计算金额的表,悄悄丢弃是最危险的。
有一个陷阱。CSV 解析器只解析需要的列。像 count() 这样完全不需要列的行动算子不会解析行,所以用 DROPMALFORMED 读取再计数,连损坏的行也会全部数进去,FAILFAST 也不会停下。如果用 count() 来做验证,就等于什么也没有验证。
步骤
- 把原始文件的七个列写成 DDL 保存到 /root/spk/schema/schema.ddl(
qty是 INT,order_ts是 TIMESTAMP,其余是 STRING)。 - 以应用名
spk-schema-infer创建 /root/spk/schema/infer.py,用inferSchema=True读取,并把推断出的 schema 的simpleString()写入 /root/spk/schema/out/inferred.txt。 - 以应用名
spk-schema-load创建 /root/spk/schema/load.py,在第 1 步 schema 的末尾加上_corrupt STRING,用 PERMISSIVE 读取,并把全部内容以 Parquet 写入 /root/spk/schema/out/permissive。 - 以应用名
spk-schema-drop创建 /root/spk/schema/dropmal.py,用 DROPMALFORMED 读取,把count()得到的值,以及以 Parquet 写入 /root/spk/schema/out/dropped 之后再计数的值,以{"count_only": 정수, "written": 정수}(占位符为整数)写入 /root/spk/schema/out/drop_count.json。 - 以应用名
spk-schema-failfast创建 /root/spk/schema/failfast.py,把用 FAILFAST 读取的内容以 Parquet 写出,并把导致失败的解析错误条件名写到 /root/spk/schema/out/failfast.txt 的第一行。 - 以应用名
spk-schema-width创建 /root/spk/schema/width.py,按行读取原始文件(不含表头行),把按逗号拆分后的字段数对应的行数,以{"칸 수": 줄 수}(占位符依次为字段数、行数)写入 /root/spk/schema/out/width.json。 - 以应用名
spk-schema-clean创建 /root/spk/schema/clean.py,在解析成功的行中,只把qty大于或等于 1 的行以 Parquet 写入 /root/spk/schema/out/clean,其余的加上reason列(解析失败为parse,数量问题为qty)写入 /root/spk/schema/out/quarantine。 - 在 /root/spk/schema/report.md 中写出
## 추론이 틀린 곳(韩文,意为“推断出错的地方”)、## 세 가지 모드(韩文,意为“三种模式”)、## 격리한 줄(韩文,意为“被隔离的行”)三节。第二节放入损坏行数和用 DROPMALFORMED 写出的行数,第三节放入干净行数和被隔离的行数,都用数字写出。
参考
- 原始文件:
/data/shop/orders_dirty.csv(1 行表头行 + 6,000 行)。时间格式是yyyy-MM-dd HH:mm:ss,所以要用timestampFormat告诉它。 - 损坏记录列必须直接放进 schema。名字必须与
columnNameOfCorruptRecord选项相同。 - 只针对损坏记录列的查询,需要重新解析原始文件,可能会被禁止。把读取一次的结果
cache()起来再过滤。 - FAILFAST 的异常在 Python 里可能不会直接返回条件名。消息里的方括号
[...]中包含着条件名。 - 常见错误:用
count()来确认模式的效果、没有把损坏记录列放进 schema、把空数量(null)送进干净的行。 - 官方文档:CSV Files · Data Types · Datetime Patterns · Error Conditions
先确定 schema
按原始文件 /data/shop/orders_dirty.csv 表头行的顺序,把七个列写成一行 DDL 字符串,保存到 /root/spk/schema/schema.ddl。qty 是 INT,order_ts 是 TIMESTAMP,其余是 STRING。
DDL 的形式是 이름 형, 이름 형, ...(占位符依次为列名、类型)。用 head -1 看表头行,并照着它的顺序写——CSV 是按位置而不是按名字来对应列的。
看看推断漏掉了什么
以应用名 spk-schema-infer 创建 /root/spk/schema/infer.py,用 inferSchema=True、header=True 读取原始文件,并把 df.schema.simpleString() 写入 /root/spk/schema/out/inferred.txt。
推断是把数据多扫一遍的作业。在事件日志里看作业数,比确定了 schema 时要多。看看结果里 qty 和 order_ts 被识别成什么类型、为什么,去 grep 一下原始文件。
PERMISSIVE——不丢弃,把原文装起来
以应用名 spk-schema-load 创建 /root/spk/schema/load.py,用第 1 步的 DDL 末尾加上 _corrupt STRING 所得的 schema,以 mode=PERMISSIVE、columnNameOfCorruptRecord=_corrupt、timestampFormat=yyyy-MM-dd HH:mm:ss 读取,并把全部内容以 Parquet 写入 /root/spk/schema/out/permissive。
PERMISSIVE 不会改变行数。6,000 行原样保留,类型对不上或字段数不同的行,其原文会进入 _corrupt,对应字段变成 null。评分器会重新读取原始文件数出损坏的行,并与你的 _corrupt 比较。
DROPMALFORMED——count() 会说谎
以应用名 spk-schema-drop 创建 /root/spk/schema/dropmal.py,用第 1 步的 schema 和 mode=DROPMALFORMED 读取,把(一)count() 得到的值、(二)以 Parquet 写入 /root/spk/schema/out/dropped 之后再读取计数得到的值,以 {"count_only": 정수, "written": 정수}(占位符为整数)写入 /root/spk/schema/out/drop_count.json。
两个数字不一样才是正常的。count() 不需要任何列,所以 CSV 解析器不解析行,而不解析就不知道它是否损坏。只有做需要所有列的行动算子,才会真正丢弃。
FAILFAST——让它停下,并读出原因
以应用名 spk-schema-failfast 创建 /root/spk/schema/failfast.py,把用第 1 步的 schema 和 mode=FAILFAST 读取的内容以 Parquet 写出(路径随意),捕获失败,把解析失败的错误条件名(以 MALFORMED_… 开头的那个)写到 /root/spk/schema/out/failfast.txt 的第一行。
写入会解析所有列,所以作业会在第一个损坏行处失败。异常是层层包裹的,外层是读取文件失败,解析失败在原因一侧。把消息中方括号里的名字全部提取出来看看。评分器还会检查该应用里是否确实有失败的作业。
数一数字段数不同的行
以应用名 spk-schema-width 创建 /root/spk/schema/width.py,用 spark.read.text 按行读取原始文件,对去掉表头行后的行,按逗号拆分后的字段数来计数,以 {"칸 수": 줄 수}(占位符依次为字段数、行数)写入 /root/spk/schema/out/width.json(键是字符串)。
F.size(F.split("value", ",")) 就是字段数。这个原始文件的引号里没有逗号,所以简单拆分就行(如果是带引号的 CSV,这个方法就是错的)。比较一下字段数不是 7 的行,占第 3 步损坏行中的多少。
再加上业务规则——干净表和隔离表
以应用名 spk-schema-clean 创建 /root/spk/schema/clean.py,用 PERMISSIVE 读取之后,只把解析成功(_corrupt 为 null)且 qty 大于或等于 1 的行写入 /root/spk/schema/out/clean(不含损坏记录列),其余的加上 reason 列(parse 或 qty)以 Parquet 写入 /root/spk/schema/out/quarantine。
有些行格式是对的,但业务上是错的——负数数量、空数量。解析器不把它们当作损坏的行,所以规则要自己写。如果把原因留在列里,退回给合作方时,该让他们改什么马上就清楚了。
留下丢弃了什么、为什么丢弃
在 /root/spk/schema/report.md 中写出 ## 추론이 틀린 곳(韩文,意为“推断出错的地方”)、## 세 가지 모드(韩文,意为“三种模式”)、## 격리한 줄(韩文,意为“被隔离的行”)三节。第二节放入第 3 步的损坏行数和第 4 步写出后数得的行数,第三节放入第 7 步的干净行数和被隔离的行数,都用数字写出。
数字要从你自己的结果文件里重新数过再抄。可以当作是要发给合作方的邮件——收到了多少行,有多少行、因为什么损坏,又有多少行因业务规则退回。