Spark 数据倾斜排查:别一上来就加机器
Spark 数据倾斜排查:别一上来就加机器
Spark 任务跑得慢时,很多人的第一反应是加 executor、加内存、加并行度。资源加上去后,任务可能稍微快一点,也可能完全没改善。数据倾斜是最常见的原因之一:少数 key 拿走了大部分数据,某几个 task 跑很久,整体任务被拖住。倾斜不解决,加机器只是给问题铺更大的地毯。
Spark 数据倾斜排查要先找到倾斜点,再决定处理策略。
一、先看 Stage 和 Task
flowchart TD
A[Spark Job] --> B[Stage]
B --> C[Task Duration]
B --> D[Shuffle Read]
B --> E[Skewed Key]
Spark UI 是排查入口。看哪个 stage 最慢,task duration 是否长尾明显,shuffle read 是否集中在少数 task。正常任务几十秒,个别任务十几分钟,基本就要怀疑倾斜。
不要只看总耗时。总耗时告诉你慢了,stage 和 task 才告诉你慢在哪里。定位到具体 shuffle stage 后,再看上游 join、group by 或 distinct 操作。
为什么不直接看 SQL 而是先看 Stage 和 Task? Spark 是一个分布式执行引擎,同一行 SQL 经过优化器后会变成多个 Stage,每个 Stage 包含上百个 Task。SQL 只告诉你"做了什么",Stage 和 Task 才告诉你"谁在拖后腿"。举个例子,你写了一个三表 join 的 SQL,总耗时 40 分钟。从 SQL 层面你只能知道"这个查询慢",但从 Spark UI 的 Stage 页你能看到:Stage 2(第一个 join)花了 5 分钟,Stage 5(第二个 join)花了 35 分钟——问题在第二个 join 上,而且 Stage 5 里 99% 的 task 3 分钟结束,只有 1 个 task 跑了 35 分钟。这就把问题从"SQL 慢"收敛到了"第二个 join 的某个 key 倾斜"。底层原理:Spark 的 Stage 划分依赖于 Shuffle 边界,每个 Shuffle 后的 Stage 形成了天然的排查单位,你不应该跳过它。
二、找出倾斜 key
SELECT user_id, COUNT(*) AS cnt
FROM events
GROUP BY user_id
ORDER BY cnt DESC
LIMIT 20;
倾斜通常来自少数高频 key,例如默认值、未知渠道、大客户、热点商品、空字符串。先统计 key 分布,看看头部 key 占比。如果某个 key 占了 30% 数据量,后面的优化方向就很明确。
空值和默认值要特别注意。很多表里 unknown、0、空字符串会把大量数据聚到一起。数据清洗阶段如果不处理,Spark join 时就会爆出来。
为什么空值在 join 时特别危险? Spark 的 Shuffle 机制本质是按 key 的 hash 值分发数据。空字符串和 NULL 的 hash 值完全相同,这意味着所有 NULL 值都会被分到同一个 partition、同一个 task 去处理。如果一张表有 30% 的行 user_id 是 NULL,那这 30% 的数据全部塞给一个 task——这个 task 可能处理几百 GB 数据,而旁边 199 个 task 都在摸鱼。更阴险的是,NULL 值常常意味着"未登录用户""数据采集缺失",你在业务上根本不 care 这些行,但它们把你的计算资源全吃光了。我见过的最极端的案例:一个 day-level 的 ETL 任务跑了 6 小时,排查发现 95% 的时间都在给一个 task 做 NULL join——加一句
WHERE user_id IS NOT NULL,任务压缩到 15 分钟。
三、Join 倾斜要分情况处理
from pyspark.sql.functions import rand
salted = large_df.withColumn("salt", (rand() * 10).cast("int"))
大表 join 小表,可以考虑 broadcast join;大表 join 大表且少数 key 倾斜,可以对倾斜 key 加盐,把单个大 key 拆成多个分片;如果倾斜 key 没有分析价值,也可以单独过滤或归类处理。
加盐不是万能药。它会增加数据量和复杂度,也需要对另一侧数据做对应扩展。只有确认是 key 倾斜导致长尾任务时,才值得使用。
为什么加盐能解决倾斜,但会让结果验证变复杂? 加盐的原理是把一个大 key "A" 通过随机数拆成 "A_0"、"A_1"…"A_9" 十个子 key,分别分到不同 partition 处理。Shuffle 层面确实解决了——每个 task 的数据量降到原来的 1/10。但加盐的代价是数据膨胀,如果小表侧也需要对应展开(每个匹配行复制 10 份),那 join 中间结果的数据量是原来的 10 倍。更隐蔽的问题是,加盐后的 join 结果包含了重复的统计值,聚合时需要二次聚合才能还原真实数字——这一步很容易出错。我帮同事排查过一个"GMV 直接涨了 3 倍"的线上事故,原因就是加盐后的中间结果直接入库,没有把重复行合并。不是加盐不好,而是加盐的每一次数据膨胀都要配上确认步骤。
四、参数优化是最后一步
spark.sql.shuffle.partitions
spark.sql.adaptive.enabled
spark.sql.adaptive.skewJoin.enabled
Spark AQE 的 skew join 优化能自动处理部分倾斜,调整 shuffle partitions 也能改善任务粒度。但参数不是魔法。如果数据分布极端,或者 SQL 写法让大量数据提前膨胀,参数只能缓解,不能根治。
优化顺序应该是:确认慢 stage,识别倾斜 key,调整逻辑或数据,再调参数。反过来先调参数,容易把问题变得更难解释。
为什么先调参数会把问题变得更难解释? Spark AQE 的 skewJoin 优化确实是好东西,但它是一种"透明优化"——你改完参数后看到 task 分布"看起来正常了",就以为问题解决了。实际上,AQE 只是把一个大 key 自动拆开了,但数据膨胀、二次聚合这些底层问题一个没少。更致命的是,如果你先调了 AQE 再去看别的问题,你已经看不清"原始的数据分布长什么样"了——AQE 帮你掩盖了分布异常,你失去了诊断的基准线。就像体检前吃退烧药,体温是正常了,但病还在。正确的顺序一定是:先在不做任何优化的情况下跑一次,看清原始分布;然后再针对性地做数据层优化;最后才是参数。这样你才能建立"原因→手段→效果"的完整逻辑链,而不是"我改了某个参数→它快了→不知道为啥"的玄学调优。
处理完成后要留一份对比记录。优化前后 task 长尾、shuffle read、总耗时和资源消耗分别是多少,写进任务说明。Spark 优化很容易靠经验口口相传,记录下来才能变成团队资产。
踩坑提醒
-
不要在 Spark SQL 里用
DISTINCT来"顺手去重"——很多人在 join 之前习惯性加个SELECT DISTINCT以为能减少数据量,但实际上 Spark 的DISTINCT本身就是一个全量 Shuffle 操作。如果去重后的数据量并没有显著减少(比如只少了 5%),那你等于多跑了一次沉重的 Shuffle 去换了一次几乎无用的优化。先GROUP BY看 key 的重复率,低于 10% 就别去重。 -
Broadcast Join 的内存阈值不要设太大——把
spark.sql.autoBroadcastJoinThreshold调到 500MB 看起来很安全,但如果你的 Driver 只分配了 2G 内存,广播一个 400MB 的表就可能让 Driver OOM。Driver 不仅要存广播表,还要负责 Task 调度、指标收集,内存压力是全方位的。规则很简单:广播表大小不能超过 Driver 内存的 1/3。 -
AQE 的两个参数不能只开一个——
spark.sql.adaptive.enabled=true和spark.sql.adaptive.skewJoin.enabled=true必须同时开。很多人只开了前面那个,以为 AQE 已经生效了,其实 skew join 优化是单独控制的。不开 skew join 的 AQE 只会优化 partition 数量和 join 策略,不会自动处理倾斜 key 拆分——你等于开了半个优化,然后纳闷为什么倾斜还在。
五、总结
Spark 数据倾斜排查要从 Spark UI 的 stage 和 task 入手,找到 shuffle 长尾和倾斜 key,再根据 join 类型选择 broadcast、加盐、过滤或 AQE。
别一上来就加机器。数据分布不讲道理时,资源再多也会被少数倾斜 key 拖住。