为什么Spark中数据倾斜会引发问题?左外关联场景执行原理咨询
Spark 大Key左连小表卡顿原因及执行流程拆解
核心认知误区澄清
Spark 的分区自动拆分规则不会作用于shuffle阶段的同key数据,shuffle的重分区逻辑严格按照key的哈希值映射分区编号,同一key的所有记录一定会被分配到同一个分区处理,不会自动拆分。
整个Join执行流程的卡顿点定位
你当前使用常规shuffle join模式执行left outer join时,流程和故障点如下:
- 阶段1:大表、小表分别读取生成分区,各自完成map阶段预处理,生成带key的中间数据
- 阶段2:shuffle分桶环节
- 所有key为
true的5000万条记录全部哈希映射到同一个分区编号,无论数据量多大都不会拆分 - 所有key为
false的5万条记录哈希映射到另一个分区 - 极小右表的对应key数据也会按哈希映射到上述两个分区
- 所有key为
- 阶段3:Join计算环节
单个task需要处理true对应分区的5000万条记录+右表对应key的所有数据,该task的计算耗时是其他task的数百倍,整个任务必须等待最慢的这个task执行完成才能结束,直观表现就是任务卡顿。
为什么你预期的分区拆分没有生效
你提到的“分区负载过高自动拆分”仅在两个场景下生效:
- 数据源读取阶段:读取大文件时如果单个分片超过阈值,会自动拆分多个小分片读取,该阶段还未做key分组,不会限制同一key的分布
- AQE动态合并分区:仅会合并数据量过小的分区,默认不会拆分单个key对应的大分区——如果强制拆分,会破坏同一key必须在同一分区处理的join计算逻辑
如果开启了Spark 3.x的AQE倾斜优化,框架会自动检测大key并拆分,但默认倾斜阈值较高,如果你的大key对应数据量未达到阈值也不会触发,需要手动调整spark.sql.adaptive.skewJoin.skewedPartitionThresholdInBytes参数适配场景。
该场景最优解决方案
因为你是和极小数据集做左连接,直接用广播Join(Map Join)避免shuffle即可彻底解决倾斜问题:
- 提前将小表广播到所有Executor节点:
val broadcastSmallTable = spark.sparkContext.broadcast(smallTable.collectAsMap()) - 大表在map阶段直接读取本地的广播小表做关联,不需要走shuffle流程,自然不会出现同key聚集的倾斜问题。
如果是小表太大无法广播的场景,可以给倾斜的大key加上随机前缀拆分后再做关联。
内容的提问来源于stack exchange,提问作者Kate
相关产品推荐
相关产品推荐

