Apache Spark — 慢作业的答案在执行计划和事件日志里
写入无法撤销,所以先定好模式和文件布局
一句话总结
Spark 的写入由三件事决定:保存模式(已经存在时怎么办)、分区目录(按什么形状摆放)、提交协议(怎样告知已经写完),这三样只要选错,就会在没有任何错误的情况下,让数据消失或膨胀。
为什么要单独学写入
读取和转换,不管错多少次,原始数据都原封不动。写入就不同了。一次覆盖写入就会把三个月的分区删掉,重试的追加写入会把同一天放进两次。而且这两种情况都会以成功结束。作业是绿灯,只有报告是错的。
加载与保存文档在介绍保存模式时,首先就警告了这一点。保存模式不使用锁,也不是原子的。而且覆盖写入会在写新数据之前删除已有数据。这意味着,覆盖写入的过程中作业如果挂了,就会出现既没有旧数据、也没有新数据的时刻。
工作原理——四种保存模式
同一份文档列出的模式有四个。
- errorifexists(或 error)——默认值。路径里已经有数据时抛出异常。在实验里会以
PATH_ALREADY_EXISTS条件暴露出来。 - append——在已有数据旁边追加新文件。用同样的数据运行两次,恰好变成两倍。
- overwrite——删除已有数据,重新写入。
- ignore——已经存在时什么也不做。与 SQL 的
CREATE TABLE IF NOT EXISTS相似。
默认是报错,这是个好设计。它能防止没有说明要做什么的写入动到别人的数据。问题出在,人们为了消除那个错误,习惯性地加上 overwrite。
分区目录与文件数
用 partitionBy("day") 写入,每个值会产生一个像 day=2026-01-03/ 这样的目录,这个列不在文件里面,而是在路径里。读取一方从路径里还原出值,如果有条件作用在这个列上,就会整个目录地跳过。同一份文档写道,这种方式会创建目录结构,所以不适合取值非常多的列。按用户 ID 做 partitionBy,目录会有用户数那么多。
文件数是更隐蔽的陷阱。一个写入任务会对自己持有的行的每个分区值写一个文件。如果 shuffle 之后的 200 个任务每个都持有 30 天的一小部分,文件最多会有 200 × 30 = 6,000 个。是每天 200 个只有几 KB 的文件。
(df.repartition("day") # 같은 날은 한 태스크로 — 날마다 파일 하나
.write.partitionBy("day")
.option("maxRecordsPerFile", 50000) # 너무 큰 날은 5만 줄씩 끊는다
.mode("overwrite")
.option("partitionOverwriteMode", "dynamic")
.parquet("/data/lake/orders"))
写入之前按分区列做 repartition,同一天的行就会汇集到一个任务里,每天一个文件。作为代价,大的那一天会变成一个巨大的文件。给它加上限的,是配置文档里的 spark.sql.files.maxRecordsPerFile。它是一个文件最多写入的记录数,默认的 0 表示没有限制。也可以作为写入选项传入。
想按取值很多的列来划分时,替代方案是分桶。根据同一份文档,bucketBy 与取值个数无关,它通过哈希把数据分装到固定个数的桶里,所以也可以用在唯一值无限增长的列上。作为代价,分桶和排序只适用于持久表(saveAsTable)。用只往路径里写文件的 save(),是留不下分桶的。可以这样区分记忆:像日期这样取值少、而且总是出现在查询条件里的列用 partitionBy,像用户 ID 这样取值多、又被当作连接键的列用 bucketBy。
静态覆盖与动态覆盖
对分区过的路径做 overwrite,会删掉什么?配置文档里的 spark.sql.sources.partitionOverwriteMode 决定这一点,默认是 STATIC。静态模式在写入之前,会先把对应目标的分区删掉。如果用 DataFrame 覆盖整个路径,目标就是整个路径。为了修正一天的数据,用只装着那天行的 DataFrame 去覆盖,结果其余日期全都消失,就是这种事故。
动态(dynamic)模式不预先删除,只替换实际写入了数据的分区。用一天的 DataFrame 覆盖,就只有那一天改变,其他天不变。文档写道,写入选项 partitionOverwriteMode 的优先级高于会话配置。所以像上面的例子那样,在写入的位置用选项写出来更安全。会话配置别人可能会改,而写在代码里的选项是跟着那次写入一起走的。
提交协议与 _SUCCESS
多个任务同时往一个目录里写,读取一方什么时候才能相信已经写完了?任务不是写在最终位置,而是写在 _temporary/ 下属于自己的位置。任务成功时,它的输出被提交,所有任务结束后,作业提交再把输出移动到最终位置。Hadoop 的 mapred-default.xml 说明,算法版本 1 的作业提交会把任务输出合并移动、删除 _temporary 之后,写入 _SUCCESS。所以 _SUCCESS 是“这个目录已经完整写完”的标记。失败作业的目录里是没有的。
Spark 通过配置来选择这个提交算法。配置文档中 spark.hadoop.mapreduce.fileoutputcommitter.algorithm.version 的默认值是 1,而 2 可能引发 MAPREDUCE-7282 这样的正确性问题。版本 1 依赖移动(重命名)。云集成文档警告说,在对象存储上重命名非常慢,而且一旦失败,状态就变得无法确定。这就是在本地磁盘或 HDFS 上很便宜的事情,到了 S3 上变得很贵的原因。
在现场相遇的样子
第一,想修正一天,结果全删了。这是静态覆盖的典型。覆盖分区数据湖的代码,要把动态模式以选项的形式写死。
第二,重试造成了两倍。用 append 写入的作业中途失败再重新运行,在先前已提交的部分之上,同样的数据又放进去了一遍。对于可以按天重新运行的作业,比起 append,对当天分区做动态覆盖更安全。因为运行多少次,结果都一样。
第三,下一个作业读到了写到一半的目录。如果后续作业一看到文件就开始读取,看到的就是写入过程中的状态。只要有一条“确认 _SUCCESS 再读取”的规则,就能防住这一点。还要记住,它不是出现在每一个分区目录,而是只在写入的最上层路径出现一个。读取一方要在数据湖最上层,而不是日期目录里找这个标记。
实际工作中真正重要的事
- 默认模式是报错。不要习惯性地加上 overwrite。
- 保存模式不是原子的。overwrite 是先删后写。
- 覆盖分区数据湖时,要在写入的位置用选项写明动态模式。
- 文件数是任务数 × 分区值。按分区列做 repartition,并用 maxRecordsPerFile 设置上限。
- _SUCCESS 是“已完整写完”的标记。读取一方要确认它。
下一项实验要做什么
按日期对订单做 partitionBy,写出数据湖,确认在同一路径上用默认模式再次写入时出现的错误条件。在 overwrite 之后再用 append 写一次,数一数行数变成两倍、而唯一订单号不变;再看用一天的数据做静态覆盖时,其他日期是怎么消失的,然后用动态覆盖只修正那一天。最后,比较把点击打散成 40 份写出的数据湖,与按分区列重新拆分写出的数据湖的文件数,并用 maxRecordsPerFile 给一个文件里装的记录数设置上限。