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

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

只有一个任务结束不了,原因在于一个键

在 TT Lab 中继续学习

一句话总结

shuffle 会把相同的键送到同一个任务,所以一个键如果占了数据的很大一块,收到这个键的那一个任务就会拖住整个阶段。这就是数据倾斜(skew),解决的办法有三种:AQE 倾斜连接、加盐、把热点键单独拆出。

为什么一个任务就成了问题

阶段要等最慢的那个任务结束才算结束。即使 199 个任务两秒就完成,只要有一个用了三分钟,阶段就是三分钟。这期间其余的核心闲着。把执行器增加一倍也没用——慢的那一个任务仍然在一个核心上运行。

为什么只有一个慢?RDD 编程指南把 shuffle 描述为一种 all-to-all 运算:为了把某个键的所有值汇集到一处,要读取所有分区。无论是连接还是聚合,都是由键的哈希来决定目标分区,所以相同的键必然去往同一个任务。这既是保证 shuffle 正确性的规则,同时也是数据倾斜的原因。24 万次点击里,如果有一个机器人制造了 9.6 万次,那么 40% 就会集中到收到这个用户 ID 的那一个分区。增加分区数也解决不了。因为一个键是切不开的。

如何识别

数据倾斜从平均值上是看不出来的。整个阶段读取的字节数看上去完好。要看的是按任务的分布。Web UI 文档介绍的阶段详情页面里有所有任务的汇总指标,其中 Duration 和 Shuffle Read Size / Records 会给出最小值、中位数、最大值。最大值是中位数的多少倍,就是数据倾斜的程度。如果是几十倍,就意味着一个任务几乎在独自干活。

事件日志里也有同样的数字。每个任务结束时留下的记录里,包含该任务通过 shuffle 读取的记录数,所以只要有一个日志文件,不用 UI 也可以亲手数出最大值和中位数。

第一条路——AQE 倾斜连接

性能调优文档中的倾斜连接优化,会把排序合并连接里发生倾斜的分区拆成大小相近的多个任务,并把另一侧与之配对的分区按需复制。spark.sql.adaptive.enabled 和 spark.sql.adaptive.skewJoin.enabled 必须都开启,两者默认值都是 true。

哪个分区发生了倾斜,要同时满足两个条件才能确定。大小要超过中位数的 skewedPartitionFactor 倍(默认 5.0),同时还要超过 skewedPartitionThresholdInBytes(默认 256MB)。因为有第二个条件,在小的实验数据上,不管多么倾斜,默认设置下什么也不会发生。这就是实验中要把阈值调低的原因。拆分时作为目标的大小是 advisoryPartitionSizeInBytes(默认 64MB)。

一旦生效,最终计划里会出现 SortMergeJoin(skew=true),其下的 shuffle 读取会变成 AQEShuffleRead … coalesced and skewed。有一条附带说明:如果拆分需要多产生一次 shuffle,AQE 默认不会应用。即便如此也想做的时候,要开启的是 spark.sql.adaptive.forceOptimizeSkewedJoin(默认 false)。

第二条路——加盐

没有 AQE,或者不是连接的地方发生的倾斜,要由人来解决。加盐是在大表一侧的键上,附加一个 0 到 N-1 之间的随机数字(盐),把一个热点键变成 N 个互不相同的键。这样哈希就会把它们分散到 N 个分区。作为代价,较小的一侧为了不失去配对,必须按所有盐值各复制一份。

from pyspark.sql import functions as F
N = 8
clicks_s = clicks.withColumn("salt", (F.rand(7) * N).cast("int"))
users_s = users.crossJoin(spark.range(N).withColumnRenamed("id", "salt"))
joined = clicks_s.join(users_s, ["user_id", "salt"]).drop("salt")

它的代价很清楚。较小的一侧会变成 N 倍。所以 N 只设到能化解倾斜的程度,并且在较小的一侧真的很小时才使用。

第三条路——把热点键单独拆出

如果热点键已知且只有几个,还有更简单的办法。只把这些键的行过滤出来单独处理,其余的照常连接,然后用 union 把两部分拼起来。热点键一侧,在另一张表里对应该键的行只有几行,所以用广播来连接,就完全没有 shuffle。其余部分则是倾斜已经消失的均匀数据。找热点键,只要数一数每个键的个数、看前几名就够了。

三者之中选哪个

顺序通常是这样的。首先在最终计划里确认AQE 是否已经在解决了。如果是排序合并连接,并且发生倾斜的分区超过了两个条件,那么连一个设置都不用改就能解决。想一想 AQE 拆分的方式,也能看出局限。它把发生倾斜一侧的分区拆成多个碎片,并给每个碎片整个地配上另一侧对应的分区。如果另一侧对应的分区也很大,复制成本就会很高;如果两侧以相同的键一起发生倾斜,即使把碎片拆开,配对的数量本身也不会减少。

如果 AQE 解决不了,就看热点键是不是固定的几个。像一个机器人账号、三个大客户这样能叫出名字的话,拆出来单独处理最简单,结果也最容易解释。如果热点键每天都在变,或者有几十个,加盐更好。因为即使不知道哪个键是热点,它也会把所有的键均匀地分散开。无论哪条路,结束之后都要重新测量每个任务的最大值和中位数,确认比值是否真的降低了。

倾斜被掩盖的情形——部分聚合

用同一份机器人数据运行 groupBy("user_id").count(),奇怪的是几乎看不到倾斜。看看计划,原因就在那里。shuffle 之前有 HashAggregate(partial_count),每个 Map 任务对每个键只发送一行部分合计。机器人的 9.6 万次,在 shuffle 之前就已经缩成了 Map 任务数那么多的几个数字。RDD 指南建议,按键求和或求平均时不要用 groupByKey,而要用 reduceByKey、aggregateByKey,道理也一样。

所以倾斜会在连接中、以及没有部分聚合的运算中暴露出来。窗口函数需要把分区键的所有行汇集到一个任务里,而收集列表的聚合,其部分结果本身就和原来的行一样大。不能因为聚合没问题,就推测连接也没问题。

在现场相遇的样子

第一,进度条停在 99%。只剩一个任务,要跑几十分钟。如果那个任务的 shuffle 读取记录数是其他任务的几十倍,那就是倾斜。

第二,null 键就是热点键。如果在缺失值的列里堆了几百万个 null,null 在哈希看来也是一个键,会汇集到一个分区。在连接条件中,null 和任何值都不相等,所以不会产生配对。但外连接要在结果里保留没有配对的行,所以这些行会原样被 shuffle。先把 null 键的行拆出来,连接之后再用 union 拼回去,结果不变,倾斜却消失了。

第三,倾斜会长大。从出现一个机器人或一个大客户的那天起,昨天还好好的作业就变慢了。代码没变,所以不看数据分布就找不到原因。如果每天记录连接键的前几名的行数,就能在作业变慢之前发现热点键在长大的苗头。

实际工作中真正重要的事

下一项实验要做什么

在一个机器人制造了 40% 点击的数据里,先找出热点用户。关掉广播和 AQE,用排序合并来连接,然后从事件日志里取出读取了 shuffle 的每个任务的记录数,计算最大值和中位数。接着开启 AQE 倾斜连接,但要针对小的实验数据调低倾斜阈值,确认最终计划里出现倾斜标记。再用加盐和把热点键单独拆出这两种方式重新解决同一个连接,比较三种处方的答案是否相同;最后用同样的两个数字确认,在按用户的 groupBy 聚合中,多亏了部分聚合,几乎没有倾斜。