Spark数据倾斜实战:从定位到优化的六个解决方案
1. 一个让Spark任务卡死6小时的真实问题前阵子接手了一个数据清洗任务跑的是公司内部一个日级ETL逻辑不算复杂两张表关联后做聚合再写回Hive分区。代码写完测试环境跑得飞快可一到生产任务就像被什么东西卡住一样Stages一直卡在99%不动Executor日志里全是GC告警和OOM。刚开始以为是集群资源不够连着加了两轮资源结果毫无变化任务照样超时失败。后来点开Spark UI一看某个Stage的Shuffle Read行的数据量分布极其离谱24个Task里有1个Task读取了2.1GB数据其余23个Task每个只读了不到50MB。那一刻基本可以确诊任务撞上了Spark开发里最经典也最难缠的性能杀手——数据倾斜Data Skew。这不是我第一次遇到数据倾斜也不会是最后一次。从实习时被倾斜问题折磨到凌晨两点到后来在线上一眼就能从Spark UI里定位倾斜这中间踩过的坑、试过的方法、总结出的经验值得好好写一篇实战记录。本文适合所有正在用Spark做离线计算、实时计算或者准备Spark面试的读者尤其适合那些已经被“某个Task跑得特别慢”折磨过的朋友。2. 数据倾斜的本质和识别方法2.1 数据倾斜到底是怎么发生的数据倾斜官方定义叫“数据分布不均”本质是Spark在做Shuffle、Join、GroupBy等操作时某个或某几个Key对应的数据量远超其他Key。打个比方一个班级里有40个学生老师让大家分组搬运同样大小的纸箱大部分组每人搬2箱但有个组被分配了100箱这个组就是数据倾斜里的那个“倾斜Key”。在Spark的物理执行计划里数据从上游Task流向下游Task时如果某个Key的Hash值让大部分数据都落到同一个Reducer上这个Reducer的执行时间就会远大于其他Reducer。最后整个Stage的完成时间等于最慢的Task时间于是任务整体被拖垮。倾斜产生的原因多种多样最常见的是业务Key分布本身就不均匀比如某个热门商家的订单量是普通商家的几百倍其次是数据清洗时漏了空值处理导致大量null被聚合成同一个Key还有一种情况是Join时小表里的关联Key重复度过高。我在实际排查中发现很多倾斜问题都是几个因素叠加产生的单一因素的倾斜反而好解决。2.2 如何在Spark UI中快速定位倾斜Task定位数据倾斜最直接的手段就是打开Spark UI看Executors标签页和Stages标签页。如果你发现某个Stage里有极少数Task的Shuffle Read或Input数据量是其他Task的几十倍同时这些Task的Duration也异常长基本就可以判定这个Stage存在数据倾斜。还有几个辅助信号值得关注Spark UI中某个Task的GC Time特别长说明这个Task的内存压力已经很大Event Timeline里某个Task的进度条在99%停滞很久说明Reducer在拉取最后一批数据时被卡住日志中出现频繁的Shuffle FetchFailedException或者BlockManager报错往往也和倾斜Task的资源耗尽有关。我之前有个偏好就是无论任务报不报错只要有性能问题先看Spark UI的“SQL”标签页。在SQL执行计划图里每个Stage下面会标注Shuffle Read和Shuffle Write的数据量一眼就能扫出哪个Stage的倾斜嫌疑最大。这一步做完再往下定位具体是哪个Key用一句group by count就能查出来。定位到具体倾斜Key之后需要判断倾斜的“粒度”。是整个Key空间分布就不均衡还是只有一两个超级Key导致的局部倾斜这两种情况的处理策略完全不同。判断方法很简单SQL里执行一下select key, count(*) as cnt from table group by key order by cnt desc limit 20;如果前几个Key的count远远大于其他Key就是典型的“热点Key”问题。如果所有Key的count都很高但个别Key的值域特别宽那就要从业务逻辑上分析是否可以把宽Key拆分。3. 六个久经考验的解决方案3.1 方案一增加Reduce并行度低成本缓解倾斜调整并行度是最廉价也最容易被想到的办法。通过spark.sql.shuffle.partitions或spark.default.parallelism参数把Shuffle的下游并行度调大。比如说原来200个Reducer每个Reducer要处理500万条数据调成2000个Reducer后每个Reducer只处理50万条单个Task的压力自然降下来。但这种方法有个致命缺陷它只能“摊薄”压力不能“消除”热点。如果倾斜Key的绝对数据量特别大并行度调到一万也没用因为不管怎么分那个倾斜Key的数据都会被Hash到同一个Reducer上。所以这个方法只适合倾斜程度不严重的场景比如数据倍数在几倍以内或者集群资源比较充裕的情况。在实际操作中我通常会在SQL脚本里写成这样set spark.sql.shuffle.partitions 500;或者用SparkSession设置spark.conf.set(spark.sql.shuffle.partitions, 500)调完以后记得重新看一眼Spark UI确认数据分配是否均匀。如果倾斜Task的瓶颈在数据计算本身单纯调大并行度可能带来副作用Task数量变多后Shuffle的中间文件数量也增加了会提升网络和磁盘IO的负担甚至反而变慢。这一点很多新手容易忽略。3.2 方案二过滤异常Key和空值收尾很多时候倾斜的根源是数据本身的问题而不是计算框架的问题。我在处理用户行为日志时遇到过一种情况埋点数据采集异常导致约30%的日志event_id为空这些空值全部落到了同一个NullKey上。下游做GroupBy统计时NullKey所在Task要处理的数据量是其他Task的20倍整个任务动弹不得。这个时候最直接的办法就是在计算前先过滤掉或者单独处理异常Key。比如对null值的处理可以提前过滤select * from table where event_id is not null;或者保留它们但单独统计不和正常Key混合-- 正常Key的统计 select event_id, count(*) from table where event_id is not null group by event_id; -- 异常数据单独统计 select null_key as event_id, count(*) from table where event_id is null;还有一种少见但更容易踩雷的情况业务数据里混入了脏数据比如某个Web攻击脚本在日志里伪造了大量自定义User-Agent导致这个UA的count值异常高。这类数据如果不过滤它会始终霸占某个Task不管你怎么调整资源都白搭。所以遇到倾斜我的第一反应不是去改代码而是先查源数据质量。去重、过滤、清洗这些步骤虽然看起来“脏活累活”但往往能解决最棘手的倾斜问题。3.3 方案三加盐随机前缀打破热点Key如果热点Key无法过滤业务上又必须统计它们加盐Salting是目前最经典的解决方案。核心思想很简单给倾斜的Key加上一个随机前缀让原本集中在同一个Reducer的数据被分散到多个Reducer上处理处理完后再去掉前缀聚合。拿上面的订单统计举例假设订单表里商家id为10086的订单量特别大而其他商家id分布正常。处理方式是分两步第一步给倾斜Key加随机前缀聚合一次。这一步把几百万条数据均匀打散到多个Task上-- 第一步给每个商家id加一个1到10之间的随机前缀 select concat_ws(_, prefix_ floor(rand() * 10), cast(merchant_id as string)) as salted_key, count(*) as cnt from order_table group by salted_key;第二步把加了前缀的Key去掉再次聚合-- 第二步去掉随机前缀得到最终统计结果 select split(salted_key, _)[1] as merchant_id, sum(cnt) as total_order_cnt from ( select concat_ws(_, prefix_ floor(rand() * 10), cast(merchant_id as string)) as salted_key, count(*) as cnt from order_table group by salted_key ) t group by merchant_id;这里需要特别注意加盐的粒度一定要控制好。盐值太少打散效果不明显盐值太多两阶段聚合的中间结果会膨胀反而增加网络IO和内存压力。根据我的经验盐值数量建议取倾斜Task数量的2到3倍或者按倾斜数据量除以期望的单个Task数据量来计算。这个方案对Group聚合类的倾斜非常有效但Join场景要用另一种方式。3.4 方案四Join场景下的广播变量优化如果是两张表的Join导致倾斜情况比GroupBy要复杂一些。GroupBy的倾斜只需要处理一个维度Join要考虑两边的数据分布。当大表和小表Join时首选方案是把小表广播Broadcast出去这样每个Executor都保存一份小表的副本大表的数据可以在本地完成匹配完全绕过Shuffle也就没有Shuffle倾斜的问题。Spark 2.0之后默认的广播阈值是10MB超过这个阈值就需要手动开启spark.conf.set(spark.sql.autoBroadcastJoinThreshold, 104857600) # 100MB这里补充一句广播模式下小表数据会复制到每个Executor的内存里。如果小表太大比如超过200MB广播本身就会带来内存压力这时候需要评估一下收益和成本。我通常建议阈值设到100MB以内太大容易把Executor的内存撑爆。如果两张大表Join都有热点Key那就没法单纯靠广播解决了。比较实用的做法是对热点Key单独处理。先把大表A和B中倾斜严重的Key挑出来对这些数据做加盐处理再把原本的表数据按热点Key和非热点Key拆开非热点部分直接Join热点部分加盐后Join最后union结果。这相当于用代码的复杂度换取执行性能也是大型Join场景下最常用的手段。3.5 方案五两阶段聚合解决GroupBy倾斜的最优解两阶段聚合其实是加盐方案的完整版核心思路是“局部聚合全局聚合”。先把数据按照加盐后的Key做一次聚合此时相同业务Key的数据被打散到各个Task每个Task只处理一部分数据然后去掉盐值再次聚合得到最终结果。这个方法特别适合业务上必须按某个维度做精确去重统计的场景。比如要统计每个用户的浏览PV但某个用户的浏览记录特别多如果直接count这个用户所在的Task就会成为热点。我在实际项目里遇到过更极端的情况某个头部用户的浏览记录占了全表20%的数据量加10个盐值都不够。后来我调整了策略加盐时用了动态盐值根据热点Key的数据量在运行时决定盐值数量。这种方法需要用一个小作业先统计热点Key的数据量然后动态构造加盐逻辑。虽然代码变复杂了一些但效果很明显把原来要跑40分钟的任务优化到了8分钟。需要注意的是两阶段聚合只对“可加性”的聚合函数有效。count、sum、max、min这些都没有问题但如果是count(distinct)或者走近似去重算法就要小心了。两阶段聚合会把去重逻辑打散聚合结果可能不对。这种情况建议改用精确去重函数或者走窗口函数另想办法。3.6 方案六RDD级别手动分区最底层的兜底方案如果以上方案都不适用比如你要处理的数据本身就不是DataFrame而是RDD算子比较复杂倾斜点不好从SQL层面规避那只能回到RDD层面手动设置分区器。RDD的Shuffle可以自定义Partitioner实现数据倾斜Key的多路分发。比如重写getPartition方法对热点Key单独指定分区让它们分散到多个分区去class SkewPartitioner(partitions: Int) extends Partitioner { override def numPartitions: Int partitions override def getPartition(key: Any): Int { val k key.asInstanceOf[String] if (k hot_key) { // 热点Key随机分发到前10个分区 Random.nextInt(10) } else { // 其他Key正常Hash (k.hashCode Int.MaxValue) % partitions } } }用的时候在reduceByKey或者groupByKey的时候传入这个分区器rdd.map(record (record.merchantId, record)) .reduceByKey(new SkewPartitioner(100), (a, b) merge(a, b))这个方法的好处是灵活度高可以精确控制每个Key的去向。坏处是代码侵入性强维护成本高而且需要你对Spark的分区机制有足够深的理解。如果对源码不熟调起来容易出各种意想不到的问题非必要不建议直接用。4. 三招组合拳生产环境里的实战定位流程很多朋友看完上面的方案会有一个困惑我知道这些方法了但真到了生产环境到底应该怎么系统地定位和选择用哪个方案这里分享一套我在生产环境里反复验证过的组合拳流程。第一招先看Spark UI锁定Stage和Task。打开应用的Stages页面找到耗时最长的Stage点进去看每个Task的数据量分布。配一张柱状图或者表重点看是否有Task的数据量是Median值的好几倍。一般3倍以上就值得警惕10倍以上基本就是要处理的大问题了。第二招查SQL逻辑分析倾斜Key的业务含义。拿到Stage对应的执行计划定位到具体的SQL操作。用explain命令看一下物理计划明确这个Stage是Join、Aggregation还是Shuffle产生的。然后回到数据源执行上面提到的配额查询找出具体的倾斜Key。这一步需要结合业务理解比如某个活动在搞促销那么活动id相关数据大概率会爆量节假日日志里某个渠道id特别多等等。第三招依据倾斜类型选择方案优先零侵入的手段。如果只是并行度不够果断调大shuffle.partitions如果是空值和脏数据先做过滤如果是业务的天然热点Key判断数据规模后决定用广播变量还是加盐。通常我按这个优先级执行先参数调优再过滤数据再广播小表最后加盐和两阶段聚合。按照这个顺序能不动代码解决的问题尽量不动代码毕竟改线上SQL是有风险的。这套流程我把它做成了表格方便日常对照倾斜现象可能原因首选方案替代方案某个Task数据量暴增并行度不足调大shuffle.partitions加盐Null/空值导致的倾斜数据质量问题过滤或单独处理Null加盐后聚合单个Key过热业务热点Key加盐/两阶段聚合动态盐值大表Join小表倾斜Join Key热点广播小表热点Key单独处理多Key值域不均数据模型问题重新设计Key分拆参数维度复杂算子倾斜RDD层面热点自定义Partitioner改SQL/过滤数据这个流程用熟了以后定位倾斜问题的时间能从小时级压缩到分钟级处理方案也能一步到位不用反复试错。5. 高频面试考点和数据倾斜面试题精讲Spark数据倾斜是面试中最高频的性能优化考点几乎每一轮技术面试都会涉及。面试官通常从“你遇到过数据倾斜吗”开始然后层层深入如何定位、如何解决、为什么这样能解决以及不同方案在不同场景下的优劣。我梳理了几个面试中最高频的题目和答题思路按照面试官追问的逻辑整理如下高频问题一Spark任务卡在99%如何判断是否存在数据倾斜答题要点先说Spark UI的观察方法Stage中个别Task执行时间过长、Shuffle Read数据量分布极不均衡再补充GC时间异常、FetchFailedException等辅助信号最后强调要结合两个维度的证据综合判断而不是单看Task数量。高频问题二请说出至少三种解决数据倾斜的方法。答题要点过滤异常Key、增加并行度、广播小表、加盐、两阶段聚合、自定义Partitioner。答的时候每说一种都要描述适用场景和底层原理显得有深度。面试官更在意你是不是真的理解为什么要这么做。高频问题三加盐和两阶段聚合有什么区别答题要点加盐是一种数据重分布手段主要解决Join中的热点Key问题两阶段聚合是加盐在GroupBy场景的完整实现先局部聚合再全局聚合通过两次Shuffle彻底打散热点。核心区别是前者只改变数据的物理位置后者改变了聚合的执行逻辑。高频问题四广播变量为什么能解决数据倾斜答题要点广播变量把Shuffle变成了“读本地副本”小表数据分发到每个Executor大表在每个Executor本地就能完成关联绕过Shuffle这个倾斜的产生源头。同时要提到broadcast Join只有在表大小小于阈值的情况下才适合超过阈值要评估内存开销。高频问题五如果你发现倾斜是由Null值导致的怎么处理答题要点先排查数据质量问题过滤Null或单独统计如果Null在业务上有含义不能随便过滤就给Null加随机盐值让它分散到多个Reducer。同时要补充一个实际经历很多倾斜问题的源头是上游数据管道解析错误检查源头比在Spark层面修补更有效。这些内容基本覆盖了我面试时遇到的考题范围。Spark面试题里也有大量数据倾斜相关的变形题万变不离其宗你只要深刻理解了上面的原理和方案答题时自然会言之有物。6. 踩坑实录和避坑指南6.1 加盐后结果不对原来踩了distinct的坑有一次我做UV统计用户表关联行为表然后count(distinct user_id)。为了优化性能我对用户id加了盐打算局部去重后汇总。结果跑完发现UV数比真实值大了好几倍排查了半天才发现两阶段聚合时第一阶段的局部count(distinct)丢掉了用户id本身的信息第二阶段再sum的时候重复的用户被算了多遍。解决这个问题有两种思路一种是放弃两阶段聚合改用加盐Join的方式先给热点用户id加盐打散Join完成后做全局精确去重另一种是使用近似去重算法比如approx_count_distinct牺牲一点精度换取性能。但要记住业务上对UV的准确性要求很高时近似算法也不能用。从这里我吸取了一条经验动手优化之前先确认聚合函数是否具备“可加性”不确定的话先用小数据集做一遍正确性验证再放到生产上跑。6.2 广播变量爆炸Executor内存被撑爆另一个印象深刻的坑是接了一个明细表Join维度表的场景明细表2亿条维度表80MB看起来广播方案很合适。结果线上任务频繁Executor LostOOM日志一大片。后来排查发现维度表实际占用内存远不止80MB行式存储转列式存储后膨胀了好几倍广播到每个Executor后和Spark执行本身需要的内存叠加直接把Executor整垮了。这个教训让我意识到一点广播阈值不是只看表存储大小还要看表的行数。行数越多Java对象的Header、指针等额外开销越大实际内存占用可能达到数据文件大小的3-5倍。现在我的经验是行数超过300万或者表大小超过50MB就不太建议用广播了宁愿走SortMergeJoin或者热点Key拆分。6.3 动态盐值自适应解决固定盐值不均衡固定盐值有个天然问题如果热点Key的数据量特别大而盐值数量设置得不够每个盐前缀下的数据量依然不均衡。一次线上促销活动中某个商家的订单量涨了20倍我原先给他加的10个盐值依然不够任务在促销期间照旧倾斜。后来我改成动态盐值方案先采样统计热点Key的数据量再根据目标Task处理能力算出需要拆分到多少分区然后给这个Key动态分配盐值数量。代码逻辑不复杂但需要先跑一个采样任务。虽然多花了1到2分钟但倾斜问题在全链路得到了根治。这个方法我现在一直在用尤其是在可预知的业务高峰比如大促、秒杀到来之前提前跑一次数据探查动态调整盐值策略。7. 一个完整案例的优化全过程最后用一个真实项目完整串一遍数据倾斜的定位与调优流程帮大家把前面的方法论落地。项目背景一个每日用户行为分析任务输入表包含用户浏览日志约5亿条/日和用户维度表约800万条。业务需求是统计每个用户在当天每个频道的浏览次数、浏览时长、独立访问页面数。上线初期任务运行时间约为55分钟勉强能接受。但一到周末流量高峰期运行时间飙升到2小时以上经常超时失败。我接手后按照上述流程做了系统排查。第一步Spark UI定位。在Stages页面发现第三个StageJoin后的聚合Stage存在明显倾斜1000个Task中最慢的Task跑了45分钟而中位Task只需要3分钟。第二步查看执行计划发现这个Stage包含了一次大表与小表的JoinJoin Key是user_id。第三步倾斜原因分析。通过count查询发现少数头部用户top 1%占了整体日志量的15%。他们的user_id本身就是热点Key。同时部分用户维度表里不存在该user_id外键失效导致大量日志在Join后落到同一个默认分区。明确了问题后我的调优方案如下第一轮连续调参过滤无效Key。把spark.sql.shuffle.partitions从默认200调到1000同时添加过滤条件把user_id为null或不在维度表里的日志单独剥离出来写入脏数据表不再参与主链路的聚合。这一轮完成后任务从2小时降到了65分钟但距离目标还是不够。第二轮针对头部热点用户做加盐Join。从日志表中提取top 100热点user_id给这部分数据加上随机盐值同时维度表针对这100个用户构建成小表并在广播出去之后同样加盐。两边的Join Key统一为“盐值_user_id”从而把热点用户分散到多个Task上。非热点用户走常规Join。这一轮优化后任务从65分钟降到了28分钟。第三轮SQL层面微调。把原先的count(distinct page_id)改为先group by user_id, page_id再去重实际上就是把窗口逻辑换成聚合逻辑减少单个Task的计算压力。任务最终稳定在22分钟左右。这个案例的优化过程中耗时最长的是定位阶段真正动手改代码只用了半天。这再次印证了开头的观点数据倾斜问题的核心难点不在“怎么解”而在“怎么找”。一旦准确定位到倾斜的Stage和Key解决方案基本都是套路。希望这篇文章能把“找倾斜”的方法论和“解倾斜”的武器库都完整交给你以后线上再遇到类似问题能少走很多弯路。