Spark SQL AQE 优化策略一:合并Shuffle分区源码剖析 本文剖析 Spark SQL AQEAdaptive Query Execution自适应查询执行中动态合并 Shuffle 分区的核心逻辑包括CoalesceShufflePartitions.apply规则的整体流程ShufflePartitionsUtil.coalescePartitions的合并算法细节一、背景与核心思想Shuffle分区的数量会大大影响Spark的性能。如果shuffle分区太少每个分区都包含了太多的数据可能导致内存不足发生OOM从而kill掉作业无法充分发挥并行计算的优势另一方面过多的shuffle分区会导致过多的小任务在任务调度和管理上造成过多的开销尤其当数据量较小时。正常情况下shuffle分区的数量可以通过调整spark.sql.shuffle.partition参数进行配置。参考Spark SQL Shuffle 分区数生成机制详解但数据的大小在整个任务的不同阶段会有很大的变化。例如某些阶段涉及过滤器或聚合操作但shuffle分区号是为Spark作业中的所有阶段预定义的对一个阶段来说是最佳的shuffle分区号可能在另一个阶段的性能就显得不是很理想。因此需要在作业的运行过程中根据运行时特定阶段的数据量动态调整shuffle分区的数量。在这一背景下AQE提供了动态配置shuffle分区数量规则CoalesceShufflePartitoins。原理是 利用运行时真实的 Map 输出统计信息把连续的小分区合并成接近目标大小的大分区动态减少分区数从而提升性能。二、CoalesceShufflePartitions 规则解析CoalesceShufflePartitions的转化逻辑在apply(plan: SparkPlan)中实现会输出经过合并分区优化后的物理计划。org/apache/spark/sql/execution/adaptive/CoalesceShufflePartitionsoverridedefapply(plan:SparkPlan):SparkPlan{// 1、总开关当spark.sql.adaptive.coalescePartitions.enabled未开启时则直接返回原计划spark.sql.adaptive.coalescePartitions.enabled默认为true。if(!conf.coalesceShufflePartitionsEnabled){returnplan}defcollectShuffleStages(plan:SparkPlan):Seq[ShuffleQueryStageExec]planmatch{casestage:ShuffleQueryStageExecSeq(stage)case_plan.children.flatMap(collectShuffleStages)}// 2、递归遍历计划树收集所有 ShuffleQueryStageExecvalshuffleStagescollectShuffleStages(plan)// 3、只有所有 Shuffle 都支持合并时才继续,否则返回原计划// 由 repartition 引入的 ShuffleExchange 不支持改变分区数用户显式指定AQE 不应擅自改动。if(!shuffleStages.forall(ssupportCoalesce(s.shuffle))){plan}else{// ShuffleQueryStageExec#mapStats returns None when the input RDD has 0 partitions,// 3.1 收集所有map stage的运行时信息.valvalidMetricsshuffleStages.flatMap(_.mapStats)// 3.2 获取map stage的所有分区数valdistinctNumPreShufflePartitionsvalidMetrics.map(statsstats.bytesByPartitionId.length).distinct// 3.3 要求 validMetrics 非空且map stage 所有 Shuffle 的分区数一致去重后只有 1 个值。if(validMetrics.nonEmptydistinctNumPreShufflePartitions.length1){// 获取最小分区数valminPartitionNumconf.getConf(SQLConf.COALESCE_PARTITIONS_MIN_PARTITION_NUM).getOrElse(session.sparkContext.defaultParallelism)// 合并 shuffle 分区逻辑valpartitionSpecsShufflePartitionsUtil.coalescePartitions(validMetrics.toArray,advisoryTargetSizeconf.getConf(SQLConf.ADVISORY_PARTITION_SIZE_IN_BYTES),minNumPartitionsminPartitionNum)// This transformation adds new nodes, so we must use transformUp here.valstageIdsshuffleStages.map(_.id).toSet plan.transformUp{// even for shuffle exchange whose input RDD has 0 partition, we should still update its// partitionStartIndices, so that all the leaf shuffles in a stage have the same// number of output partitions.casestage:ShuffleQueryStageExecifstageIds.contains(stage.id)CustomShuffleReaderExec(stage,partitionSpecs)}}else{plan}}}2.1 获取最小分区数和建议目标分区大小valminPartitionNumconf.getConf(SQLConf.COALESCE_PARTITIONS_MIN_PARTITION_NUM).getOrElse(session.sparkContext.defaultParallelism)最小分区数优先取配置spark.sql.adaptive.coalescePartitions.minPartitionNum未设置则回退到defaultParallelismadvisoryTargetSizeconf.getConf(SQLConf.ADVISORY_PARTITION_SIZE_IN_BYTES)建议目标分区大小spark.sql.adaptive.advisoryPartitionSizeInBytes默认值未64MB。2.2 计算 shuffle 分区合并方案coalescePartitions调用coalescePartitions计算真正的合并方案。valpartitionSpecsShufflePartitionsUtil.coalescePartitions(validMetrics.toArray,advisoryTargetSizeconf.getConf(SQLConf.ADVISORY_PARTITION_SIZE_IN_BYTES),minNumPartitionsminPartitionNum)2.2.1 方法签名defcoalescePartitions(mapOutputStatistics:Array[MapOutputStatistics],advisoryTargetSize:Long,minNumPartitions:Int):Seq[ShufflePartitionSpec]输入mapOutputStatistics每个 Shuffle 的 Map 端输出统计核心字段bytesByPartitionId: Array[Long]第 i 个 reduce 分区的字节数。advisoryTargetSize建议目标分区大小。minNumPartitions合并后允许的最小分区数。输出Seq[CoalescedPartitionSpec]每个形如CoalescedPartitionSpec(startReducerIndex, endReducerIndex)表示合并后的一个新分区由原始 reduce 分区区间[start, end)拼成。2.2.2 第一步计算真正的目标大小 targetSize为了合并shuffle分区CoalesceShufflePartitoins规则首先计算出真正的合并分区的目标大小target size。valtotalPostShuffleInputSizemapOutputStatistics.map(_.bytesByPartitionId.sum).sumvalmaxTargetSizemath.max(math.ceil(totalPostShuffleInputSize/minNumPartitions.toDouble).toLong,16)valtargetSizemath.min(maxTargetSize,advisoryTargetSize)totalPostShuffleInputSize获取map输出统计中的总shuffle数据大小。maxTargetSize 所有shuffle分区总量 / minNumPartitions为凑够最小分区数每个分区平均能分到的最大字节数。外层math.max(..., 16)防止空表时为 0。targetSize min(maxTargetSize, advisoryTargetSize)取较小值保证最终分区数不少于minNumPartitions避免并行度不足。一句话targetSize 会在“64MB”和“为了凑够最小分区数而必须的更小值”之间取小。2.2.3 第二步分区数一致性判断valdistinctNumShufflePartitionsmapOutputStatistics.map(statsstats.bytesByPartitionId.length).distinct assert(distinctNumShufflePartitions.length1,...)valnumPartitionsdistinctNumShufflePartitions.head获取并要求所有 Shuffle 的 reduce 分区数一致numPartitions为统一分区数。2.2.4 第三步贪心合并主循环valpartitionSpecsArrayBuffer[CoalescedPartitionSpec]()varlatestSplitPoint0// 当前合并区间的起点varcoalescedSize0L// 当前合并区间已累计的大小vari0while(inumPartitions){// (1) 计算第 i 个 reduce 分区在“所有 Shuffle”上的总大小vartotalSizeOfCurrentPartition0Lvarj0// 对同一 reduce 下标 i把所有 Shuffle 中该分区大小相加。用于处理 Join 等多输入 Shuffle 场景同分区号会在同一 reduce 任务一起读取。while(jmapOutputStatistics.length){totalSizeOfCurrentPartitionmapOutputStatistics(j).bytesByPartitionId(i)j1}// (2) 用连续的索引将shuffle分区打包成一个合并的分区直到再增加一个分区将超过目标大小if(ilatestSplitPointcoalescedSizetotalSizeOfCurrentPartitiontargetSize){partitionSpecsCoalescedPartitionSpec(latestSplitPoint,i)latestSplitPointi coalescedSizetotalSizeOfCurrentPartition}else{coalescedSizetotalSizeOfCurrentPartition}i1}// (3) 收尾最后一段区间partitionSpecsCoalescedPartitionSpec(latestSplitPoint,numPartitions)注i latestSplitPoint保证区间至少含一个分区避免空区间也意味着单个分区即使超标也不拆分独占一个合并分区保证 reduce 端读取的连续性。2.2.5 具体示例2.2.5.1 单 Shuffle 示例假设numPartitions 8单个 Shuffle各分区大小MB分区: 0 1 2 3 4 5 6 7 大小: 20 30 10 70 5 5 5 50设targetSize 64MB逐步执行i分区大小coalescedSize(加之前)是否切分结果0200否累计201302050≤64否累计502105060≤64否累计603706013064切输出[0,3)累计7045707564切输出[3,4)累计555510≤64否累计10651015≤64否累计15750156564切输出[4,7)累计50收尾输出[7,8)最终合并方案8 个分区 → 4 个[0,3) - 203010 60MB [3,4) - 70MB单个大分区独占不拆分 [4,7) - 555 15MB [7,8) - 50MB2.2.5.2 Join 场景示例多个 Shuffle 纵向求和上面的单 Shuffle 示例无法体现算法中“纵向求和”那段逻辑。Join 是最典型的多输入 Shuffle 场景下面用一个 SortMergeJoin 的例子说明。考虑如下 SQLSELECT*FROMt1JOINt2ONt1.keyt2.keySortMergeJoin 会对t1、t2各做一次 Shuffle两个 Shuffle 使用相同的分区数且相同 key 落到相同的分区号。因此 join 时reduce 任务i会同时读取两个 Shuffle 的第 i 个分区并做归并连接。所以在合并分区时判断第i个分区的负载必须把两个 Shuffle 中该分区的大小加起来——这正是主循环里内层j循环做的事vartotalSizeOfCurrentPartition0Lvarj0while(jmapOutputStatistics.length){// j 遍历 t1、t2 两个 ShuffletotalSizeOfCurrentPartitionmapOutputStatistics(j).bytesByPartitionId(i)j1}假设numPartitions 6两个 Shuffle 各分区大小MB如下分区: 0 1 2 3 4 5 t1 Shuffle: 10 15 40 5 5 30 t2 Shuffle: 20 10 50 5 5 20 ------------------------------------------------ 纵向求和: 30 25 90 10 10 50设targetSize 64MB用“纵向求和”后的大小逐步执行i求和大小coalescedSize(加之前)是否切分结果0300否累计301253055≤64否累计552905514564切输出[0,2)累计903109010064切输出[2,3)累计104101020≤64否累计20550207064切输出[4,5)累计50收尾输出[5,6)最终合并方案6 个分区 → 4 个[0,2) - t1(1015) t2(2010) 55MB [2,3) - t1(40) t2(50) 90MB单分区超标独占不拆分 [4,5) - t1(5) t2(5) 10MB [5,6) - t1(30) t2(20) 50MB2.3 应用转换到计划树valstageIdsshuffleStages.map(_.id).toSet plan.transformUp{casestage:ShuffleQueryStageExecifstageIds.contains(stage.id)CustomShuffleReaderExec(stage,partitionSpecs)}用transformUp自底向上遍历把匹配的ShuffleQueryStageExec外包一层CustomShuffleReaderExec。即使某 Shuffle 输入 RDD 有 0 分区也统一套用partitionSpecs保证一个 stage 内所有叶子 Shuffle 输出分区数一致。