深入Spark SQL内核:从源码片段看JOIN策略选择器的‘内心戏’(以3.3.1版本为例)
Spark SQL JOIN策略选择器的深度解析:从源码视角看优化器决策逻辑
在分布式计算领域,JOIN操作始终是性能优化的关键战场。当我们查看Spark SQL的执行计划时,经常会好奇:为什么Spark在这个查询中选择了Broadcast Hash Join,而在另一个看似相似的查询中却使用了Sort Merge Join?要真正理解这些选择背后的逻辑,我们需要深入Spark Catalyst优化器的核心——JoinSelection策略选择器。
1. JOIN策略选择器的架构设计
Spark 3.3.1版本中的JoinSelection是一个典型的策略模式实现,它继承自Strategy特质,负责将逻辑计划转换为物理计划。这个转换过程本质上是一个多条件决策系统,其核心逻辑位于org.apache.spark.sql.execution.SparkStrategies.JoinSelection类中。
策略选择器的工作流程可以概括为:
- 解析JOIN类型和条件(等值/非等值连接)
- 检查是否存在显式的JOIN提示(hints)
- 评估参与JOIN的表大小和分布情况
- 根据配置参数和启发式规则选择最优策略
让我们看一个简化的决策流程图:
def apply(plan: LogicalPlan): Seq[SparkPlan] = plan match {
case j @ ExtractEquiJoinKeys(joinType, leftKeys, rightKeys, nonEquiCond, left, right, hint) =>
// 等值连接的处理逻辑
createBroadcastHashJoin(true)
.orElse(createSortMergeJoinIfHinted(hint))
.orElse(createShuffleHashJoinIfHinted(hint))
.getOrElse(createJoinWithoutHint())
case _ => // 非等值连接的处理逻辑
}
2. 五种JOIN策略的源码级实现
2.1 Broadcast Hash Join的触发条件
Broadcast Hash Join是效率最高的JOIN策略,其触发条件在createBroadcastHashJoin方法中明确体现:
def createBroadcastHashJoin(onlyLookingAtHint: Boolean): Option[Seq[SparkPlan]] = {
getBroadcastBuildSide(left, right, joinType, hint, onlyLookingAtHint, conf).map {
buildSide => Seq(joins.BroadcastHashJoinExec(
leftKeys, rightKeys, joinType, buildSide, nonEquiCond,
planLater(left), planLater(right)))
}
}
关键判断逻辑包括:
- 表大小是否小于
spark.sql.autoBroadcastJoinThreshold(默认10MB) - JOIN类型是否支持(不支持Full Outer Join)
- 是否显式指定了BROADCAST提示
实际案例:当维表大小<10MB且为等值内连接时,即使没有hint也会自动选择广播。
2.2 Sort Merge Join的默认选择机制
作为Spark默认的JOIN策略(当spark.sql.join.preferSortMergeJoin=true),其实现核心是:
def createSortMergeJoin(): Option[Seq[SparkPlan]] = {
if (RowOrdering.isOrderable(leftKeys)) {
Some(Seq(joins.SortMergeJoinExec(
leftKeys, rightKeys, joinType, nonEquiCond,
planLater(left), planLater(right))))
} else None
}
关键特性:
- 要求JOIN key可排序(
RowOrdering.isOrderable) - 需要预先对两侧数据集进行shuffle和排序
- 适合大表与大表的连接
2.3 Shuffle Hash Join的权衡取舍
当Broadcast不可行但表仍然较小时可能选择的策略:
def createShuffleHashJoin(onlyLookingAtHint: Boolean): Option[Seq[SparkPlan]] = {
getShuffleHashJoinBuildSide(left, right, joinType, hint, onlyLookingAtHint, conf).map {
buildSide => Seq(joins.ShuffledHashJoinExec(
leftKeys, rightKeys, joinType, buildSide,
nonEquiCond, planLater(left), planLater(right)))
}
}
触发条件包括:
spark.sql.join.preferSortMergeJoin=false- 一侧表足够小以构建内存哈希表
- 没有禁用该策略的配置
2.4 非常用策略的兜底方案
对于Cartesian Product和Broadcast Nested Loop Join,它们通常作为最后的选择:
def createJoinWithoutHint(): Seq[SparkPlan] = {
createBroadcastHashJoin(false)
.orElse(if (!conf.preferSortMergeJoin) createShuffleHashJoin(false) else None)
.orElse(createSortMergeJoin())
.orElse(createCartesianProduct())
.getOrElse(/* BroadcastNestedLoopJoin */)
}
3. 决策优先级与参数调优
Spark的JOIN策略选择遵循明确的优先级链:
- Hint优先:显式指定的JOIN提示具有最高优先级
- 大小评估:基于表大小和广播阈值判断
- 配置参数:
preferSortMergeJoin等开关控制 - 默认回退:最终选择Broadcast Nested Loop Join
关键配置参数及其影响:
| 参数 | 默认值 | 影响范围 |
|---|---|---|
| spark.sql.autoBroadcastJoinThreshold | 10MB | 广播JOIN的触发阈值 |
| spark.sql.join.preferSortMergeJoin | true | 是否优先选择Sort Merge Join |
| spark.sql.shuffle.partitions | 200 | 影响Shuffle操作的并行度 |
调优建议:
- 对于星型模型,适当增加
autoBroadcastJoinThreshold - 在内存充足时,可尝试关闭
preferSortMergeJoin以使用Shuffle Hash Join - 监控JOIN策略选择情况,通过EXPLAIN验证实际选择
4. 从执行计划反推优化器决策
通过分析物理执行计划,我们可以逆向理解优化器的选择逻辑。典型计划示例:
== Physical Plan ==
*(5) SortMergeJoin [id#10L], [id#20L], Inner
:- *(2) Sort [id#10L ASC NULLS FIRST], false, 0
: +- Exchange hashpartitioning(id#10L, 200)
: +- *(1) Filter isnotnull(id#10L)
: +- Scan parquet table1
+- *(4) Sort [id#20L ASC NULLS FIRST], false, 0
+- Exchange hashpartitioning(id#20L, 200)
+- *(3) Filter isnotnull(id#20L)
+- Scan parquet table2
这个计划告诉我们:
- 选择了Sort Merge Join(默认策略)
- 两侧都进行了shuffle(Exchange节点)
- 都进行了排序(Sort节点)
- 分区数使用默认的200
如果发现非预期策略,可以检查:
- 表统计信息是否准确(ANALYZE TABLE)
- 是否存在数据倾斜问题
- 配置参数是否被覆盖
理解这些底层机制,我们就能更准确地预测和优化Spark SQL的JOIN性能,而不再被"为什么选择这个JOIN策略"的问题所困扰。
更多推荐
所有评论(0)