ARTICLE DETAIL

资讯详情

深耕网站建设与运营推广的一线实战洞察。

Spark数据倾斜调优实战:从定位到解决的完整指南

Spark数据倾斜调优实战:从定位到解决的完整指南 做Spark到现在数据倾斜是我遇到最多、也最让人头疼的性能问题。不管你是写ETL、跑报表还是在数仓里做离线任务只要碰上group by、join、reduceByKey这些会发生shuffle的算子就有可能踩中“绝大多数task几秒钟跑完剩下几个task跑半小时还没结束”的经典坑。这篇文章不打算照搬官方文档我就拿自己实际调优的经验把Spark数据倾斜从识别、定位到解决的完整路径拆一遍同时给出可以直接抄的代码、参数和踩坑记录。适合正在接手Spark作业、发现作业越跑越慢却不知道从哪下手的开发同学参考。下面的内容我会按照“现象与根因、定位手段、解决方案、实战复盘、选型对照”这几个重点来展开。前两章帮你看清问题第三章是全文最核心的部分后面两章告诉你不同方案到底怎么选、落地时有哪些坑。1. 先弄明白Spark数据倾斜到底是怎么发生的1.1 数据倾斜的现场长什么样数据倾斜最直观的表现就是同一个Stage里task运行时间严重不均。正常情况下几十个task一起跑完成时间应该差不多整体进度条是均匀推进的。一旦出现倾斜你会看到99%的task早就FINISHED了只剩一两个task卡在RUNNING状态进度条在它身上一动不动。严重的时候这几个task还会因为内存不够直接OOM把整个Application拖垮。我在生产环境里见过的典型场景一个是统计订单量时按城市group by结果某个一线城市的订单量比其他城市高出几个数量级整个任务就变成“一个task拖着所有task走”。另一个是join场景订单明细表关联用户维度表但是join key里有大量null值这些null值经过hash分区后全部落到同一个分区直接把那个分区对应的executor打爆。为什么会出现这种情况这里要回到Spark的shuffle机制。Spark在做group by、join这些操作时需要把具有相同key的数据通过网络拉取到同一个task上这个动作叫shuffle。分区的规则通常是根据key的hash值对分区数取模所以只要某个key特别“重”它所在的分区数据量就会远超其他分区。换句话说数据倾斜的本质不是资源不够而是key分布不均导致分区数据量失衡。1.2 倾斜的根源key分布不均和shuffle机制我把常见的倾斜来源总结成以下几类排查的时候可以对着看数据本身分布不均这是最普遍的。比如按地域统计热门城市的数据量天然就大按品牌统计头部品牌的订单量天然就高。这类倾斜是业务属性决定的没法从源头消除只能在计算层处理。空值或默认值大量堆积比如某些表的外键经常为空join时这些空key全部hash到同一个分区。还有一种情况是业务上用了魔法值比如“-9999”“unknown”这类填充字段也会造成同样的效果。聚合key粒度过粗当group by的key本身区分度很低时比如按是否VIP分组、按性别分组很容易一个key里堆积大量数据。输入数据不均衡某些上游产生的大文件特别大或者读取小文件时分区数量设置不合理导致map端就已经严重不均后续计算跟着一起倾斜。这类问题的共同点在于它们都发生在“数据按照key重新分布”的环节。所以解决方案也基本围绕一个思路让倾斜的key不再集中在同一个分区或者干脆从计算模型上绕过shuffle。2. 数据倾斜怎么定位别一上来就调参数2.1 Spark UI里看task分布一眼锁定嫌疑task很多同学一遇到作业慢就先去调executor内存、调并行度这其实是本末倒置。定位倾斜最直接的办法是看Spark Web UI。进入Application详情页后点击Stages标签找到运行时间最长的那个Stage再点到这个Stage的Task列表页面。正常情况下列表里所有task的Duration和Shuffle Read Size都差不多如果有某个task的Duration是其他task的几十倍而且Shuffle Read Size也明显偏大那基本就锁定数据倾斜了。我自己的习惯是这样的在Task列表页面按“Shuffle Read Size”从大到小排序再按“Duration”从大到小排序两个指标都排前面的task就是重点怀疑对象。同时看一眼它的Locality Level和GC Time如果GC Time也很高说明这个task在拉取数据后内存压力很大更坐实了倾斜的判断。把那几个task记下来点进去可以看它们处理的partition id后面写代码定位倾斜key时会用到。另外如果作业里有多个Stage要注意看每个Stage的Shuffle Write Size和Shuffle Read Size。有时候倾斜发生在Reduce阶段之前但Map阶段就已经有人把任务写歪了这也需要追踪。2.2 用采样脚本确认倾斜key定位到具体是哪个Stage倾斜之后下一步就是找出到底是哪个key导致的问题。这个阶段不需要全量扫描可以先采样快速验证。以PySpark为例我常用的方式是对疑似倾斜的key字段做一次分组统计按数量倒序排取前几个看看分布情况from pyspark.sql.functions import col, count # 抽样减少扫描开销 sample_df df.sample(withReplacementFalse, fraction0.01, seed42) sample_df.groupBy(key_column) \ .agg(count(*).alias(cnt)) \ .orderBy(col(cnt).desc()) \ .show(10, False)如果这前几个key的cnt加起来占了抽样总量的大部分那倾斜key基本就水落石出了。注意抽样比例不要太小否则可能漏掉真正的热点key我一般先用1%的样本看趋势确认后再用全量统计单独提取那几个热点key。还有一种情况是lift join时两个表key关联度不同。比如左表的key分布正常但右表里某些key对应的行数特别多join结果就会膨胀。这时候就要分别统计左右两个表的key分布找到重叠的“大头key”。这个动作在上线前做一次能省下后面无数次排查时间。3. Spark数据倾斜解决方案五类招数从入门到进阶3.1 提高shuffle并行度最简单的尝试但别指望根治先聊门槛最低的一种调大shuffle并行度。Spark里有两个关键参数一个是spark.sql.shuffle.partitions控制SQL和DataFrame操作中shuffle产生的分区数默认是200另一个是spark.default.parallelism主要影响RDD操作默认取executor总核数。调大分区数之后原本一个重分区里的数据会被拆成更多小分区单个task处理的数据量就变小了。常见的做法是spark-submit \ --conf spark.sql.shuffle.partitions1000 \ ...或者直接在代码里设置spark.conf.set(spark.sql.shuffle.partitions, 1000)并行度设置多少比较合理我习惯按shuffle的数据量来估算目标是让每个task处理的数据量控制在128MB到256MB之间。比如shuffle读的总数据量大概200GB那并行度设在800到1600之间比较合适。算完可以留点余量但也不要盲目设到一两万task太多会带来额外的调度开销而且还可能产生大量小文件影响下游读取。说实话调并行度在很多场景下只是“缓解”而不是“根治”。如果倾斜key本身数据量特别大哪怕每个分区只放它的一部分数据由于其他task很快都跑完了只剩那几个分区还在处理整体耗时依然被拖住。所以这个方法适合倾斜程度不严重、或者临时救急的场景。真正要解决问题还得靠下面几种思路。3.2 两阶段聚合加盐聚合类场景最常用两阶段聚合是处理group by、reduceByKey这类聚合操作用得最多的方案思路非常直白既然一个key太重那就先给它加个随机前缀让这个key拆散成多个“假key”分别聚合最后再把前缀去掉做第二次聚合。以“按城市统计销售额”为例假如某个城市销售额占了80%我们可以先给city_id加一个0到9的随机后缀这样同一个城市的数据会被分到10个不同的分区里聚合最后再把这10个部分结果加起来。PySpark代码如下from pyspark.sql.functions import col, rand, sum as _sum, lit # 假设df有 city_id 和 revenue 两列 # 第一步加随机盐把原来的key拆散 salted_df df.withColumn(salt, (rand() * 10).cast(int)) # 第一步聚合按 (city_id, salt) 粗粒度聚合 partial_df salted_df.groupBy(city_id, salt) \ .agg(_sum(revenue).alias(part_sum)) # 第二步去掉salt再做一次聚合 result_df partial_df.groupBy(city_id) \ .agg(_sum(part_sum).alias(total_revenue))用列的方式加盐有个好处不需要把盐拼到字符串里再拆出来避免了原始key里带分隔符导致截取错误的坑。我早期用“city_id _ 随机数”的做法结果key里居然有下划线第二阶段拆分前缀的时候直接取错了字段数据都算错了一次。改成单独加一列之后就再也没出过这种问题这个细节大家可以放心直接用。盐的粒度决定了打散效果。如果倾斜非常严重建议盐值取大一点比如100甚至500如果倾斜不严重取10到50就够。盐值太小倾斜key还是可能集中在少数分区盐值太大第二阶段聚合的task数量会变多虽然通常还好但也要避免无谓开销。这个方法有一个前提聚合操作本身需要满足交换律和结合律像sum、count、max、min都没问题但如果是求平均值这类需要中间状态的就得先求和和计数最后再相除不能直接对结果求平均。3.3 倾斜key单独处理精准打击两阶段聚合虽然通用但它会把所有key都加一遍盐哪怕不倾斜的key也跟着多了一次shuffle。如果倾斜key很集中通常就是那么一两个更好的做法是“单独拎出来处理”。具体思路是这样的先用采样或者全量统计找出少数热点key把这些key对应的数据过滤出来单独走加盐的两阶段聚合或者别的方案剩余不倾斜的数据走正常的聚合逻辑最后再把两部分结果union起来。这样做的好处是大部分数据不受加盐影响只有少数倾斜key付出了额外的shuffle代价。代码大概长这样from pyspark.sql.functions import col, rand, sum as _sum # 热点key列表这里是假设已经定位到两个 hot_keys [hot_city_a, hot_city_b] # 拆分数据 hot_df df.filter(col(city_id).isin(hot_keys)) normal_df df.filter(~col(city_id).isin(hot_keys)) # 热点部分加盐两阶段聚合 salted_hot hot_df.withColumn(salt, (rand() * 50).cast(int)) partial_hot salted_hot.groupBy(city_id, salt).agg(_sum(revenue).alias(part_sum)) result_hot partial_hot.groupBy(city_id).agg(_sum(part_sum).alias(total_revenue)) # 正常部分直接聚合 result_normal normal_df.groupBy(city_id).agg(_sum(revenue).alias(total_revenue)) # 合并结果 final_result result_hot.union(result_normal)注意一点热点key列表不要写死最好从统计结果里动态取这样即使数据分布发生变化任务也不会因为key列表过期而失效。跳过热点key的判断可以用isin或join来实现热点key多的时候用join更高效。这个方法在“热点key个数少但非常突出”的场景下收益最大比如一个头部商家占全网订单量的一半处理起来立竿见影。但如果倾斜是“长尾式”的比如top100的key都很大单独处理就有点麻烦了这时候更适合用下面要说的通用加盐扩容方案。3.4 map join广播变量能不开shuffle就不开join场景和聚合场景又不太一样。聚合场景是shuffle之后单个key太重join场景很多时候根本可以不shuffle那就是map join也叫broadcast join。原理很简单把小表广播到每个executor的内存里map端直接拿大表的每一行去和小表的本地副本做匹配完全不产生shuffle自然也就没有数据倾斜的问题。触发map join的方式有两种。一种是Spark自动判断当小表大小小于spark.sql.autoBroadcastJoinThreshold默认10MB时Spark SQL优化器会自动把它广播出去。如果小表稍微大一点可以调高阈值但我不建议粗暴地调到几百MB因为广播表要复制到每个executor内存压力是成倍增长的很容易把整个集群的内存占满。另一种方式是显式指定代码里用broadcast函数from pyspark.sql.functions import broadcast result_df big_df.join(broadcast(small_df), key, left)或者在SQL里用hintSELECT /* BROADCAST(b) */ a.*, b.* FROM big_table a LEFT JOIN small_table b ON a.key b.keymap join最爽的地方在于它把问题从“怎么均衡分区”直接变成了“怎么不开shuffle”所以只要小表能放得进内存我优先推荐这个方案。一个小经验是即使小表只有20MB、50MB只要executor内存充足我就手动加broadcast主动帮优化器做决定毕竟自动判断偶尔也会因为表实际大小估算不准而没触发。不过要注意广播会占用每个executor的存储内存如果表到100MB以上就要仔细评估内存够不够了。3.5 随机前缀扩容join场景的通用后手如果两张表都有数据倾斜或者倾斜key太多、没法单独处理而小表又不能直接广播那就要用到“随机前缀扩容”这个通用方案了。思路是用一个随机前缀让大表的倾斜key分散到多个分区同时把小表在相同key上进行N倍扩容让两边能对上号。下面以订单表和用户表按user_id做join为例from pyspark.sql.functions import col, rand, explode, array, lit, concat, cast, sum as _sum N 10 # 大表user_id 加随机前缀比如 123_3 big_salted big_df \ .withColumn(salt, (rand() * N).cast(int)) \ .withColumn(join_key, concat(col(user_id).cast(string), lit(_), col(salt))) # 小表把每条记录扩容成N份前缀从0到N-1都来一遍 small_expanded small_df \ .withColumn(salt_range, explode(array([lit(i) for i in range(N)]))) \ .withColumn(join_key, concat(col(user_id).cast(string), lit(_), col(salt_range))) # 现在按join_key进行等值join result_df big_salted.join(small_expanded, join_key) \ .groupBy(user_id) \ .agg(_sum(amount).alias(total_amount))这里有几个细节要注意。第一扩容倍数N决定了最终小表的数据量如果N太大小表原本很小的优势就没了反而可能因为自身膨胀变成新的倾斜源一般N取10到100比较合适。第二扩容后的join结果可能由于同一个user_id关联出了多行所以需要加一个groupBy来还原业务口径就像上面代码做的那样。第三这个方案会增加shuffle的数据量毕竟小表被放大了N倍所以它是“花钱买均衡”需要看投入产出比。如果只对大表中的热点key加盐对小表也只扩容热点key对应的行性能会更好代码复杂度也会上升。生产上我一般是“倾斜key单独处理”和“随机前缀扩容”结合着用先拆出热点key热点key走扩容逻辑非热点key走普通join最后union。这样既控制了小表膨胀的体积也解决了倾斜问题是join场景比较完善的落地姿势。3.6 Spark 3 AQE系统帮你自动处理倾斜如果你的Spark版本是3.0以上别忘了还有个开了就能白嫖的功能Adaptive Query Execution也就是AQE。它里面带了动态数据倾斜join优化开启之后Spark会在shuffle阶段收集每个分区的真实数据量发现哪些分区明显偏大就自动把它们拆成多个小分区让CPU和内存能够并行处理。开启方式很简单在spark-submit或者代码里设置这几个参数spark-submit \ --conf spark.sql.adaptive.enabledtrue \ --conf spark.sql.adaptive.skewJoin.enabledtrue \ --conf spark.sql.adaptive.skewJoin.skewedPartitionFactor5 \ --conf spark.sql.adaptive.skewJoin.skewedPartitionThresholdInBytes256MB \ --conf spark.sql.adaptive.coalescePartitions.enabledtrue参数含义是需要解释一下的spark.sql.adaptive.enabledAQE总开关Spark 3.0以后默认就是true但最好显式确认一下。spark.sql.adaptive.skewJoin.enabled倾斜join优化开关默认true。spark.sql.adaptive.skewJoin.skewedPartitionFactor当一个分区的大小超过“中位数分区大小的这个倍数”时才认为是倾斜默认是5。spark.sql.adaptive.skewJoin.skewedPartitionThresholdInBytes触发倾斜判断的最小分区大小默认256MB。也就是说某个分区数据量要同时满足“大于所有分区中位数的5倍”和“大于256MB”这两个条件AQE才会动手拆它。这种自动化的处理对join场景特别友好省掉了手写加盐扩容的麻烦。不过AQE目前主要优化join对于group by聚合阶段的倾斜还是得靠前面说的两阶段聚合来处理。我的建议是Spark 3用户先把AQE打开让系统解决一部分问题剩下的再针对性地手动优化。4. 实战复盘一次完整的订单数据倾斜修复记录4.1 作业背景和现场之前有个线上任务每天凌晨跑一次统计前一天的订单数据按商家ID聚合出GMV和订单量结果表再供报表使用。刚开始跑得挺正常大概15分钟能跑完。后来业务上涨某天开始作业变成40多分钟有时候甚至一个多小时。用户开始投诉报表更新不及时。我接到反馈后第一件事就是去Spark UI里看Application的运行情况。打开后一眼就发现Stage 3的运行时间特别长点进Task列表发现一共400个task里有398个在30秒内完成了但有两个task跑了将近40分钟还没结束而且它们的Shuffle Read Size分别高达20GB和18GB。这个差距太夸张了不用统计也知道是数据倾斜。4.2 定位过程接下来我去确认倾斜的key。因为数据是商家维度的我直接在任务代码里临时加了一段采样统计看按商家ID聚合后各个key的数据量分布from pyspark.sql.functions import col, count sample df.sample(withReplacementFalse, fraction0.01, seed123) sample.groupBy(merchant_id).agg(count(*).alias(cnt)) \ .orderBy(col(cnt).desc()) \ .show(20, False)结果显示抽样的数据里有两个merchant_id明显偏大一个占了抽样总量的12%另一个占了8%其余的都在0.x%以下。再一查这两个商家其实就是平台上最大的两家连锁品牌日订单量比其他中小商家高出一两个数量级。这种“少数巨头撑起半边天”的业务结构最容易让按商家粒度聚合的作业倾斜。4.3 落地修复由于聚合场景明确热点key也就两个我先尝试简单地调大spark.sql.shuffle.partitions到2000结果只缩短了5分钟那两个热点key对应的task依然要跑很久。这就说明光加分区数不解决本质问题。然后我改成了“倾斜key单独处理”的两阶段聚合方案。步骤大致是先统计出key分布把top2的热点merchant_id动态提取出来过滤出热点数据加盐聚合非热点数据正常聚合最后union。代码结构和3.3节里的示例基本一致只是盐值取了50因为那两个热点key的数据量非常大需要多打散几个分区。改完再提交作业从40多分钟降到了9分钟左右。最直观的变化是Spark UI里所有task的Shuffle Read Size都在几百MB以内再没有20GB的巨无霸出现。4.4 前后对比指标优化前优化后整体耗时42分钟9分钟最长task耗时接近40分钟约1分半Shuffle Read峰值约20GB约800MB是否OOM偶发无这个案例本质上并不复杂麻烦主要在前期的定位环节花了大半天时间确认热点key真正改代码反而半小时就搞定了。所以还是想说那句话优化数据倾斜先花时间定位别急着改参数。5. 方案怎么选一张对照表和几个避坑心得5.1 不同场景下的方案选择不同场景适合不同方案我整理了一张对照表可以直接对着选场景推荐方案注意点group by / 聚合类热点key不明显两阶段聚合加盐聚合算子要满足交换律和结合律平均值要拆成sumcount聚合类热点key很少且集中倾斜key单独处理热点key列表建议动态获取不要写死小表join大表map join / broadcast广播阈值量力而行别盲目调大导致OOMjoin类两张表都有热点随机前缀扩容扩容倍数别太大控制小表膨胀体积join类有少量热点key倾斜key单独处理普通join union代码稍复杂但性价比最高null值或魔法值造成的倾斜给null加随机前缀或过滤注意别改变业务统计口径Spark 3以上且是join场景开启AQE自动倾斜优化省心但聚合倾斜仍需手动处理输入数据本身大小不均调整分区策略或读取参数和shuffle倾斜不是一回事别混着处理选型的时候既要看场景也要看团队对代码的维护成本。两阶段聚合相对通用代码也简单我会优先给新同学推荐倾斜key单独处理虽然性能最好但dp拆分流和union逻辑会增加维护负担热点key变化频繁时还要保证动态更新需要结合实际情况权衡。5.2 我踩过的坑和一些实操细节最后分享几个我踩过多次的坑也算不上高深但每一个都真实花过时间去排。第一个坑是“调大shuffle.partitions就能解决”实际上对于极端倾斜的key分区数再大也没用因为数据是按key聚合后再分区的key本身的重量没有变。第二个坑是两阶段聚合时把盐拼到key字符串上结果key里本来就有分隔符第二阶段还原key时取到了错误的值导致数据算错。后来我改成单独加salt列彻底规避了这个问题。这里提醒一句在做任何“加盐-还原”的操作时务必先确认原始key的字符集别想当然。第三个坑是map join时把小表广播阈值调得过高。我有一次为了省事把autoBroadcastJoinThreshold调到了500MB结果每个executor多了一大块常驻内存集群压力骤增虽然join不shuffle了但整体效率反而更差。广播表是内存换shuffle的思路阈值要克制。第四个坑发生在AQE的推广阶段。有些同事以为开了spark.sql.adaptive.skewJoin.enabled就万事大吉结果发现group by场景依然慢。这里要明确AQE的自动倾斜优化目前主要针对shuffle join聚合类的倾斜还是要用两阶段聚合。第五个坑是sql没有处理null值。统计里经常多个表的key为空如果不处理null会全部hash到同一个分区。最简单的做法是给null一个随机后缀让它们分散到不同分区比如when(col(key).isNull(), concat(lit(null_), (rand()*100).cast(int))).otherwise(col(key))。我个人的体会是数据倾斜优化没有“一招鲜”的银弹每类方案都有自己的适用边界核心是先定位准确再选择题型。只要你肯在Spark UI和采样统计上花一点时间绝大多数倾斜问题都能在半小时内找到方向。最后再分享一个日常习惯每次上线新的shuffle类作业前我会先跑一份数据量占比统计模块化地评估key的分布是否健康这样很多倾斜问题在代码评审阶段就被提前拦下来了。
返回列表