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类中。

策略选择器的工作流程可以概括为:

  1. 解析JOIN类型和条件(等值/非等值连接)
  2. 检查是否存在显式的JOIN提示(hints)
  3. 评估参与JOIN的表大小和分布情况
  4. 根据配置参数和启发式规则选择最优策略

让我们看一个简化的决策流程图:

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策略选择遵循明确的优先级链:

  1. Hint优先:显式指定的JOIN提示具有最高优先级
  2. 大小评估:基于表大小和广播阈值判断
  3. 配置参数preferSortMergeJoin等开关控制
  4. 默认回退:最终选择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

这个计划告诉我们:

  1. 选择了Sort Merge Join(默认策略)
  2. 两侧都进行了shuffle(Exchange节点)
  3. 都进行了排序(Sort节点)
  4. 分区数使用默认的200

如果发现非预期策略,可以检查:

  • 表统计信息是否准确(ANALYZE TABLE)
  • 是否存在数据倾斜问题
  • 配置参数是否被覆盖

理解这些底层机制,我们就能更准确地预测和优化Spark SQL的JOIN性能,而不再被"为什么选择这个JOIN策略"的问题所困扰。

更多推荐