Apache Spark — 慢作业的答案在执行计划和事件日志里
写出分区数据湖,覆盖它,并减少小文件
目标
把订单按日期拆分,写成 Parquet 数据湖,确认四种保存模式(errorifexists、append、overwrite、ignore)中的三种实际做了什么。通过目录数看静态覆盖与动态分区覆盖的区别,并尝试减少小文件的两个旋钮(按分区列汇集、每个文件的行数上限)。
为什么重要
Spark 的写入不是一次就结束的事。每个任务先写到临时路径,作业成功后,提交协议把结果移到原位并留下 _SUCCESS 标记。所以写到一半的结果不会被读到。但如果保存模式选错了,这台诚实的机器就会诚实地酿成事故。
append 在一次重试中,会把同样的数据放进两次。overwrite 的默认是静态的,所以为了只修正一天而写的那一次,会把目标路径下所有的日期都删掉。动态分区覆盖(spark.sql.sources.partitionOverwriteMode=dynamic)只会替换所写数据中包含的分区。不了解这个差别,往生产数据湖里 overwrite 一次,几个月的数据就会一下子消失。
文件数是写入的任务数 × 该任务持有的分区值的个数。如果 shuffle 之后的任务都持有一小部分所有日期,那么每个日期目录里就会产生任务数那么多的小文件。先按分区列汇集(repartition("칼럼"),占位符为列名),或给每个文件的行数设置上限(maxRecordsPerFile),就可以控制大小。
步骤
- 在 /root/spk/write/common.py 中放入读取订单并加上
order_date的函数,再用 /root/spk/write/lake.py(应用spk-write-lake)把按order_date拆分的 Parquet 写入 /root/spk/write/lake/orders。 - 用 /root/spk/write/exists.py(应用
spk-write-exists)在同一路径上不指定模式再写一次,把错误条件名写入 /root/spk/write/out/exists.txt。 - 用 /root/spk/write/append.py(应用
spk-write-append)把 1 月的订单用 overwrite 写一次、再用 append 写一次到 /root/spk/write/lake/append,把重新读取的行数和唯一order_id数写入 /root/spk/write/out/append.json。 - 用 /root/spk/write/static.py(应用
spk-write-static)把全部数据写入 /root/spk/write/lake/static 之后,只用 2026-02-10 这一天的数据做 overwrite 重新写一遍,把剩下的分区目录数写入 /root/spk/write/out/static.json。 - 用 /root/spk/write/dynamic.py(应用
spk-write-dynamic,partitionOverwriteMode=dynamic)把全部数据写入 /root/spk/write/lake/dynamic 之后,只把 2026-02-10 的订单去掉取消订单的部分用 overwrite 写入,把剩下的分区目录数写入 /root/spk/write/out/dynamic.json。 - 用 /root/spk/write/small.py(应用
spk-write-small)把点击repartition(40)后写入 /root/spk/write/lake/clicks_many,repartition("page")后写入 /root/spk/write/lake/clicks_few,各自按page拆分,并把两个数据湖的数据文件数写入 /root/spk/write/out/files.json。 - 用 /root/spk/write/maxrec.py(应用
spk-write-maxrec)把订单repartition("channel")之后,以maxRecordsPerFile=20000按channel拆分写入 /root/spk/write/lake/by_channel。 - 在 /root/spk/write/report.md 中以
## 저장 모드(韩文,意为“保存模式”)、## 정적과 동적 덮어쓰기(韩文,意为“静态覆盖与动态覆盖”)、## 파일 수(韩文,意为“文件数”)三节写成报告。第二节放入第 4、5 步的分区数,第三节放入第 6 步的两个文件数。
参考
- 脚本放在
/root/spk/write里并在那里运行(from common import …)。 - 数据湖的数据文件以
part-开头。_SUCCESS是提交标记,.crc是校验和(不计数)。 - 动态覆盖通过会话配置(
spark.sql.sources.partitionOverwriteMode)或写入选项(.option("partitionOverwriteMode", "dynamic"))来开启。 - 常见错误:在多个步骤里改写同一个路径,把前面步骤的结果删掉了(每一步的路径都不同);以为静态覆盖只会改变一天;以为 maxRecordsPerFile 决定的是文件大小(字节)。
- 官方文档:Generic Load/Save Functions — Save Modes · Bucketing, Sorting and Partitioning · Parquet Files · Configuration — spark.sql.sources.partitionOverwriteMode
写出按日期拆分的数据湖
在 /root/spk/write/common.py 中放入按 schema 读取订单并加上 order_date = to_date(order_ts) 的函数,再以应用名 spk-write-lake 创建 /root/spk/write/lake.py,用 partitionBy("order_date") 把 Parquet 写入 /root/spk/write/lake/orders(mode("overwrite"))。
90 天的日期就是 90 个目录。作业成功才会产生 _SUCCESS——没有这个标记的目录可能是写到一半的,所以读取一方要先确认。评分器会检查目录数、行数和标记。
默认模式会停下——errorifexists
以应用名 spk-write-exists 创建 /root/spk/write/exists.py,在与第 1 步相同的路径上不指定模式再写一次,把异常的 getCondition() 写到 /root/spk/write/out/exists.txt 的第一行。第 1 步的数据湖必须原样保留。
默认保存模式是“已经存在就报错”。这是防止误写到同一路径的安全装置,也是在生产作业里需要养成明确指定模式的习惯的原因。
append——重试会造成两倍
以应用名 spk-write-append 创建 /root/spk/write/append.py,把 2026 年 1 月的订单用 overwrite 写一次、再用 append 写一次到 /root/spk/write/lake/append,把重新读取的行数和唯一 order_id 数,以 {"rows": 정수, "distinct_orders": 정수}(占位符均为整数)写入 /root/spk/write/out/append.json。
append 不看已有的文件,只追加新文件。用同样的输入重试作业,行数就会翻倍,而且文件名不同,表面上看起来完好无损。要做到幂等写入,就要用覆盖(可以的话用动态)或按键合并的方式。
静态覆盖——想修一天,结果全删了
以应用名 spk-write-static 创建 /root/spk/write/static.py,把全部订单按 order_date 拆分写入 /root/spk/write/lake/static 之后,只把 2026-02-10 这一天的数据,在同一路径上用 overwrite(配置保持默认值)再写一遍。把剩下的 order_date= 目录数,以 {"partitions_after": 정수}(占位符为整数)写入 /root/spk/write/out/static.json。
partitionOverwriteMode 的默认值是 static。overwrite 会在写入之前,把目标路径下整个删掉——不看所写数据里有哪些日期。这一步的结果是一次事故。目的就是亲眼看到它。
动态分区覆盖——只改那一天
以应用名 spk-write-dynamic、设置 spark.sql.sources.partitionOverwriteMode=dynamic 创建 /root/spk/write/dynamic.py,把全部订单按 order_date 拆分写入 /root/spk/write/lake/dynamic 之后,只把 2026-02-10 的订单去掉取消(cancelled)的部分,在同一路径上用 overwrite 写入。把剩下的 order_date= 目录数,以 {"partitions_after": 정수}(占位符为整数)写入 /root/spk/write/out/dynamic.json。
在动态模式下,只替换所写数据中包含的分区(这里是一天)。其余 89 天保持不变。评分器会检查目录数、那一天里是否已没有取消订单,以及其他天是否与原始数据一致。
小文件——任务数 × 分区值
以应用名 spk-write-small 创建 /root/spk/write/small.py,把 /data/clicks/clicks.jsonl 执行 repartition(40) 之后写入 /root/spk/write/lake/clicks_many,执行 repartition("page") 之后写入 /root/spk/write/lake/clicks_few,各自用 partitionBy("page") 写出,并把两个数据湖的数据文件(page=*/part-*)数,以 {"many": 정수, "few": 정수}(占位符均为整数)写入 /root/spk/write/out/files.json。
一个写入任务会为自己持有的每个分区值各打开一个文件。如果 40 个任务都持有七个页面,最多就是 280 个。先按分区列汇集,一个页面就只进入一个任务,每个页面一个文件(作为代价,如果某个页面很大,那个任务就会变重)。
给每个文件的行数设置上限
以应用名 spk-write-maxrec 创建 /root/spk/write/maxrec.py,把订单 repartition("channel") 之后,用 option("maxRecordsPerFile", 20000),以 partitionBy("channel") 写入 /root/spk/write/lake/by_channel。
按渠道汇集后,一个渠道进入一个任务,成为一个文件,而最大的渠道超过 17 万行。设置上限后,任务每写 2 万行就会新开一个文件。评分器会通过 Parquet 元数据,检查每个渠道的文件数是否等于向上取整(行数 ÷ 20000),以及每个文件的行数是否没有超过 2 万。
把写入规则变成团队规则
在 /root/spk/write/report.md 中写出 ## 저장 모드(韩文,意为“保存模式”)、## 정적과 동적 덮어쓰기(韩文,意为“静态覆盖与动态覆盖”)、## 파일 수(韩文,意为“文件数”)三节。第二节放入第 4、5 步剩下的分区数,第三节放入第 6 步的两个文件数。
如果是往生产数据湖里写的作业,试着把默认用哪种模式、禁止什么写成规则。数字要从你自己的结果文件里抄。