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

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

减少 200 个洗牌分区产生的文件和任务

在 TT Lab 中继续学习

目标

对按客户的销售额汇总,分别用关闭 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 便宜,但只能减少——不能把两个碎片变成四个。

步骤

  1. 在 /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。
  2. 数一数该结果文件夹里的数据文件(以 part- 开头)个数,以整数写入 /root/spk/shuffle/out/files_200.txt。
  3. 用 /root/spk/shuffle/eight.py(应用 spk-shuffle-8,关闭 AQE,spark.sql.shuffle.partitions=8)把同样的结果写入 /root/spk/shuffle/out/by_customer_8。
  4. 用 /root/spk/shuffle/aqe.py(应用 spk-shuffle-aqe,设置保持默认值)把同样的结果写入 /root/spk/shuffle/out/by_customer_aqe,并把该应用中读取了 shuffle 的任务数,以整数写入 /root/spk/shuffle/out/aqe_tasks.txt。
  5. 把点击原始数据用 /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。
  6. 从日志里求出第 3 步应用通过 shuffle 写出的字节数之和,以 {"app": "spk-shuffle-8", "shuffle_write_bytes": 정수}(占位符为整数)写入 /root/spk/shuffle/out/shuffle_bytes.json。
  7. 用 /root/spk/shuffle/narrow.py(应用 spk-shuffle-narrow)给已完成支付订单加上 day 列,只把 order_id、customer_id、qty、day 写入 /root/spk/shuffle/out/narrow。这个应用里不能有 shuffle。
  8. 在 /root/spk/shuffle/report.md 中以 ## 200 개의 파티션(韩文,意为“200 个分区”)、## AQE 가 합친 것(韩文,意为“AQE 合并了什么”)、## repartition 과 coalesce(韩文,意为“repartition 与 coalesce”)三节写成报告。在每一节里放入第 2、4、5 步的数字。

参考

关闭 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 字节数。