Spark中预排序分桶表执行Sort Merge Join仍存在Sort步骤的原因
分桶预排序表自Join仍出现Sort步骤的原因分析
问题描述
创建按id分桶并预排序的Parquet表bucketed_table_H,执行自Join时物理计划依然包含Sort步骤,相关代码及物理计划如下:
建表代码
data = [(1, "Alice", "A"), (3, "Charlie", "A"), (2, "Bob", "B"), (4, "David", "B")] schema = ["id", "name", "partition_key"] df = spark.createDataFrame(data, schema=schema) df.repartition(2, f.col("id")).write.mode("overwrite")\ .format("parquet") \ .bucketBy(2, "id") \ .sortBy("id") \ .option("compression", "snappy") \ .saveAsTable("bucketed_table_H")
自Join查询
select * from bucketed_table_H a join bucketed_table_H b on a.id = b.id
物理计划
== Physical Plan == AdaptiveSparkPlan isFinalPlan=false +- SortMergeJoin [id#24L], [id#27L], Inner :- Sort [id#24L ASC NULLS FIRST], false, 0 : +- Filter isnotnull(id#24L) : +- FileScan parquet default.bucketed_table_h[id#24L,name#25,partition_key#26] Batched: true, Bucketed: true, DataFilters: [isnotnull(id#24L)], Format: Parquet, Location: InMemoryFileIndex(1 paths)[file..., PartitionFilters: [], PushedFilters: [IsNotNull(id)], ReadSchema: struct<id:bigint,name:string,partition_key:string>, SelectedBucketsCount: 2 out of 2 +- Sort [id#27L ASC NULLS FIRST], false, 0 +- Filter isnotnull(id#27L) +- FileScan parquet default.bucketed_table_h[id#27L,name#28,partition_key#29] Batched: true, Bucketed: true, DataFilters: [isnotnull(id#27L)], Format: Parquet, Location: InMemoryFileIndex(1 paths)[file.., PartitionFilters: [], PushedFilters: [IsNotNull(id)], ReadSchema: struct<id:bigint,name:string,partition_key:string>, SelectedBucketsCount: 2 out of 2
原因解析
- 预排序元数据未被持久化:写表时用
sortBy("id")指定了桶内排序,但Spark不会把这个排序信息存入Catalog元数据。读表时,Spark只知道表是分桶的,无法确认桶内数据有序,因此会插入Sort步骤满足SortMergeJoin的输入要求。 - 小数据量导致优化逻辑未触发:测试数据仅4条,排序开销可以忽略,Spark优化器不会特意跳过Sort步骤来利用预排序特性。换成百万级以上大数据量时,优化器更可能跳过Sort步骤。
- 自适应执行未完成最终优化:物理计划显示
AdaptiveSparkPlan isFinalPlan=false,说明还处于自适应执行早期阶段,后续虽可能基于运行时统计优化,但数据量太小的情况下,最终也不会触发跳过Sort的逻辑。 - 桶内排序的利用缺乏元数据支撑:Spark的分桶Join优化主要是避免Shuffle,但要利用桶内预排序,需要明确的有序元数据,而当前Spark并不支持持久化这类信息,因此无法自动跳过Sort。
内容的提问来源于stack exchange,提问作者nnqh
相关产品推荐
相关产品推荐

