Apache Spark — 慢作业的答案在执行计划和事件日志里
减少 200 个洗牌分区产生的文件和任务
目标
对按客户的销售额汇总,分别用关闭 AQE 的默认值(shuffle 分区 200)、手工减小的值(8)、AQE 运行三次,用事件日志测出 shuffle 之后的任务数和结果文件数有什么变化。并确认 repartition 和 coalesce 在是否 shuffle 与文件数上有何不同。
为什么重要
groupBy、join 这样的宽转换需要把相同的键汇集到一个任务里,所以会产生 shuffle。前一个阶段把结果按分区写成文件,后一个阶段把这些碎片汇集起来读取。后一个阶段的任务数由 spark.sql.shuffle.partitions 决定,默认是 200。
200 是以大集群为基准的数字。对几 MB 的数据用 200,200 个任务各自只做一点点工作,而把结果写成文件,就会产生 200 个小文件。反过来,对几百 GB 用 200,一个任务要扛几 GB,还会溢写到磁盘。所以这个数字必须按数据大小来定,而 AQE 会在 shuffle 结束之后,根据实际大小合并较小的碎片,替你做这件事。
改变分区数的方法有两种。repartition(n) 通过 shuffle 重新生成碎片,coalesce(n) 则不经 shuffle 把相邻的碎片粘在一起。coalesce 便宜,但只能减少——不能把两个碎片变成四个。
步骤
- 在 /root/spk/shuffle/common.py 中放入生成按客户销售额 DataFrame 的函数,用 /root/spk/shuffle/noaqe.py(应用
spk-shuffle-noaqe,spark.sql.adaptive.enabled=false)把结果以 Parquet 写入 /root/spk/shuffle/out/by_customer_200。 - 数一数该结果文件夹里的数据文件(以
part-开头)个数,以整数写入 /root/spk/shuffle/out/files_200.txt。 - 用 /root/spk/shuffle/eight.py(应用
spk-shuffle-8,关闭 AQE,spark.sql.shuffle.partitions=8)把同样的结果写入 /root/spk/shuffle/out/by_customer_8。 - 用 /root/spk/shuffle/aqe.py(应用
spk-shuffle-aqe,设置保持默认值)把同样的结果写入 /root/spk/shuffle/out/by_customer_aqe,并把该应用中读取了 shuffle 的任务数,以整数写入 /root/spk/shuffle/out/aqe_tasks.txt。 - 把点击原始数据用 /root/spk/shuffle/repart.py(应用
spk-shuffle-repart)执行repartition(4)后写入 /root/spk/shuffle/out/clicks_repart,再用 /root/spk/shuffle/coalesce.py(应用spk-shuffle-coalesce)执行coalesce(4)后以 Parquet 写入 /root/spk/shuffle/out/clicks_coalesce,并把两个文件夹的数据文件数,以{"repartition_files": 정수, "coalesce_files": 정수}(占位符均为整数)写入 /root/spk/shuffle/out/repart.json。 - 从日志里求出第 3 步应用通过 shuffle 写出的字节数之和,以
{"app": "spk-shuffle-8", "shuffle_write_bytes": 정수}(占位符为整数)写入 /root/spk/shuffle/out/shuffle_bytes.json。 - 用 /root/spk/shuffle/narrow.py(应用
spk-shuffle-narrow)给已完成支付订单加上day列,只把order_id、customer_id、qty、day写入 /root/spk/shuffle/out/narrow。这个应用里不能有 shuffle。 - 在 /root/spk/shuffle/report.md 中以
## 200 개의 파티션(韩文,意为“200 个分区”)、## AQE 가 합친 것(韩文,意为“AQE 合并了什么”)、## repartition 과 coalesce(韩文,意为“repartition 与 coalesce”)三节写成报告。在每一节里放入第 2、4、5 步的数字。
参考
- 脚本放在
/root/spk/shuffle,并在那里执行spark-submit(from common import …)。 - 设置用
SparkSession.builder.config("키", "값")或spark-submit --conf 키=값(两处的占位符依次为键、值)来给。无论哪种,都会留在事件日志里。 - 如果做一个从日志里取数字的小工具
logtool.py,第 4、6 步就是一行的事了。读取了 shuffle 的任务,是SparkListenerTaskEnd的Shuffle Read Metrics.Total Records Read大于 0 的任务;通过 shuffle 写出的字节数,是Shuffle Write Metrics.Shuffle Bytes Written之和。 - 常见错误:用同一个名字运行多次,却数了旧日志(要数最近的那份);把
_SUCCESS或.crc也算进文件数;以为用 coalesce 可以增加分区。 - 官方文档:RDD Programming Guide — Shuffle operations · Performance Tuning — Adaptive Query Execution · Configuration · Monitoring — REST API·metrics
关闭 AQE,使用默认值 200
在 /root/spk/shuffle/common.py 中放入生成已完成支付订单按客户销售额(customer_id、revenue=qty×price 之和)DataFrame 的函数,再以应用名 spk-shuffle-noaqe、设置 spark.sql.adaptive.enabled=false 创建 /root/spk/shuffle/noaqe.py,把结果以 Parquet 写入 /root/spk/shuffle/out/by_customer_200。
关闭 AQE 后,shuffle 之后的阶段恰好用 spark.sql.shuffle.partitions 个任务来运行。评分器会从日志里检查那个阶段的任务数是不是 200,以及结果是否与从原始数据算出的按客户销售额相同。
数一数 200 产生的文件
数一数 /root/spk/shuffle/out/by_customer_200 里的数据文件(以 part- 开头的)个数,以一个整数写入 /root/spk/shuffle/out/files_200.txt。
shuffle 之后的一个任务写一个文件(空分区可能不写文件)。两万个客户的销售额只有几百 KB,看看文件有几个——这是小文件问题最常见的来源。_SUCCESS 和 .crc 不计入。
手工减到 8
以应用名 spk-shuffle-8、设置 spark.sql.adaptive.enabled=false、spark.sql.shuffle.partitions=8 创建 /root/spk/shuffle/eight.py,把同样的结果以 Parquet 写入 /root/spk/shuffle/out/by_customer_8。
这次 shuffle 之后的任务是 8 个,文件也是 8 个以内。结果必须与第 1 步一行都不差——分区数是划分工作的方法,而不是答案。
AQE 在 shuffle 之后合并
以应用名 spk-shuffle-aqe(设置保持默认值)创建 /root/spk/shuffle/aqe.py,把同样的结果写入 /root/spk/shuffle/out/by_customer_aqe,并从日志里数出该应用中读取了 shuffle 的任务数,以整数写入 /root/spk/shuffle/out/aqe_tasks.txt。
AQE 在 shuffle 的 Map 一侧结束之后,根据各分区的实际大小,把达不到目标大小(spark.sql.adaptive.advisoryPartitionSizeInBytes)的相邻碎片合并起来。shuffle 分区仍然是 200,但读取的任务会少得多。最终计划里会看到 AQEShuffleRead … coalesced。
repartition 与 coalesce——有 shuffle 和没有 shuffle
用 /root/spk/shuffle/repart.py(应用 spk-shuffle-repart)把 /data/clicks/clicks.jsonl 执行 repartition(4) 后写入 /root/spk/shuffle/out/clicks_repart,用 /root/spk/shuffle/coalesce.py(应用 spk-shuffle-coalesce)把同一份原始数据执行 coalesce(4) 后以 Parquet 写入 /root/spk/shuffle/out/clicks_coalesce。把两个文件夹的数据文件数,以 {"repartition_files": 정수, "coalesce_files": 정수}(占位符均为整数)写入 /root/spk/shuffle/out/repart.json。
先看看 18MB 的原始文件在 local[2] 下会被读成几块(rdd.getNumPartitions())。repartition 通过 shuffle 恰好生成四个,而 coalesce 不经 shuffle 只是把相邻的粘起来,所以不会比原来的碎片数更多。评分器还会通过日志检查是否只有 repart 应用存在 shuffle 写入。
shuffle 写到磁盘上的字节数
在第 3 步的应用(spk-shuffle-8)最近一份日志里,把所有任务的 Shuffle Write Metrics.Shuffle Bytes Written 相加,以 {"app": "spk-shuffle-8", "shuffle_write_bytes": 정수}(占位符为整数)写入 /root/spk/shuffle/out/shuffle_bytes.json。
shuffle 写入,就是 Map 一侧的任务把结果按分区拆开写到本地磁盘。多亏了部分聚合(HashAggregate 的 partial),只有 2 万个客户 × Map 任务数那么多的行会传过去,所以比原始数据小得多。这个数字就是 shuffle 的实际成本。
只有窄转换,就只有一个阶段
以应用名 spk-shuffle-narrow 创建 /root/spk/shuffle/narrow.py,给已完成支付订单加上 day = to_date(order_ts),只选出 order_id、customer_id、qty、day,以 Parquet 写入 /root/spk/shuffle/out/narrow。这个应用里一个 shuffle 都不能有。
where、withColumn、select 是看一行、出一行。不需要其他任务的数据,所以在一个阶段内连续运行(流水线化)。哪怕只加入一个 orderBy 或 distinct,也会产生 shuffle。
留下确定分区数的依据
在 /root/spk/shuffle/report.md 中写出 ## 200 개의 파티션(韩文,意为“200 个分区”)、## AQE 가 합친 것(韩文,意为“AQE 合并了什么”)、## repartition 과 coalesce(韩文,意为“repartition 与 coalesce”)三节。第一节用数字放入第 2 步的文件数,第二节放入第 4 步的任务数,第三节放入第 5 步的两个文件数。
每项各补一句:按这份数据的大小,你会把 shuffle 分区设为多少;既然有 AQE,是否还有必要手工去定。在生产环境中,这个决定的依据是结果文件数和 shuffle 字节数。