You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

Spark SQL分桶Join优化未生效,请求排查原因

关于Spark 3.x分桶表Join未触发分桶优化的问题

问题背景

使用Spark 3.x编写测试代码学习分桶特性,两张表均按b列分2桶存储,但执行Join后物理计划显示采用BroadcastHashJoin(存在广播Shuffle),未触发预期的分桶Join优化(无需Shuffle)。

测试代码

test("bucket join 1") {
    val spark = SparkSession.builder().master("local").enableHiveSupport().appName("test join 1").config("spark.sql.codegen.wholeStage", "false").getOrCreate()
    import spark.implicits._
    val data1 = (0 to 100).map {
      i => (i, ('A' + i % 6).asInstanceOf[Char].toString)
    }

    val t1 = "t_" + System.currentTimeMillis()
    data1.toDF("a", "b").write.bucketBy(2, "b").saveAsTable(t1)

    val data2 = (0 to 5).map {
      i => (('A' + i % 6).asInstanceOf[Char].toString, ('A' + i % 6).asInstanceOf[Char].toString)
    }

    val t2 = "t_" + System.currentTimeMillis()
    data2.toDF("a", "b").write.bucketBy(2, "b").saveAsTable(t2)

    val df = spark.sql(
      s"""
      select  t1.a ,t1.b,t2.a, t2.b from $t1 t1 join $t2 t2 on t1.b = t2.b
      """.stripMargin(' '))

    df.explain(true)
    df.show(truncate = false)
    spark.sql(s"describe extended $t1 ").show(truncate = false)
    spark.sql(s"describe extended $t2 ").show(truncate = false)
}

执行后的物理计划

BroadcastHashJoin [b#27], [b#29], Inner, BuildRight
:- Project [a#26, b#27]
:  +- Filter isnotnull(b#27)
:     +- FileScan parquet default.t_1700986792232[a#26,b#27] Batched: false, DataFilters: [isnotnull(b#27)], Format: Parquet, Location: InMemoryFileIndex[file:/D://spark-warehouse/t_17009867..., PartitionFilters: [], PushedFilters: [IsNotNull(b)], ReadSchema: struct<a:int,b:string>, SelectedBucketsCount: 2 out of 2
+- BroadcastExchange HashedRelationBroadcastMode(List(input[1, string, true])), [id=#35]
   +- Project [a#28, b#29]
      +- Filter isnotnull(b#29)
         +- FileScan parquet default.t_1700986813106[a#28,b#29] Batched: false, DataFilters: [isnotnull(b#29)], Format: Parquet, Location: InMemoryFileIndex[file:/D:/spark-warehouse/t_17009868..., PartitionFilters: [], PushedFilters: [IsNotNull(b)], ReadSchema: struct<a:string,b:string>, SelectedBucketsCount: 2 out of 2

分桶配置验证

通过describe extended确认两张表分桶配置生效:

Bucket Columns: b
Num Buckets: 2

原因分析

  1. 广播Join优先级更高:Spark默认会对小表触发BroadcastHashJoin,只要表大小低于spark.sql.autoBroadcastJoinThreshold(默认10MB)就会自动触发。测试中的t2表仅6条数据,远小于阈值,因此Spark优先选择广播Join而非分桶Join。
  2. 分桶Join的触发逻辑:分桶Join是SortMergeJoin的优化版本,只有当Spark无法触发广播Join时,才会检查分桶条件是否满足(分桶列与Join列一致、分桶数相同),进而触发无Shuffle的分桶Join。

解决方案

1. 禁用自动广播优化

在SparkSession配置中添加spark.sql.autoBroadcastJoinThreshold=-1,强制Spark放弃自动广播,转而考虑分桶Join:

val spark = SparkSession.builder()
  .master("local")
  .enableHiveSupport()
  .appName("test join 1")
  .config("spark.sql.codegen.wholeStage", "false")
  .config("spark.sql.autoBroadcastJoinThreshold", -1) // 禁用自动广播
  .getOrCreate()

2. 验证优化效果

修改后重新运行,物理计划会变为SortMergeJoin且无Shuffle操作,示例如下:

SortMergeJoin [b#27], [b#29], Inner
:- * Project [a#26, b#27]
:  +- * Filter isnotnull(b#27)
:     +- * FileScan parquet default.t_xxxx[a#26,b#27] Batched: false, ... SelectedBucketsCount: 2 out of 2
+- * Project [a#28, b#29]
   +- * Filter isnotnull(b#29)
      +- * FileScan parquet default.t_yyyy[a#28,b#29] Batched: false, ... SelectedBucketsCount: 2 out of 2

3. 可选:调整广播阈值

若不想完全禁用广播,可将阈值调至极小值(如1KB),让t2表超过阈值,从而跳过广播Join:

.config("spark.sql.autoBroadcastJoinThreshold", 1024) // 设置为1KB

内容的提问来源于stack exchange,提问作者Tom

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.05 12:22:19