如何识别并解决Spark中的倾斜分区问题?(多Join任务场景)
解决Spark SQL多表Join的数据倾斜问题
一、识别倾斜对应的Join环节
- 查看Spark UI的Stage统计:定位Shuffle Read/Write量远超平均值的Task,记录对应的Shuffle ID,再关联到SQL执行计划中的Join节点,确定是哪次Join导致的倾斜。
- 生成详细执行计划:执行
EXPLAIN EXTENDED查看查询的物理执行计划,梳理所有Join操作的依赖关系,结合Shuffle Key定位倾斜源头。 - 采样统计倾斜Key:对疑似倾斜的表执行采样查询,比如:
找出出现频次极高的Key,确认这些Key是否对应到倾斜的分区。SELECT join_key, COUNT(*) AS cnt FROM large_table GROUP BY join_key ORDER BY cnt DESC LIMIT 10
二、拆分/消除倾斜分区的方法
1. 拆分倾斜Key单独处理
针对高频倾斜Key,将其从主查询中分离,采用不同Join策略后合并结果:
- 先通过采样获取Top N倾斜Key;
- 对大表过滤出这些Key的数据,与小表(如3-4GiB的表)做广播Join(
/*+ BROADCAST(small_table) */),避免Shuffle; - 主查询处理非倾斜Key的数据,使用常规Join;
- 最后用
UNION ALL合并两部分结果。
2. 加盐(Salting)打散数据
通过给倾斜Key添加随机后缀,将单个大拆分为多个小分区:
- 对大表的倾斜Key添加随机后缀:
concat(join_key, '_', cast(rand() * 100 as int))(100可根据倾斜程度调整); - 对关联的小表,将对应Key复制100份,每份添加不同的后缀(0到99);
- 使用加盐后的Key执行Join,完成后去掉后缀合并数据。
3. 预分区优化
- 对1-2TB的大表提前执行
repartition(2000, join_key)(分区数根据集群规模调整,建议为核心数的2-3倍),让数据按Join Key均匀分布后再执行查询。 - 如果是常用表,可创建Bucketed表,按Join Key分桶,Bucket数设为集群核心数的倍数,后续Join时无需全量Shuffle,直接按Bucket关联。
4. 过滤无效倾斜Key
若倾斜Key是NULL、空字符串等无效值,可提前过滤:
- 在查询中添加
WHERE join_key IS NOT NULL AND join_key != ''; - 若需保留这类数据,可单独处理(如不关联直接输出,或与小表做左Join时单独处理)。
内容的提问来源于stack exchange,提问作者Asif Khan
相关产品推荐
相关产品推荐

