以下是 Spark 3.1.1 自适应执行(AE)的核心配置清单,包含必开参数、场景化配置及避坑点,可根据实际需求调整:

一、必开基础参数(开启 AE 核心功能)

参数名默认值推荐值说明
spark.sql.adaptive.enabledfalsetrue总开关:启用自适应执行框架
spark.sql.adaptive.coalescePartitions.enabledtrue(依赖总开关)true动态合并小分区(解决 Shuffle 后小文件问题)
spark.sql.adaptive.join.enabledtrue(依赖总开关)true动态调整 Join 策略(自动切换 Broadcast/Shuffle Hash/Sort Merge)
spark.sql.adaptive.skewJoin.enabledfalsetrue启用动态倾斜 Join 优化(检测并拆分倾斜 Key)

二、关键调优参数(按场景优化)

1. 动态合并分区(解决小文件/数据不均)

参数名作用推荐配置
spark.sql.adaptive.coalescePartitions.minPartitionNum合并后最小分区数建议设为集群核数的 1~2 倍(如 100 核集群设为 100~200)
spark.sql.adaptive.coalescePartitions.maxPartitionNum合并后最大分区数避免分区过多导致资源浪费,建议 500~2000(视数据量)
spark.sql.adaptive.coalescePartitions.targetBytes目标分区大小(字节)默认 64MB(67108864),可根据集群块大小调整(如 128MB 设为 134217728)

2. 动态 Join 优化(自动选择最优策略)

参数名作用推荐配置
spark.sql.autoBroadcastJoinThreshold自动广播小表的阈值(字节)默认 10MB(10485760),可根据内存调整(如 50MB 设为 52428800),-1 禁用广播
spark.sql.adaptive.join.shuffleHashJoinThreshold启用 Shuffle Hash Join 的阈值(字节)默认 200MB(209715200),当表大小介于广播阈值和此值之间时,自动用 Shuffle Hash

3. 倾斜 Join 优化(解决数据倾斜)

参数名作用推荐配置
spark.sql.adaptive.skewJoin.skewedPartitionThresholdInBytes判定倾斜分区的阈值默认 10MB,建议设为目标分区大小的 2~3 倍(如目标 64MB 则设为 134217728)
spark.sql.adaptive.skewJoin.skewedPartitionFactor倾斜因子(分区大小/平均大小)默认 5,即当某分区大小是平均值的 5 倍以上时判定为倾斜

三、高级功能参数(按需开启)

参数名作用适用场景
spark.sql.adaptive.forceApply强制对所有查询启用 AE复杂查询(多阶段 Join/聚合),默认 false(简单查询可能有额外开销)
spark.sql.adaptive.dynamicPartitionPruning.enabled动态分区裁剪(减少扫描数据量)分区表多表 Join 场景,默认 false,需手动开启
spark.sql.adaptive.nestedLoopJoin.enabled启用嵌套循环 Join(小表关联大表)极小规模表(<1 万行)与大表 Join,默认 false

四、避坑点

  1. 小数据集性能损耗:AE 会在 Shuffle 后额外计算统计信息,小数据查询(<1GB)可能变慢,可通过 spark.sql.adaptive.forceApply=false 自动跳过。
  2. 倾斜检测延迟:极端倾斜场景(如单 Key 占比 90%)可能检测不及时,需结合 spark.sql.shuffle.partitions 预调整 Shuffle 分区数。
  3. 参数冲突:手动设置 spark.sql.join.preferSortMergeJoin=true 会覆盖 AE 的动态 Join 策略,建议保持默认 false
  4. 内存压力:广播 Join 阈值过大会导致 Driver/Executor 内存溢出,需结合 spark.driver.memoryspark.executor.memory 调整。

五、配置示例(离线批处理场景)

// 启用核心功能
spark.conf.set("spark.sql.adaptive.enabled", "true")
spark.conf.set("spark.sql.adaptive.skewJoin.enabled", "true")
spark.conf.set("spark.sql.adaptive.dynamicPartitionPruning.enabled", "true")

// 调整分区合并策略(目标 128MB 分区)
spark.conf.set("spark.sql.adaptive.coalescePartitions.targetBytes", 134217728)
spark.conf.set("spark.sql.adaptive.coalescePartitions.minPartitionNum", 200)

// 优化 Join 策略(广播阈值 50MB,Shuffle Hash 阈值 500MB)
spark.conf.set("spark.sql.autoBroadcastJoinThreshold", 52428800)
spark.conf.set("spark.sql.adaptive.join.shuffleHashJoinThreshold", 524288000)

// 倾斜检测(分区>256MB 或是平均的 8 倍则判定为倾斜)
spark.conf.set("spark.sql.adaptive.skewJoin.skewedPartitionThresholdInBytes", 268435456)
spark.conf.set("spark.sql.adaptive.skewJoin.skewedPartitionFactor", 8)

根据任务类型(离线/近实时)和数据特征微调,通常能获得 10%~50% 的性能提升,复杂场景效果更明显。

更多推荐