以下是 Spark 3.1.1 自适应执行(AE)的核心配置清单,包含必开参数、场景化配置及避坑点,可根据实际需求调整:
一、必开基础参数(开启 AE 核心功能)
| 参数名 | 默认值 | 推荐值 | 说明 |
|---|
spark.sql.adaptive.enabled | false | true | 总开关:启用自适应执行框架 |
spark.sql.adaptive.coalescePartitions.enabled | true(依赖总开关) | true | 动态合并小分区(解决 Shuffle 后小文件问题) |
spark.sql.adaptive.join.enabled | true(依赖总开关) | true | 动态调整 Join 策略(自动切换 Broadcast/Shuffle Hash/Sort Merge) |
spark.sql.adaptive.skewJoin.enabled | false | true | 启用动态倾斜 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 |
四、避坑点
- 小数据集性能损耗:AE 会在 Shuffle 后额外计算统计信息,小数据查询(<1GB)可能变慢,可通过
spark.sql.adaptive.forceApply=false 自动跳过。 - 倾斜检测延迟:极端倾斜场景(如单 Key 占比 90%)可能检测不及时,需结合
spark.sql.shuffle.partitions 预调整 Shuffle 分区数。 - 参数冲突:手动设置
spark.sql.join.preferSortMergeJoin=true 会覆盖 AE 的动态 Join 策略,建议保持默认 false。 - 内存压力:广播 Join 阈值过大会导致 Driver/Executor 内存溢出,需结合
spark.driver.memory 和 spark.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")
spark.conf.set("spark.sql.adaptive.coalescePartitions.targetBytes", 134217728)
spark.conf.set("spark.sql.adaptive.coalescePartitions.minPartitionNum", 200)
spark.conf.set("spark.sql.autoBroadcastJoinThreshold", 52428800)
spark.conf.set("spark.sql.adaptive.join.shuffleHashJoinThreshold", 524288000)
spark.conf.set("spark.sql.adaptive.skewJoin.skewedPartitionThresholdInBytes", 268435456)
spark.conf.set("spark.sql.adaptive.skewJoin.skewedPartitionFactor", 8)
根据任务类型(离线/近实时)和数据特征微调,通常能获得 10%~50% 的性能提升,复杂场景效果更明显。
所有评论(0)