Apache Spark — 慢作业的答案在执行计划和事件日志里
在一个机器人占 40% 点击的数据上测量并解决倾斜
目标
把由一个机器人制造了 40% 点击的数据与用户表连接,用事件日志测出记录集中到 shuffle 之后某一个任务上的样子。用 AQE 倾斜连接、加盐、把热点键单独拆出这三种方式来解决,并确认同一个热点键在 count 聚合中为什么不会成为问题。
为什么重要
“作业卡在 99%”这类报告,大多是数据倾斜。shuffle 由键的哈希来决定分区,所以相同的键必然去往同一个任务。如果一个键占了整体的 40%,一个任务就要扛下 40%,其余任务都结束之后,就只剩这一个在运行。再加核心也没用——那个任务是切不开的。 数据倾斜看分布比看时间更准确。在阶段里比较每个任务读取记录数的最大值和中位数,无论机器快还是慢,得到的都是同样的数字。本实验也是按这个比值,而不是按时间来判定。 处方有三种。AQE 会根据 shuffle 之后的实际大小,把大分区拆成多个碎片,并复制另一侧的配对(只要设置合适,不需要改代码)。加盐是给热点键附上随机标签,把它分散到多个分区,并把另一侧放大相应的份数。最简单的,是把热点键拆出来单独处理。另外,在聚合中,有些情况下倾斜原本就不是问题——因为有先在 Map 一侧进行缩减的部分聚合。
步骤
- 在 /root/spk/skew/common.py 中放入读取点击表和用户表的函数,用 /root/spk/skew/hot.py(应用
spk-skew-hot)把点击最多的前 5 名用户(并列时user_id升序)以 CSV(user_id、clicks)写入 /root/spk/skew/out/hot。 - 用 /root/spk/skew/plain.py(应用
spk-skew-plain,阈值 -1,关闭 AQE,shuffle 分区 8)把点击和用户连接,把按 segment 统计的点击数和 ms 之和以 CSV(segment、clicks、ms)写入 /root/spk/skew/out/by_segment。 - 从第 2 步应用的日志里,在读取了 shuffle 的阶段中,选出每个任务读取记录数最大值最大的那个阶段,把该阶段的编号、最大值、中位数(取较小的一个)写入 /root/spk/skew/out/skew.json。
- 用 /root/spk/skew/aqe.py(应用
spk-skew-aqe,阈值 -1,AQE 开启,shuffle 分区 16,倾斜阈值 256k,建议大小 64k)把同样的结果写入 /root/spk/skew/out/by_segment_aqe。 - 用 /root/spk/skew/salt.py(应用
spk-skew-salt,与第 2 步相同的设置),只给机器人的点击附上盐 0–7,并把用户一侧的机器人行放大成八份,按user_id和salt连接,把结果写入 /root/spk/skew/out/by_segment_salt。 - 用 /root/spk/skew/split.py(应用
spk-skew-split,与第 2 步相同的设置),让机器人的点击做广播连接,其余的做普通连接,再用unionByName拼起来,把结果写入 /root/spk/skew/out/by_segment_split。 - 用 /root/spk/skew/agg.py(应用
spk-skew-agg,关闭 AQE,shuffle 分区 8)把按用户的点击数以 Parquet 写入 /root/spk/skew/out/per_user,并把该应用中读取了 shuffle 的阶段里每个任务记录数的最大值和中位数(取较小的一个)写入 /root/spk/skew/out/agg_skew.json。 - 在 /root/spk/skew/report.md 中以
## 쏠림의 모양(韩文,意为“倾斜的样子”)、## 세 가지 처방(韩文,意为“三种处方”)、## 부분 집계(韩文,意为“部分聚合”)三节写成报告。第一节放入第 3 步的两个数字,第三节放入第 7 步的两个数字。
参考
- 原始数据:
/data/clicks/clicks.jsonl(user_id, page, ts, ms——24 万行)、/data/clicks/users.csv(user_id, segment——机器人的 segment 是bot)。 - 每个任务读取的记录数是
SparkListenerTaskEnd的Task Metrics.Shuffle Read Metrics.Total Records Read。中位数(较小的一个)是把 n 个值按升序排列后的第(n-1)//2个(从 0 开始)值。 - 脚本放在
/root/spk/skew里,并在那里运行。如果做一个读取日志的小工具(skewtool.py),第 3、7 步就是一行的事。 - 常见错误:在第 2 步开着 AQE 或广播,让倾斜消失了;把盐用随机数附加在两侧,导致配不上对(另一侧必须放大);各处方得到的 segment 结果不相同(答案必须相同)。
- 官方文档:Optimizing Skew Join · Adaptive Query Execution · Monitoring — Executor Task Metrics · Built-in Functions
找出热点键
在 /root/spk/skew/common.py 中放入读取点击表(/data/clicks/clicks.jsonl)和用户表(/data/clicks/users.csv)的函数,再以应用名 spk-skew-hot 创建 /root/spk/skew/hot.py,把点击数前 5 名的用户(点击数降序,相同时 user_id 升序)以带表头行的 CSV(user_id、clicks)写入 /root/spk/skew/out/hot。
化解倾斜的第一步,是知道哪个键是热点。看看第一名和第二名的差距。如果第一名是其余的好几倍,这个键一个就能拖住一个任务。
直接连接,会集中到一个任务上
以应用名 spk-skew-plain、设置 spark.sql.autoBroadcastJoinThreshold=-1、spark.sql.adaptive.enabled=false、spark.sql.shuffle.partitions=8 创建 /root/spk/skew/plain.py,把点击和用户按 user_id 连接,把按 segment 统计的点击数(clicks)和 ms 之和(ms)以带表头行的 CSV(segment、clicks、ms)写入 /root/spk/skew/out/by_segment。
关掉广播和 AQE,是为了故意让倾斜暴露出来。排序合并连接会把两侧都按 user_id 做 shuffle,所以机器人的 9 万多行点击会全部去往同一个分区。评分器会检查那个阶段里最大值是中位数的多少倍。
用数字表示倾斜——最大值和中位数
在第 2 步应用(spk-skew-plain)最近一份日志里,对读取了 shuffle 的每个阶段,汇总各任务的 Total Records Read,选出最大值最大的阶段,以 {"stage_id": 정수, "max_records": 정수, "median_records": 정수}(占位符为整数)写入 /root/spk/skew/out/skew.json。中位数是升序排列的 n 个值中第 (n-1)//2 个(从 0 开始)值。
连接阶段会同时读取两侧的 shuffle,所以装着机器人的那个分区的任务,要读取机器人的全部点击和机器人用户的一行。最大值 ÷ 中位数,就是倾斜的程度。这个比值超过 3,通常就称为“倾斜了”。
处方 1——AQE 倾斜连接
以应用名 spk-skew-aqe、设置 spark.sql.autoBroadcastJoinThreshold=-1、spark.sql.shuffle.partitions=16、spark.sql.adaptive.skewJoin.skewedPartitionThresholdInBytes=256k、spark.sql.adaptive.advisoryPartitionSizeInBytes=64k(保持 AQE 开启)创建 /root/spk/skew/aqe.py,把与第 2 步相同的结果写入 /root/spk/skew/out/by_segment_aqe。
当分区超过中位数的若干倍(skewedPartitionFactor,默认 5)并且大于阈值字节数时,AQE 就认为发生了倾斜。默认阈值(256MB)不适合这份小数据,所以调低了。如果在最终计划里看到 SortMergeJoin(skew=true),就是没有改代码而被拆开了。
处方 2——给热点键撒盐
以应用名 spk-skew-salt、与第 2 步相同的设置(阈值 -1,关闭 AQE,shuffle 分区 8)创建 /root/spk/skew/salt.py:点击一侧,如果是机器人就附加 0–7 之间的随机数 salt,否则附加 0;用户一侧,只把机器人的行按 salt 0–7 放大成八份(其余为 0),然后按 user_id 和 salt 连接,把结果写入 /root/spk/skew/out/by_segment_salt。
盐进入了连接键,所以机器人的点击会被分散到八个分区。另一侧的机器人行必须有八份,这样无论落到哪个盐上都有配对。结果必须与第 2 步完全相同,评分器会检查倾斜比值是否降到了第 2 步的一半以下。
处方 3——把热点键单独拆出
以应用名 spk-skew-split、与第 2 步相同的设置创建 /root/spk/skew/split.py:机器人的点击与机器人用户的一行做 broadcast 连接,其余的点击做普通连接,再用 unionByName 拼起来,写入 /root/spk/skew/out/by_segment_split。
如果热点键只有一个,并且知道它的名字,这就是最简单、最可靠的处方。那个键在另一侧只有一行,所以广播等于白送,其余部分则均匀分布。计划里必须能同时看到 BroadcastHashJoin、SortMergeJoin、Union。
count 为什么不倾斜——部分聚合
以应用名 spk-skew-agg、设置 spark.sql.adaptive.enabled=false、spark.sql.shuffle.partitions=8 创建 /root/spk/skew/agg.py,把按用户的点击数(user_id、clicks)以 Parquet 写入 /root/spk/skew/out/per_user,并把该应用中读取了 shuffle 的阶段里每个任务记录数的最大值和中位数(取较小的一个),以 {"max_records": 정수, "median_records": 정수}(占位符为整数)写入 /root/spk/skew/out/agg_skew.json。
同样有那个机器人,这次任务却分得很均匀。计划里的第一个 HashAggregate 是 partial_count——每个 Map 任务先按用户数一遍,所以通过 shuffle 传过去的,是每个 Map 任务每个用户一行。倾斜成为问题的,是这样无法先缩减的运算(连接、collect_list 之类)。
用数字留下倾斜和处方
在 /root/spk/skew/report.md 中写出 ## 쏠림의 모양(韩文,意为“倾斜的样子”)、## 세 가지 처방(韩文,意为“三种处方”)、## 부분 집계(韩文,意为“部分聚合”)三节。第一节用数字放入第 3 步的最大值和中位数,第三节用数字放入第 7 步的最大值和中位数。
第二节各用一行写出三种处方分别改变了什么(是代码还是设置,复制了什么)。也可以写一下在生产环境里会先尝试哪一种。