同一Spark DataFrame按不同键分别Join多表的性能优化问询
问题结论先行
- 你提到的将df1分别按
key1、key2分桶两次是完全可行的方案,在存储成本、预处理成本可接受的前提下,能拿到两次join的最优性能。 - 也存在无需存储两份df1的优化策略,可根据你的业务场景选择。
可选优化方案
方案1:双分桶持久化(性能最优)
- 操作步骤:
- 提前对静态df1做两次分桶存储,分别按
key1、key2分桶,分桶数要和df2按key1的分桶数、df3按key2的分桶数完全一致 - 开启Spark分桶配置:
spark.sql.sources.bucketing.enabled=true
- 提前对静态df1做两次分桶存储,分别按
- 适用场景:df1为冷数据,更新频率极低,join为高频常驻任务,存储空间充裕
- 收益:两次join都可完全避免shuffle,直接做本地桶内join,性能提升幅度最大
方案2:复合分桶+桶内排序(平衡存储和性能)
- 操作步骤:
- 对df1按
key1, key2复合键分桶,分桶数和df2按key1的分桶数一致,同时每个分桶内按key2排序存储 - df3按
key2分桶时,分桶数设置为df1分桶数的整数倍
- 对df1按
- 适用场景:存储资源有限,无法存储两份全量df1的场景
- 收益:
- 按
key1的join1完全不需要shuffle,直接桶内join - 按
key2的join2虽然仍需要shuffle,但df1的每个分桶已按key2排序,shuffle后无需全量排序,直接做归并join,计算开销可降低40%~60%
- 按
方案3:预重分区缓存(适合单次批量任务)
- 操作步骤:
单次任务启动时,先对df1做两次重分区并持久化,后续所有join复用分区后的结果:import org.apache.spark.storage.StorageLevel val df1ByKey1 = df1.repartition(col("key1")).persist(StorageLevel.MEMORY_AND_DISK) val df1ByKey2 = df1.repartition(col("key2")).persist(StorageLevel.MEMORY_AND_DISK) // 后续join直接使用重分区后的df val join1 = df1ByKey1.join(df2, "key1") val join2 = df1ByKey2.join(df3, "key2") - 适用场景:单次任务内需要多次执行两种join,df1在任务生命周期内不会更新
- 收益:单次任务内仅需要对df1做两次shuffle,后续所有join无需重复shuffle,不需要长期占用外部存储资源
方案4:开启AQE自适应优化(Spark 3.0+ 无侵入优化)
- 操作步骤:开启以下配置即可,无需修改业务代码:
spark.sql.adaptive.enabled=true spark.sql.adaptive.coalescePartitions.enabled=true spark.sql.adaptive.join.enabled=true spark.sql.adaptive.localShuffleReader.enabled=true - 适用场景:无法修改数据预处理逻辑、代码改造成本高的场景
- 收益:运行时自动优化shuffle分区大小,动态调整join策略,即使不提前做预分桶,也能大幅降低join的额外开销,部分场景下可自动将小表join转广播join。
内容的提问来源于stack exchange,提问作者Amit Joshi
相关产品推荐
相关产品推荐

