
做Spark调优这几年最折磨人的问题就是数据倾斜。明明整个Job几百个Task前面跑得飞快最后卡在几个Task上下不来Stage进度停在99%点开Task看又没报错日志里刷的全是GC。更气人的是按网上教程把executor内存调大、并行度调高甚至重启了好几遍问题原封不动。如果你也遇到过这种场景那你踩的大概率是数据倾斜而不是单纯的资源不够。这篇文章我把自己在线上环境处理数据倾斜的完整经验梳理一遍包括现象判断、根因分析、定位方法和六类可落地的解决方案覆盖生产环境九成以上的倾斜场景特别适合被ETL任务折磨得凌晨爬起来看日志的同学参考。1. 先识别敌人数据倾斜的真实表现和很容易误判的几种情况1.1 最典型的三个信号数据倾斜的第一个信号不用看任何监控就能发现任务进度条卡住。一个Stage明明有几百个Task前90%都秒完最后几十个Task里有一两个一直转圈整体进度停在99%甚至直接显示失败。这不是偶发慢任务而是倾斜的标准姿势。第二个信号是Task之间的处理时间分布极度不均。打开Spark Web UI的Stage详情页你会看到绝大多数Task的Duration只有几秒到几十秒而某个或某几个Task的Duration是几百秒甚至几万秒。这个对比非常直观有时候甚至用不着看数据量光看时间对比就知道有问题。第三个信号藏在Executor日志里。倾斜Task通常伴随着频繁的Full GC、GC Overhead Limit Exceeded、或者直接OOM。因为单个Task要处理的数据量远超过其他Taskshuffle read阶段会把大量数据拉进内存内存扛不住就不断GCGC又拖慢速度形成恶性循环。1.2 假倾斜四种容易误判的情况不过我要泼盆冷水看到Task卡住别急着断定就是数据倾斜线上我见过太多假倾斜。第一种源文件本身大小不均。比如读HDFS时一个分区目录下某些文件是几MB另一些是几个GB这在读文件阶段就会造成Task时间不均衡。但它发生在输入阶段没有经过Shuffle不算严格意义上的数据倾斜。这种情况用Spark的maxSplitBytes相关参数调整读文件分区或者提前做文件合并就能缓解。第二种代码里Filter或Limit导致的分区不均。比如某张表做了很重的过滤后原本均匀分布的数据在结果里变得极不均匀后续算子又直接对结果做Shuffle问题才爆发。这种情况本质上是过滤后重新分布的需求没处理好。第三种Executor资源配置不均。比如开启动态分配后某些Executor拿到了更多核心某些Executor内存配得小导致Task执行速度天然不一样。这种问题从Task输入数据量看是均匀的但Duration差距拉大很多人会误判成倾斜。第四种Shuffle分区数设置得太少。我见过有人把spark.sql.shuffle.partitions设成10在一个大Job里某个Task正好承接了特别多数据看起来像倾斜实际上只是分区粒度过粗。这种情况把分区数调大或者开启AQE的分区合并逻辑往往立刻见效。判断的核心方法看Task的Shuffle Read Size和Records。如果这两个指标里只有少数Task明显高于多数Task才是数据倾斜如果大家的数据量差不多只是时间不同优先怀疑资源和分区设置问题。2. 追根溯源Shuffle分区机制与算子如何制造倾斜2.1 HashPartitioner在工作时到底做了什么要理解数据倾斜先得搞懂Shuffle的分区逻辑。Spark默认的HashPartitioner在Shuffle时会对key做一次哈希然后用哈希值对分区数取模决定这条记录进哪个分区partitionId hash(key) % numPartitions。这样做的本意是让数据尽量均匀地散到各个分区。但问题来了如果某个key本身就囊括了海量数据比如订单表里user_iddefault占了50%的数据这个key经过哈希后只会落到一个分区这一个分区就承接了海量记录其他分区全在看戏。这就是倾斜最直接的来源——key分布不均匀而不是Shuffle本身有bug。还有一类情况比单key更隐蔽哈希冲突导致的多个不相关key扎堆。虽然平均分配能摊平大部分情况但少数key的哈希值取模后恰好落在同一个分区如果这些key各自的数据量都不小这个分区也会显著增大。这种情况不好事先预判只能靠运行时统计去发现。2.2 这些算子最容易触发倾斜并不是所有算子都会引发倾斜。触发Shuffle的算子才是重灾区常见的有groupByKey、reduceByKey、aggregateByKey按key分组或聚合所有相同key的数据都会汇聚到同一个分区key分布不均直接影响任务均衡。join两张表按key关联Shuffle阶段每个分区内要匹配相同key的数据倾斜通常发生在少数key在两边都大量出现的场景比如把一个日活几千万的用户表和一张只有几百个热门用户的表join。distinct、countByKey本质也是按key分组热点key会把大量数据带到一个分区。sortByKey会触发全局排序的Shuffle虽然底层做了范围分区但边界设置不合理时同样可能让某个分区数据量暴涨。repartition/coalesce手动修改分区时如果指定了不恰当的分区表达式也会造成数据在分区层面不均匀。实际调优里Group和Join是出现倾斜最频繁的两个操作后面的解决方案我基本也会围绕这两类场景来讲。2.3 空值、热点Key与Join的特殊性有一个非常容易被忽视的倾斜源空值。SQL里那些user_id IS NULL的记录在Shuffle时key就是null不管代码怎么写null的哈希结果都一样全部进同一个分区。如果这张表里脏数据比例很高比如10%的记录user_id是null那么光这一个分区就扛了10%的数据量任务不慢才怪。还有一个经典场景表A是用户行为明细几亿条表B是用户维度表几万条按user_id join。如果只有少数热门用户贡献了海量行为比如测试账号默认用户这种这些用户的行为数据会在同一个分区内堆积join的shuffle阶段就会倾斜。这里有个区分维度如果是大表join超大表、命中率很高属于普遍性问题如果是大表join维度表、倾斜集中在极少数key上那就非常适合用广播或定向扩容来解决。理解根因之后你会发现所有解决方案本质上都是在做同一件事让原本集中在某个分区的数据在计算过程中被打散再用局部结果合并的方式保证最终结果正确。3. 定位慢Stage的完整排查链路UI、日志与执行计划配合使用3.1 Stage页面上怎么快速圈定嫌疑Task打开Spark Web UI进入对应Application的Stages页面你有三个地方要重点看。第一个是Stage列表。找到耗时最长、显示Succeeded数量最多但进度没到100%的那个Stage点进去。如果是跑得很慢但还在跑的Task通常State是Running。第二个是Task级别的时间分布。在Task表里按Duration排序快速比对最大时间和中位数。如果中位数是10秒最大时间是6000秒差距超过两个数量级基本可以锁定倾斜。同时看Shuffle Read Size / Records列找那些输入数据量明显偏大的Task记下它的Partition ID。第三个是Shuffle Read或者Shuffle Write总量。如果一个Stage的Shuffle Read总量特别大说明有多个Task在拉取大量数据这时还要回去看上一个Stage的Shuffle Write分布看哪个分区写出去的数据特别多。Shuffle Write分布是判断倾斜发生在哪个Shuffle阶段的关键。提示判断一个Task是否因为数据倾斜变慢最硬核的指标是Shuffle Read Size/Records而不是Duration本身。Duration受GC影响很大数据量相同也可能时快时慢。3.2 日志里的关键线索Web UI能定位到哪个Task但要知道为什么这个Task慢还得看Executor日志。一个是driver端日志。Driver的Log里会输出每个Spark Job的粒度信息如果某个Stage反复重试同一批Task多半和Executor端OOM、连接超时有关。另一个是Executor端日志。进入对应Worker节点找到Executor的stdout/stderr日志。倾斜Task最常见的日志特征是大量GC日志、OutOfMemoryError、FetchFailedException拉取shuffle数据失败、Connection refused。出现FetchFailed时通常是因为某个Executor已经被数据量压垮导致其他Executor拉不到它的shuffle数据。值得注意的一个细节是如果错误信息集中在某个partition id上并且这个partition在其他Executor上也反复失败那就不仅仅是资源问题而是这个partition的数据量已经大到任何一个Executor都扛不住这时候靠调大executor内存往往没用得从数据层面拆散它。3.3 用执行计划和Key分布确认根因UI和日志定位到了Task和Stage那怎么知道这段代码到底是group还是join触发的倾斜呢最直接的办法是翻物理执行计划。在Spark SQL场景下执行df.explain(true)找到那个有多个Exchange节点的位置Exchange就是ShuffleDataSize和行数统计可以帮你判断哪个Shuffle阶段写出的数据不均衡。如果发现某个Exchange的输出分区里数据量差异巨大基本上代码里就是那条SQL的某一步有问题。另一个更接地气的办法直接对疑似倾斜的key做一次统计。比如你在join时怀疑某个user_id是罪魁祸首可以用样例数据建议取1%先做df.filter(...) .select(user_id) .groupBy(user_id) .count() .orderBy(col(count).desc) .show(10)打印出Top 10的key和记录数看看是不是集中在某几个高热度key上。这一招屡试不爽比各种监控都直观。核心技巧是一定要弄清楚倾斜发生在哪个算子之后、是哪个键在倾斜否则后面的优化大概率白做。4. 常规解法三件套两阶段聚合、空值处理与广播Join4.1 两阶段聚合加随机前缀再拆再合两阶段聚合是我在groupByKey/reduceByKey场景下第一个会试的方案思路非常朴素既然所有相同key的数据都会压到一个分区那我先给key加一个随机前缀让相同key变成很多个分身后的key数据自然就散开了子任务处理完之后去掉前缀恢复原key再做一次全局聚合。适用场景很明确聚合类算子、倾斜key有限、数据语义允许中间临时打散。不适合的场景是那些对key顺序有要求的操作比如要按时间排序取每条记录的窗口逻辑一旦加了随机前缀顺序就乱了。Scala示例// 第一步加随机前缀打散倾斜key val salted df .withColumn(salt, (rand() * 100).cast(int)) .withColumn(salted_key, concat(col(user_id), lit(_), col(salt))) // 第二步局部聚合 val partial salted .groupBy(salted_key) .agg(sum(amount).as(partial_sum)) // 第三步去掉前缀全局聚合 val result partial .withColumn(user_id, expr(substring_index(salted_key, _, 1))) .groupBy(user_id) .agg(sum(partial_sum).as(total_sum))这里有两个关键决策点。第一个是随机前缀的范围。范围太小比如只加0~9最多也就是拆成10份倾斜key如果太大还是拆不匀范围太大比如0~9999局部聚合时会产生大量临时key反而增加Shuffle数据量和计算开销。我的经验是先看key分布如果Top 1的key占30%数据量我一般用100~500的前缀范围兼顾拆散效果和中间数据膨胀。第二个是数据正确性。用随机前缀后一个原始key的聚合结果被拆成了多个分组必须先做局部聚合sum、count、max这种都可以再做一次全局聚合。如果是count(distinct)这类去重聚合两阶段聚合会失效因为同一个key的重复值被拆到不同分组后局部去重结果再合并时无法判断全局唯一性这种情况需要用别的方案。另外要提醒一个细节如果原始key本身就可能包含下划线用substring_index拆回原始key会出错。更稳妥的做法是单独保存一个原始key字段不要靠拼接字符串来反解否则数据错了很难排查。4.2 空值独立处理先隔离再合并空值倾斜是我线上遇到过最多、但解决起来最省事的一类。前面说过null key在Shuffle时会全部集中在一个分区。处理思路非常直接把null和非null的数据分开处理null不参与倾斜Shuffle。具体有三种做法。做法一如果业务允许直接丢弃null那就在Shuffle之前加过滤条件把user_id IS NULL的记录过滤掉。这个方法最简单但前提是下游不需要这些数据。做法二如果null数据有用那就给null一个随机前缀让null也散到不同分区。做法三把null单独拎出来和非null数据分别聚合后再union回去。注意这里的空值不限于SQL里的null还包括业务上约定俗成的特殊值比如unknown、default、-1这种。这些特殊key在数据里往往代表未知但会成为热点的温床。我在一次日志分析任务里就碰到过app_name未知占了一半数据量的情况处理方式和null完全一样。4.3 广播Join从源头消灭Shuffle很多倾斜问题在join场景下根子是一个超大表join一个小表。比如大表几亿条明细小表几千条维度数据。这时最理想的做法不是去拆key而是让小表广播到每个Executor上每个Executor本地直接查没有Shuffle倾斜自然消失。Spark默认有一个自动广播阈值spark.sql.autoBroadcastJoinThreshold 10485760 // 默认10MB如果小表统计大小低于这个阈值Spark会自动用Broadcast Hash Join。但很多场景下小表虽然内存放得下大小却超过10MB或表的统计信息不准自动广播没有触发这时需要手动指定val result bigDF.join(broadcast(smallDF), Seq(user_id), left)SQL方式SELECT /* BROADCAST(small_table) */ * FROM big_table JOIN small_table ON big_table.user_id small_table.user_id使用广播join的注意事项广播表不能太大。广播表要复制到每个Executor如果小表实际有200MBExecutor有100个那就是每个Executor多200MB内存容易OOM。一般建议单表1GB以下、总Executor内存充裕时才考虑强制广播。别盲信阈值自动判断。表数据经过filter后可能很小但统计信息还是原始表的大小Spark不会每次都做精确统计这种情况下自动广播可能不触发手动加hint反而更可控。广播join不适合大表join大表。如果两边都是几十GB级别广播本身就是灾难需要走后面的扩容方案。5. 顽固热点Key与高基数场景扩容Join和定向Skew Join实操5.1 扩容Join给大表撒盐、给小表做复制如果两个表都很大又不能过滤数据广播方案失效就需要用扩容的办法。思路我在第2节讲过就是把倾斜key的数据量通过随机前缀拆开到多个分区同时把小表按同样的前缀做多份复制让原本一对一的join变成多对多join最后再过滤掉多余的匹配结果。实操步骤第一步识别倾斜key。先用count统计出记录数超过阈值的key比如设定阈值是总记录数的1%或者某个固定条数比如100万条。第二步对大表倾斜key加随机后缀。把倾斜key的数据行随机附加0~N-1的编号非倾斜key保持原样。第三步对小表做全量复制N份。复制时给复制后的key分别加上编号0~N-1这样每个编号都有小表的完整数据。第四步join之后由于大表的倾斜key被随机分到N个分区而每个分区都能在小表复制里找到对应key所以不会丢数据且每个分区的数据量被摊平。// 大表倾斜key加随机后缀 val bigSalted bigDF.withColumn( join_key, when(isSkewKey(col(user_id)), concat(col(user_id), lit(_), (rand() * N).cast(int))) .otherwise(col(user_id)) ) // 小表复制N份并加上不同后缀 val smallSalted smallDF .crossJoin(Range(0, N).toDF(salt)) .withColumn(join_key, concat(col(user_id), lit(_), col(salt)))这里需要注意因为扩容后会产生数量膨胀N的取值要控制。比如大表倾斜key有1000万条N10那扩容后是1亿条参与join小表也复制了10份。在资源允许的情况下通常N取10~20已经足够。5.2 定向Skew Join只改造倾斜部分扩容join有一个让人头疼的问题如果只是少数热点key倾斜其他key都很健康但扩容方案会对所有key重建连接键、让小表整体复制N份白白增加很多数据处理量、占用大量内存。更聪明的方法是只对倾斜key做特殊处理非倾斜key走正常的join路径最后把两段结果合并。这也是大多数skew join框架的实现思路。实现框架类似这样val skewedKeys skewKeyDF.as[String].collect().toSet val bcSkew spark.sparkContext.broadcast(skewedKeys) val trimmedBig bigDF.filter(!col(user_id).isin(bcSkew.value.toSeq: _*)) val skewedBig bigDF.filter(col(user_id).isin(bcSkew.value.toSeq: _*)) // 非倾斜部分正常join val normalResult trimmedBig.join(smallDF, Seq(user_id), left) // 倾斜部分对倾斜key走扩容join val skewedResult skewJoin(skewedBig, smallDF, N) val result normalResult.union(skewedResult)这个方案有两个优势一是只有倾斜key的数据量被放大整体额外开销可控二是小表不需要整体复制N份只在倾斜分支的join里复制。但也有代价代码复杂度增加而且如果倾斜key的数量特别多collect到driver端的集合会很大反而成为瓶颈。所以当倾斜key数量少比如几个到几十个时这个方法性价比最高。5.3 扩容倍数到底怎么定很多人第一次用扩容join时会纠结N取多少。我给一个简单算法先统计倾斜key的记录数S看当前分区数P目标是让倾斜key分摊到和正常分区差不多的数据量。假设正常分区平均数据量为M则建议扩容倍数N≈S/M再上下浮动。实际操作里如果无法精确估计N取10~20通常能解决大部分问题然后通过Spark UI观察Task时间是否均衡做二次调整。另外提醒一个容易踩的坑扩容之后的join结果中大表的倾斜key会和小表的复制数据产生N份匹配最后一定要根据业务逻辑去重或者合并。比如你join后只需要维度表的属性维度表复制N份只是为了保证每个拆分后的key都能关联上结果中同一行会出现N次需要在最后按原有主键去重。不加这一步下游不仅数据量爆炸还会出脏数据。6. Spark 3.x AQE自动倾斜优化的参数与实测配置6.1 AQE之前为什么调优这么费劲Spark 3.0之前最让人头疼的地方在于优化器在生成执行计划时使用的是表的统计信息而统计信息和真实运行时的数据分布往往差得很远。比如一张表在生成计划时是均匀的但实际跑到一半发现倾斜这时候计划已经定死了你只能靠人工加各种hint、改代码来救场。AQEAdaptive Query Execution的诞生就是让Spark在运行时根据真实数据分布动态修正执行计划。对于倾斜来说它最大的价值是能自动识别Sort Merge Join中的倾斜分区并自动做拆分优化。6.2 开哪些参数、配多少合适AQE本身是一个参数开关但实际使用中只开它不够还需要配合几个子参数。我以Spark 3.1为例列一份常用的配置spark.sql.adaptive.enabledtrue spark.sql.adaptive.coalescePartitions.enabledtrue spark.sql.adaptive.coalescePartitions.parallelismFirsttrue spark.sql.adaptive.skewJoin.enabledtrue spark.sql.adaptive.skewJoin.skewedPartitionFactor5 spark.sql.adaptive.skewJoin.skewedPartitionThresholdInBytes256MB这几个参数的含义和我的经验值参数默认值建议值作用说明spark.sql.adaptive.enabledfalse(3.0)/true(3.1)trueAQE总开关spark.sql.adaptive.skewJoin.enabledtruetrue自动处理倾斜的开关spark.sql.adaptive.skewJoin.skewedPartitionFactor55~10分区中位数乘以该系数超过则视为倾斜分区spark.sql.adaptive.skewJoin.skewedPartitionThresholdInBytes256MB256MB~512MB分区数据量超过该值才会触发倾斜判定spark.sql.adaptive.advisoryPartitionSizeInBytes64MB64~128MB自动合并分区时期望的目标分区大小实际配置时要根据资源情况和分区数据量来调。skewedPartitionFactor设得太小容易把正常分区误判为倾斜导致不必要的拆分设得太大则识别不出倾斜。我一般会先看历史任务的Shuffle Read尺寸分布取中位数的5~10倍作为判定线。6.3 AQE的边界哪些倾斜它管不了虽然AQE很好用但不是所有倾斜都能自动解决。我梳理了几个它管不了、还需要人工兜底的情况聚合类倾斜AQE主要针对Sort Merge Join的倾斜分区做动态拆分对于groupByKey/reduceByKey这类聚合场景并不会自动加盐该卡还是卡需要回落到两阶段聚合方案。读取阶段的源数据不均如果HDFS上文件本身就有大文件小文件问题AQE不会帮你重新切分。广播阈值误判AQE会动态决定是否把某个join转成广播join但如果小表太大或者大表太小它的判断也不一定合理还是要靠人工统计。手动repartition引起的倾斜代码里显式调用的repartition(col)会覆盖AQE的自动分区优化。我自己的使用习惯是AQE默认打开作为第一道防线但遇到历史遗留的复杂任务、或者是写入性能瓶颈时一定会回到人工分析执行计划和Task分布不把所有希望都压在自动调优上。7. 一次线上倾斜调优全记录从卡死3小时到20分钟7.1 线上场景去年有个任务经常在凌晨挂掉是一个用户行为ETL把用户点击行为明细一天2亿条和用户标签表300万条做left join然后按标签聚合输出到结果表。最初用的是普通Sort Merge Join没加任何优化。每天凌晨跑到一半某个Stage的进度就停在99%等了一个多小时都不动然后开始报OOM整个Job重试之后又重复。7.2 排查时间线第一步打开Web UI发现卡死的Stage是第三个Stage也就是join之后的聚合Stage。为什么判断是join后的聚合而不是join本身呢因为我看Shuffle Write分布时第二个Stage也就是join的map端写出的数据里有一个分区写了3GB而中位数分区只有200MB。第二步继续看Task详情卡住的Task对应的Shuffle Read达到3GB而其他Task只有200MB左右。很明显是某个key的记录特别多导致它在hash分区时把大量数据带到了同一个分区。第三步用sample统计key分布跑出来的结果让我哭笑不得Top 1的key是user_id0——就是埋点上报时用户未登录产生的默认值占当天数据的38%。这一个key直接把join后按user_id聚合的Stage打崩了。7.3 方案落地与效果对比因为业务方明确说user_id0的数据不能丢过滤不行。我没有直接用整体扩容而是分了三步走。第一步空值特殊值隔离。在join之前先把user_id0的行单独拆出来不走正常join。具体做法是先给它一个随机后缀打散并复制标签表让这一批数据完成打散关联。第二步主路径用广播join。用户标签表300万行估算压缩后不到300MB虽然超过默认10MB阈值但我的executor内存充裕直接手动广播主路径不产生shuffle。第三步对仍然存在的groupByKey聚合Stage做两阶段聚合把非倾斜key去掉后剩下的key数量并不大整体耗时可控。改造后实测结果维度优化前优化后Job运行时间约3小时约20分钟失败/重试次数每日常规OOM重试0最大Task Shuffle Read3GB400MB集群CPU利用率峰值低、空闲多稳定高水位7.4 这个case沉淀下来的几个经验第一瓶颈往往不在join本身而在join之后的聚合。看问题一定顺着执行计划往下看不要只看报错的Stage本身。第二特殊值和空值要当作头号嫌疑犯。答案一目了然但我仍然习惯性地对每个字段做一遍空值和默认值统计能帮我提前发现是不是在埋点阶段就埋下了隐患。第三方案组合比单一方案更可靠。我这次同时用了特殊值打散、广播join、两阶段聚合三种手段层层拆解单独用任何一个可能都解决不了问题。尤其是user_id0这种值一旦打散后续聚合还需要全局聚合所以两阶段聚合是必要的兜底。第四也是我踩过最痛的一次坑改完第一版后直接把spark.sql.shuffle.partitions调到了2000想通过并行度解决一切。结果数据是均匀了但小文件写出来几百份下游hive表小文件暴增查询反而更慢。后来我把分区数调到一个合理值并配合coalesce处理输出问题才真正解决。这提醒我提高并行度只能缓解表面症状真正的解法一定是把单个分区的数据量降下来同时控制好输出文件数。