ARTICLE DETAIL

资讯详情

深耕网站视觉设计与运营推广的一线实战洞察。

AQE自适应执行.md

AQE自适应执行.md Spark 的执行计划不是铁板一块跑起来之后 AQE 还会改它。提交 Spark 作业的时候控制台吐出来的那份执行计划多数人当成最终判决看完就关了。它其实只是草稿。真正跑起来之后有一套机制会拿着运行时的真实数据量把这份计划当场改掉这套机制的官方名字叫 AQE自适应查询执行。它默认开着但大多数人对它的全部认知止步于「默认开着」。这篇写给每天在集群上跑大作业的人作业里有大表 join、有 groupBy、分区数动辄几百跑一次几十分钟起步。把 AQE 的机制、它干的三件事和上手的路径一次讲清讲的都是官方文档里写得明明白白的东西。一. 计划为什么有两版写 SQL 的时候Spark 的编译器要提前定好一套执行方案join 用什么算法、分多少个区、谁先谁后。定方案的依据是统计信息表多大、有多少行、key 的分布什么样。麻烦在统计这东西天然不靠谱表可能是昨天统计的行数可能是估的临时目录里的表干脆就没有统计。统计不靠谱是常态不是意外。显式跑过 ANALYZE 的表统计停留在那一天之后的数据一概不知道。没跑过的表Spark 现场去猜猜的依据是文件大小这类粗粒度信号。最没底的是过滤之后的行数一个 where 条件的选择性编译期只能拍脑袋。统计是老的、行数是猜的、分布是拍的三样凑一起计划在开工前就已经和现实脱节。AQE 的做法是把决策推迟到证据齐全之后。带 shuffle 的查询会被 shuffle 边界切成一段一段的执行单元排在前面的单元跑完产出的分区就是落地的真实数据多少个分区、每个多大、行数几何全是刚称完的重量。AQE 拿这份真实统计把下游还没执行的部分重新规划一遍一个 stage 改一次改完再跑。和编译期的 Catalyst 优化器打个比方各管一段。Catalyst 在跑之前把逻辑计划优化成物理计划靠的是规则和静态统计信息止步于查询本身。AQE 接在物理执行里用真实落地的数据量做物理层的二次决策它不改语义只改执行姿势。AQE 挑的时机很讲究它不动编译期只在 shuffle 边界下手。每个 shuffle 阶段跑完落地的数据是多少个分区、每个分区多大、行数多少全是刚出炉的真实数字。AQE 拿这份真实统计把下游还没执行的部分重新规划一遍一个 stage 改一次改完再跑。对着EXPLAIN FORMATTED能直接看到这件事提交作业时输出的顶层是AdaptiveSparkPlan isFinalPlanfalse这一版就是草稿。作业跑完之后再取一次计划isFinalPlan变成true同屏给出 Final Plan 和 Initial Plan 两段。学会 diff 这两段比背十个参数都值。AdaptiveSparkPlan isFinalPlantrue - Final Plan - SortMergeJoin ...按真实统计重排过的版本 - Initial Plan - SortMergeJoin ...提交时的草稿输出形态各版本大同小异以你手上的版本文档为准。二. 第一件事合并小分区AQE 做的第一类事最朴素把碎分区拼起来。场景非常常见大表 join 之后产出了几百个 reduce task跑完一看每个分区就几 MB。碎分区的来源五花八门上游一次过大的 shuffle、过滤之后残存的零星数据都可能造出来。每个 task 都不是免费的调度要发指令启动要建环境结束要汇报JVM 里的固定成本一样不少。活的重量配不上这些固定成本整个 stage 的墙钟时间就耗在一轮轮空转上。拿数字算一遍最直白200 个 5MB 的碎分区按 64MB 的目标拼16 个上下就装完了task 数直接砍到不足一成。目标、下限、并行度优先三个参数的三角关系理清楚合并这一支就没有黑盒了。还有一个容易被忽略的限定AQE 合并的是相邻分区。分区在 shuffle 里有固定的先后次序只有物理上相邻的碎分区才能拼到一起这个约束保证了合并不打乱数据的连续性。advisory 这个参数其实身兼两职合并小分区时它是拼装的目标尺寸拆倾斜分区时它又是拆出来的子分区要对齐的尺寸一个值管两头。同一套合并机制在写盘场景同样救急。写分区文件的个数由分区数决定几百个碎分区落盘就是几百个小文件下游再来读就是几百次打开关闭。合并之后文件数跟着降下游的读取体验一起变好。写盘场景还有个近亲Spark 3.2 起的 rebalance 分区走的是同一套合并加拆分的逻辑专门管写盘前的分区整形3.3 又补了小分区因子的精调参数。对每天在湖里落表的作业来说这一支比查询侧更常见。spark.sql.adaptive.coalescePartitions.enabled默认trueAQE 会按目标尺寸把相邻的碎分区连续合并目标尺寸由spark.sql.adaptive.advisoryPartitionSizeInBytes指定默认 64MB。合并之后 task 数骤降每个 task 拿到的量回到合理区间。跑批前多看一眼分区数是最便宜的性能投资。这里埋着一个大多数人不知道的默认行为。spark.sql.adaptive.coalescePartitions.parallelismFirst默认true它的含义是合并时无视 64MB 那个目标只守住 1MB 的最小分区尺寸尽可能多留 task 来 maximize 并行度。也就是说默认状态下 advisory 那个 64MB 基本是个摆设。配套还有个spark.sql.adaptive.coalescePartitions.minPartitionSize默认 1MB是合并时单分区的下限。官方文档给的建议很直白忙集群上把这个参数设为false让资源利用更有效率别造出一堆小 task。三. 第二件事运行时换 join 策略编译期选 join 算法靠的是估算的表大小。估算说这张表很大计划里就写了 sort-merge join两边各自排序再归并的 join老老实实排两次。维度表偏偏是那种天天在变的对象今天的行数和上次统计时的行数可能已经不是一回事。AQE 在 join 的 shuffle 阶段跑完后量的是真实大小。这一量发现 join 的一侧其实很小低于广播阈值它就把 sort-merge join 换成 broadcast hash join小表全量广播到各节点、大表边读边配对的 join排序全省了原本压在两边的 sort 预算直接归零。判断的时机卡在 join 一侧的 shuffle 阶段完成的那一刻量的就是那一刻落地的真实大小。要是两侧的真实大小都超过广播阈值换不动计划维持 sort-merge 继续跑。spark.sql.adaptive.autoBroadcastJoinThreshold管这道判断默认沿用spark.sql.autoBroadcastJoinThreshold的值设-1可以在 adaptive 框架里禁掉广播。换成广播之后spark.sql.adaptive.localShuffleReader.enabled默认true还会让 Spark 尽量从本地读 shuffle 数据省一轮网络传输。本地读的原理不复杂按 key 重分区的需求消失之后数据落在哪个节点就从哪个节点直接读跨节点拉取的那轮固定动作被整个省掉。这份阈值默认沿用全局那份同名配置想只对 AQE 单独调就显式设这一份。判断的时机卡在 join 一侧的 shuffle 阶段完成的那一刻量的就是那一刻落地的真实大小这次转换在 SQL 页的执行图里同样有节点可查。本地读的意思拆开说产出这份 shuffle 数据的 executor本来就散在各节点上。换成广播之后不再需要按 key 重分区各节点把落在自己盘上的那份 shuffle 文件直接喂给广播哈希 join网络传输省下来。官方对这次换法的评价很诚实说它不如一开始就规划成广播高效但好过继续硬跑 sort-merge至少两边的排序和那轮网络都省了。要是两侧的真实大小都超过广播阈值换不动计划维持 sort-merge 继续跑。换不换的开关始终握在数据手里计划只是执行数据的意思。这也是它和 hint 最大的区别hint 是人拍脑袋告诉计划该怎么做转换是数据自己说话。换法不止一种除了广播哈希 join变更单里还有把 sort-merge join 换成 shuffled hash join靠哈希表配对、不用全局排序的 join的路子同样是拿真实大小做的决策。这条线默认不动触发阈值默认是 0属于知道有就行的进阶项。估错表大小不是罪跑起来愿意改就是好同志。四. 第三件事拆倾斜分区join 跑到一半某个分区的量明显碾压其它分区这就是倾斜长在运行时的样子。AQE 的处理是把这个大分区按目标尺寸拆成几个小分区另一侧对应的行复制多份去陪它配对长尾被物理摊平。判定一个分区算不算倾斜要同时过两道线尺寸超过中位数分区的 5 倍spark.sql.adaptive.skewJoin.skewedPartitionFactor默认 5.0并且尺寸本身超过 256MBspark.sql.adaptive.skewJoin.skewedPartitionThresholdInBytes默认 256MB。两道线是与的关系过一道不过另一道都算不上倾斜。拆的动作是双向的大分区自己拆成若干接近目标尺寸的小分区另一侧对应的行按需复制多份去和拆出来的每一份配对。独木桥就这么变成了几座并排的小桥。复制这个词别慌复制的是参与配对的那部分行拆分保证配对逻辑不变代价是那一小段数据多读几遍、配对多算几遍局部开销换整体不翻车。判定和拆分全程由引擎自动完成不需要任何 hint 也不需要人工介入。两道线的含义值得多看一眼。只比中位数大没用小表的中位数本来就小只超绝对值也没用大表哪个分区都大。两条同时过才动手官方还建议把 256MB 这道线保持大于 advisory 的 64MB不然拆出来的目标比判定线还小逻辑上拧巴。拿默认值算一笔账最直观。中位数 10MB 的 stage 里冒出一个 300MB 的分区30 倍于中位数过了第一道线300MB 也过了 256MB 的第二道线两道全中拆。隔壁 stage 的中位数是 1GB一个 800MB 的分区比中位数小得多第一道线就不过它是健康分区。拆不出来公平拆出来的才是真倾斜。在 Spark UI 的 SQL 页里能直接认出它的动作被拆过的分区在执行图里挂着AQEShuffleRead的节点标识描述里写着 skewed。边界同样要说清这套倾斜处理只管 sort-merge join 和 shuffled hash join聚合操作的倾斜它不碰。聚合的热点没有自动机制兜着得靠手工把热点打散再聚这笔账要在设计聚合口径的时候就算进去。join 倾斜的手工方案同样是存在的热点 key 单独拎出来处理、构造随机前缀打散重分布都是路。AQE 的位置是把这些手工活里「分区级、可自动识别」的那部分接了过去剩下的仍然要人来判断。join 倾斜的手工方案同样是存在的热点 key 单独拎出来处理、构造随机前缀打散重分布都是路。AQE 的位置是把这些手工活里「分区级、可自动识别」的那部分接了过去剩下的仍然要人来判断。五. 代价与边界AQE 不是白来的魔法它有三条边界要认。第一只在有 shuffle 的查询里起作用。重优化发生在 shuffle 边界一条 SQL 从头到尾没有 exchangeAQE 全程没有出手的机会。exchange 的有无肉眼不好断EXPLAIN 输出里数一数 Exchange 算子最稳。第二重优化本身有开销。每个 stage 结束都要做一轮统计和重规划小查询、跑得飞快的查询这份开销占比就高收益反而看不出来。它服务的对象是那些跑得久、数据量大的作业。官方没有给过重优化开销的具体数字定性地讲跑得越久、数据越重的作业这点开销越接近零头。能力不是一天长齐的。3.0 带进了合并分区、拆倾斜、join 策略转换这批主干3.2 起 AQE 整体默认开启同一版本补上了并行度优先的开关和自定义代价评估器3.3 又增强了强制倾斜优化这些边角。个别规则不想让它在某条查询上生效spark.sql.adaptive.optimizer.excludedRules支持逐条排除这是 3.1 就有的口子。能力盘子上官方文档给 AQE 列了六个子项合并分区、拆倾斜分区、sort-merge 转广播、sort-merge 转 shuffled hash、优化倾斜 join外加一节高级定制盘面还在随版本变宽。第三改的只是物理执行层。join 算法、分区数、分区尺寸这些它管SQL 的语义它一个字都不动写错的 join 条件、漏掉的过滤AQE 一个都不会帮你修。说到底它是执行层的顺风局放大器不是语义层的救命药。两个常见的误区顺手排掉。一个是把 advisory 往大了调想治倾斜方向错了治倾斜看的是那两道判定线advisory 管的是合并和拆分的目标尺寸。另一个是小作业跑得慢怪 AQE重优化的固定开销在小查询里占比高这种场景它本来就不是为你准备的。还有一笔资源账要自己盯合并让 task 变少单个 task 拿到的量就变大executor 的内存水位要看一眼别合并出撑爆内存的大 task。大分区的垃圾回收压力也是老毛病合并不当会把 GC 的老病请回来GC 表现值得多看一眼。怎么判断它有没有帮到你看三个地方就够。初始计划和最终计划的差异是它动手的直接证据。task 数量的变化合并前后的对比在 SQL 页一眼可见。stage 耗时分布有没有从长尾变均匀这是最终的效果验收。六. 上手就三步第一步确认总开关活着。这一步零成本看一眼配置的事却是后面一切的前提。spark.sql.adaptive.enabled从 3.2.0 起默认true4.x 时代这个默认没有回退过检查一下有没有人在老版本的习惯里把它显式关了。第二步拿EXPLAIN FORMATTED看计划。提交时输出只有一版草稿顶层是AdaptiveSparkPlan isFinalPlanfalse。作业跑完再取一次isFinalPlan变true同屏给出 Final Plan 和 Initial Plan 两段diff 这两段就是 AQE 的功劳簿它的功劳簿和免责声明同时到手。第三步跑完去 Spark UI 的 SQL 页找 AQE 的痕迹。每条查询都挂着执行图AQE 动过的位置会在节点描述里留字合并的、拆分的、换策略的各留各的痕。被合并的分区、被拆的倾斜分区在执行图里都有对应的节点标识肉眼可查。常用参数记两个就够入门advisoryPartitionSizeInBytes管合并与拆分的目标尺寸parallelismFirst在忙集群上建议关掉让 advisory 真正生效。把这两个的分工记住其余的参数都会自己归位。顺手把这一族参数收成一张表值都是 4.1 文档写的默认。参数默认管什么spark.sql.adaptive.enabledtrue总开关spark.sql.adaptive.coalescePartitions.enabledtrue合并小分区spark.sql.adaptive.advisoryPartitionSizeInBytes64MB合并与拆分的目标尺寸spark.sql.adaptive.coalescePartitions.minPartitionSize1MB合并时的单分区下限spark.sql.adaptive.coalescePartitions.parallelismFirsttruetrue 时无视目标尺寸保并行度spark.sql.adaptive.autoBroadcastJoinThreshold沿用同名非 adaptive 配置运行时换广播的阈值spark.sql.adaptive.localShuffleReader.enabledtrue换广播后本地读 shufflespark.sql.adaptive.skewJoin.enabledtrue倾斜分区拆分spark.sql.adaptive.skewJoin.skewedPartitionFactor5.0倾斜判定的倍数线spark.sql.adaptive.skewJoin.skewedPartitionThresholdInBytes256MB倾斜判定的绝对线三步走完这一族参数就都在你的掌控里了。
返回列表