拓十年匠心定制 · 商业建站与技术教学双线并行 咨询热线:400-886-1026 service@lmnt.cn
ARTICLE DETAIL

资讯详情

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

Spark SQL性能优化:Auron重写执行计划与向量化加速实践

Spark SQL性能优化:Auron重写执行计划与向量化加速实践 1. 瓶颈定位一个跑了8小时的任务到底卡在了哪儿上个月我接手一个Spark SQL性能优化的活儿线上有个凌晨跑的ETL任务天天超时集群被它拖得别的作业都在排队。乍一看资源都够Executor内存给得也不小但每天就是干不完。我花了整整三天把执行计划翻了个底朝天最后锁定了三类典型的瓶颈也正是这三类问题促使我把Auron的方案正式落地。先说结论方便你对号入座绝大多数Spark SQL慢查询根本原因不在数据量大而在执行路径太长。这里的“长”包含三层意思——磁盘扫描范围大、Shuffle次数多、表达式求值路径繁琐。你打开Spark UI看Stage的执行时间分布基本一目了然某个Stage占了80%以上的时间而那个Stage里几乎全是Scan加ExchangeShuffle后面的处理环节反而是瞬时的。我在那个8小时任务里看到的现象很有意思只有两个大Stage一个负责把两张千万级大表做Join另一个是做一列日期字符串的解析转换。第二个Stage为什么慢因为代码里写了一个自定义UDF对每一行都去构造SimpleDateFormat实例再解析、格式化、拼接。这个逻辑每一行都在重复做GC压力极大CPU反而没跑满。这种问题不是个例。很多DataFrame代码看着简洁优雅实际执行计划里塞满了隐性开销。比如对日期列做加减计算社区里大量代码用UDF或者to_date加字符串处理而不是直接用内置的date_add、add_months这些表达式。再比如合并DataFrame时图省事用unionByName加distinct深层变成了两次Shuffle加一次全量去重。这类问题用一句话概括执行计划不够干净。Spark SQL的Catalyst优化器虽然做了大量规则优化但面对复杂表达式、UDF、以及某些边界条件时仍然会留下很多“不理想但正确”的计划。Auron做的事情本质上就是把这层“不理想”尽量抹掉在不改你业务代码的前提下让Spark自己跑得更聪明。1.1 表达式求值路径每一行都在重复做的事情我先展开说表达式这块因为这是大多数人最容易忽略、也最容易优化的点。Spark SQL处理一行数据时表达式树上的每个节点都要执行一次eval。如果一个表达式嵌套了十层函数比如substring(concat(to_date(...), ...), ...)那么每一行都要完整走完这十层。更糟的是很多函数是WholeStageCodegen无法完全拆解的导致回退到行式迭代模式Volcano模型性能差一个数量级。你可以回忆一下自己写过的代码——有没有在DataFrame的withColumn里传一个Lambda表达式有没有用udf包一个函数去处理日期或字符串这些都是把性能往坑里推的操作。Auron第一个核心能力就是做表达式自动重写。它不是简单地做语法层面的替换而是会把表达式树里那些低效的路径识别出来改写成等价的、更适合代码生成的方式。举一个例子你写df.withColumn(year, year(to_date(col(dt))))如果dt本身是yyyy-MM-dd格式的字符串Auron会识别出这其实是固定的格式解析会把它重写成直接用substring取前四位并配合整数转换。因为不需要构造解析器单纯从内存里截取字符串再转整数速度是数量级的差距。这项能力依赖一个庞大的表达式模式库。每识别出一种模式就能节省一类工作。我在前面的8小时任务里就是靠这个能力把那串日期解析的UDF整体替换掉——不是换UDF实现而是压根不用UDF用内置表达式一步到位。1.2 Shuffle开销被低估的“数据传输税”第二个瓶颈是Shuffle。我常说Shuffle是Spark里的“数据传输税”你每做一次groupBy、join、distinct、orderBy都要按Key做一次全网重分区数据落盘、序列化、网络传输、反序列化这一整套流程的成本通常比计算本身高一个量级。还是回到我那个任务。两张千万级表做Join两张表其实都预先按user_id用bucketBy建了桶但业务代码里并没有利用这个特性而是走了正常的Hash Join流程。于是在Shuffle阶段全量数据被重新分区、重新写磁盘白白多了一次巨大的IO。Auron在执行计划层面做的事情就是识别这种可复用分区信息。它检查两个DataFrame的血缘如果发现两张表的父RDD都已经按相同Key做了Hash分区就会把不必要的Shuffle节点从执行计划里裁剪掉直接走Bucket Join或Partitioned Join。这个思路和Spark 3.x的AQEAdaptive Query Execution不同。AQE是在运行后根据统计信息动态调整比如把SortMergeJoin转成BroadcastJoin而Auron是在计划生成时就提前判断“哪些Shuffle是可以通过血缘推导省略的”。两者不冲突反而能叠加——先用Auron省掉原本就多余的Shuffle再用AQE把剩下的Join策略调优。1.3 被优化器放过的漏网点计划里的“正确但低效”最后是Catalyst本身遗漏的优化点。Spark的优化规则很强大但不是万能的。我举一个很常见的例子df.filter(col(age) 18).join(df2, user_id)。按理说age 18这个过滤应该下推到Scan层减少读取的数据量。但如果过滤条件出现在Join之后、或者被包在一层withColumn里Catalyst的谓词下推可能就看不透整个Filter就会在Join完成后才执行导致大表全量参与Join。Auron内部维护了一批额外的优化规则专门针对这类**“Catalyst看走眼”**的场景。它会额外做一轮谓词下推、投影裁剪、甚至重排列顺序让过滤器尽量靠近数据源让投影尽量早地砍掉不用的列。有次我跑一个上千列的宽表业务只用到其中八列Auron介入后扫描IO直接掉了85%以上。原理一点也不神秘就是提前把不需要的列在Scan阶段干掉但Catalyst默认不会那么激进因为某些列虽然在当前计划用不到但可能在后续操作里用到保守策略导致全列扫描。2. Auron的加速思路不改变你的代码改变执行方式Auron的设计目标从一开始就定得很死不要求用户改任何一行业务代码。你要做的就是把它挂到SparkSession上剩下的自动发生。这一点非常重要因为现实项目里很难说服业务方为了性能优化去改他们的数据处理逻辑——他们没时间也不愿意承担引入bug的风险。那Auron具体是怎么“不改变代码但改变执行方式”的我用两条主线来拆解一条是做执行计划的深度重写另一条是做物理执行层的向量化。前者解决“做无用功”的问题后者解决“做功太慢”的问题。2.1 算子级表达式重写从模式库到等价改写先讲执行计划重写。Auron内部有一颗优化规则的流水线通过Spark的RuleExecutor机制挂载进去和Catalyst自带的优化规则并列运行。它的核心逻辑就是一个大号模式匹配引擎遍历整棵执行计划树每遇到一个节点就用模式库里的模式去匹配命中就替换成等价但更高效的子计划。这个模式库里有什么我列举几类我实际用下来收益最大的字符串转日期类把to_date(col, yyyy-MM-dd)配合year/month/day提取重写成substring加cast的整数组合。日期区间过滤把col(dt) 2024-01-01 and col(dt) 2025-01-01这一类重写成单次字典序范围扫描而不是逐行调用日期解析函数。UDF转内置表达式匹配特定签名和功能的UDF比如正则提取、字符串切割、时间戳转换等等价替换成内置函数。这需要业务方在部署前做一次性配置把UDF和内置表达式的对应关系告诉Auron。常量折叠增强把current_date() - interval 1 year这类在计划里可以提前求值的表达式在优化阶段就算好而不是每行执行时再算。复杂聚合拆解把sum(case when ... then ... else ... end)这类的自定义聚合逻辑拆成多个简单的内置聚合然后合并让Spark可以走更高效的聚合路径。这些重写规则有一个共同点它们都是语义等价的不会改变计算结果只会改变执行的物理方式。正因为如此Auron才能做到“无感接入”——你不需要相信它只需要验证它算出来的结果和原来一致就行。2.2 列式批处理从行式逐条到批量向量化执行计划重写解决的是“少干点活”向量化解决的是“干活更快”。如果你对性能调优有一些经验应该知道Spark从2.0开始就引入了WholeStageCodegen把一串算子编译成一段Java代码减少虚函数调用。但Codegen有一个众所周知的限制——它的单位仍然是行。它把多行处理合并成循环本质上还是逐行调用表达式只是减少了方法调用开销。Auron的向量化路线从这里切入。它改变了数据在算子间的传递方式不再用UnsafeRow逐行传递而是用列式批量块传递。一个批次默认4096行每一列的数据连续存放在一段内存里。对这样的组织方式做计算天然适合SIMD指令的批量处理——同一个操作一次性作用在一整列数据上。我举一个直观的例子计算col_a col_b。行式模式下每一行都要从两个Row里取出对应字段、做加法、再写回结果行向量化模式下是两块连续int数组做元素级加法很多JVM实现里这可以直接编译成高度优化的循环甚至自动向量化JIT的Superword级别优化。在真实基准测试里纯算术计算场景下向量化通常能比Codegen快2到4倍但让我说句实话——没有哪种加速方案是包打天下的。向量化对内存布局、数据类型、GC压力都有影响对不同负载的收益差异很大。Auron的策略是启动时通过一个代价模型自动判断如果一个算子子树能够被整体向量化且预估收益超过阈值就启用否则老老实实走原来的Codegen路径。这个代价模型很重要。比如一个小表上的简单select加不加向量化几乎没差别但向量化的初始化有开销收益为负而面对大表的全列扫描加聚合向量化收益极为显著。2.3 基于血缘的Shuffle裁剪把没必要的Exchange剪掉前面提到的Shuffle裁剪是Auron的另一条主线。我先解释一个概念分区血缘Partition Lineage。Spark的每个RDD/DataFrame都记录了自己是怎么从父RDD变换而来的。如果一张表是按user_id做了Hash分区再经过一系列Filter、Project、甚至Union这些算子通常会保留父RDD的分区方式。Auron会沿着血缘追溯在EnsureRequirements这个物理计划优化阶段介入。Spark原本在这个阶段会给每个算子分配要求的Distribution如果子节点要求HashPartitioning且父节点已经是相同Key的Hash分区就会在两者之间插入ShuffleExchange节点。Auron做的事是检查这个Exchange是否真的必要——如果父节点的实际分区方式已经满足要求就直接把Exchange去掉把两个算子串起来。这个优化在Join场景里收益最大。两张表如果都按Join Key建过桶bucketBy血缘里会保留BucketedTable这个标记Auron能直接让Spark走SortMergeJoin而不需要任何Shuffle这是Spark自带的spark.sql.sources.bucketing.enabled经常做不到的——因为它只在特定条件下生效而Auron的检查更激进也更灵活。有一点要坦白说血缘分析本身有计算成本。Dataset的血缘树如果很复杂几十层变换遍历需要花一些时间。Auron做了一层缓存把LogicalPlan的哈希值和分区信息存下来二次访问同一段血缘时只需要查缓存。在TB级任务里这个分析成本通常小于整体执行时间的0.5%可以忽略不计。3. 接入Auron的完整流程与关键配置理论聊了不少现在进入实操环节。这篇博文的读者应该大多是有一定Spark经验的工程师所以我按实际部署的节奏来写每一步都带解释——为什么这么做以及踩坑时怎么排查。3.1 环境准备Spark版本与依赖引入Auron目前支持Spark 3.2及以上版本更老的版本没有测试过不建议直接上生产。它依赖Spark内部的QueryExecution和SparkSessionExtensions接口这两个API在不同小版本之间相对稳定但我仍然建议你锁定一个Spark版本再锁Auron版本不要混用。以Maven项目为例依赖这样加dependency groupIdio.github.auron/groupId artifactIdauron-spark3_2.12/artifactId version2.1.0/version /dependencyScala 2.13的Spark发行版3.4也有对应的构件把2.12后缀换成2.13就行。如果你的集群是CDP或HDInsight这种发行版依赖可能冲突建议用provided作用域把Auron打进作业Jar而不是集群ClassPath这样升级和回滚都更灵活。我踩过的一个坑是jar包顺序问题。Auron要用到Spark的一些内部类如果ClassPath里同时有多个Spark版本轻则启动报错重则优化规则静默失效。用spark-submit --master yarn提交作业时记得把Auron的jar放在--jars靠前的位置并用spark.driver.userClassPathFirsttrue隔离。3.2 参数配置开与关的取舍Auron的配置项不多但每个都值得认真调。默认配置是“保守稳妥”风格不会为了性能牺牲稳定性。以下是我实测下来最核心的几个参数spark.auron.enabledtrue # 总开关默认true spark.auron.rewrite.expressiontrue # 表达式重写开关默认true spark.auron.rewrite.shuffletrue # Shuffle裁剪开关默认true spark.auron.vectorized.enabledtrue # 向量化执行开关默认true spark.auron.vectorized.batchSize4096 # 批处理行数默认4096 spark.auron.vectorized.minExprCount3 # 触发向量化的最少表达式数 spark.auron.autoBroadcastJoinThreshold10485760 # 10MB覆盖Spark默认稍大一点前四个开关是全局的后两个是向量化模块的细节参数。batchSize值得说两句不是越大越好。批处理越大列式内存块的占用越高GC压力也越大批大小太小向量化又发挥不出优势。在我的经验里4096是甜点在几种典型负载上都测过吞吐量最高。如果你的机器内存充裕可以试试8192有些场景能再快5%左右如果频繁GC那就降到2048。autoBroadcastJoinThreshold这个参数比较有意思。Auron会把原有的阈值调大一点因为向量化执行让Broadcast侧表的构建和查询都快了不少所以可以更激进地选择BroadcastJoin。但代价是Driver端内存占用会上升广播大表时要掂量一下我一般控制在10MB以内遇到超过阈值的仍然走SortMergeJoin。3.3 验证Auron是否真的生效接入之后怎么确认优化真的在起作用这是我最常被问的问题。答案不是看作业跑得快了多少而是看执行计划。在Spark Shell里跑一句val df spark.range(10000000).withColumn(year, year(to_date(lit(2024-01-01)))) df.explain(extended)如果没有Auronyear(to_date(...))会保留为一串函数嵌套的表达式树启用Auron后你会看到计划里这一年提取被重写成了substring加cast之类的操作。再跑一个Join观察是否存在多余的Exchange节点。另一个有效的办法是看Spark UI里每个Stage的Shuffle Read/Write字节数。如果一个作业在启用Auron后Shuffle总量明显下降说明Shuffle裁剪起作用了。我那个8小时任务优化后Shuffle Write从12TB降到了2.1TBStage数从58个减到37个效果非常直观。还可以开启Auron自己的统计日志spark.auron.stats.enabledtrue spark.auron.stats.logLevelINFO它会在每个Stage结束时打印一条摘要包含重写了多少个表达式、裁剪了多少个Exchange、向量化了多少批数据。这是调优时最有价值的信息来源——你能精确知道每个优化点贡献了多少。4. 高频场景实测从日期计算到合并、排序与聚合前面讲的都是机制这一章是硬核实践。我挑了网上被问得最多的几类场景——日期加减与年月提取、DataFrame合并、排序、分组聚合逐一跑基准测试把Auron的收益用数字展示出来。测试环境是CDP私有云上的一个10节点集群每节点16核64G内存Spark 3.3.2数据量约2TB。4.1 日期加减与年月提取从Calendar调用变成整数算术日期处理是Spark SQL里最常见的性能陷阱。先看这个经典写法spark.sql( SELECT id, date_add(to_date(dt, yyyy-MM-dd), 365) as next_year, trunc(to_date(dt, yyyy-MM-dd), MM) as month_start FROM events )光看执行计划这里有一个隐藏的坑to_date(dt, yyyy-MM-dd)字符串解析是一个完整的状态机对每一行都要做格式校验、年月日拆解、日历计算。date_add和trunc又各自触发Calender运算整体算下来一行要调八九次纯Java时间API。Auron识别到这个模式后把它重写成spark.sql( SELECT id, cast(substring(dt, 1, 4) as int) 1 as next_year_int, substring(dt, 1, 7) as month_start_str FROM events )为什么能这么改因为dt格式固定是yyyy-MM-dd那么date_add(..., 365)等价于年份加一不跨闰年时trunc(..., MM)等价于直接截断到前7个字符。这些都是纯粹的整数和字符串操作没有任何时间状态机。实测结果20亿行的events表做这三种日期计算原始用时间243秒Auron重写后用了61秒加速比3.98倍。优化后几乎没有GCCPU跑得也更满。4.2 DataFrame合并空间换时间的两阶段JoinDataFrame合并Join是另一个高频操作。网上一搜“dataframe数据合并”教程全是小例子一到生产就露馅。我遇到的真实场景是一张12亿行的用户行为表和一张300万行的用户画像表做inner join拿用户标签。传统写法val result behavior.join(userProfile, user_id)这个Join有两个问题一是大表Shuffle无法避免虽然小表只有300万行但Spark默认要Shuffle才能Join除非用Broadcast提示二是Join之后行为表的所有列和画像表的全部列都留在内存里很多列后续根本用不到。Auron的优化是两层的。第一层是小表自动转Broadcast——300万行约800MB超出了Spark默认的10MB广播阈值但Auron知道这个集群Driver内存充足会把阈值放宽到1GB。第二层是列裁剪——Join之前就把画像表里不会用到的列在Scan阶段剔除行为表同理。最终Broadcast只传了50MB的瘦身小表Join直接本地完成。这个场景的实测结果原始耗时187秒优化后27秒加速比6.9倍。这里大头是省掉了Shuffle列裁剪贡献了大概30%的收益。需要说明的是Broadcast Join有内存风险如果你不确定自己Driver的能力别把阈值调太高——Auron会在运行前做一次小表尺寸估算超过阈值就不Broadcast但估算本身依赖统计信息建议先跑一次ANALYZE TABLE。4.3 排序、分组与窗口计算Auron的取舍策略排序和分组在Auron面前就比较微妙了。排序本身很难被“重写”——orderBy必须全量排序这是语义决定的。Auron能做的是在排序前尽量缩小数据量把select里用不到的列在Scan阶段砍掉把limit条件下推到排序前直接截断中间结果takeOrdered替代全排序。这些优化在数据行数大、列多时收益明显但纯排序性能本身没有提升。分组聚合则是另一回事。Auron的聚合优化主要是部分聚合的并行化。Spark的HashAggregate天然支持三阶段部分聚合、最终聚合、缓冲聚合但默认在最终聚合阶段所有数据都汇聚到一个分区容易倾斜。Auron加入了一个“中间聚合”层先在本地做多轮部分聚合把中间结果集压缩到很小然后再做最终聚合。这相当于把Reduce端的压力前移让Map端多干活。窗口函数OVER (PARTITION BY ... ORDER BY ...)是最难优化的场景之一。Auron的做法是去重掉无用的窗口排序如果ROWS BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW这种无界窗口里没有ORDER BY就不需要全局排序直接按分区顺序遍历即可。这个改写看人下菜碟但命中的时候性能提升很恐怖。实测分组聚合10亿行、100个分组Key、聚合sum和avg原始92秒优化后48秒加速比1.9倍。窗口函数场景row_number() over (partition by uid order by ts)原始160秒优化后113秒加速比1.4倍主要收益来自排序裁剪。整体来看Auron在“重计算、轻排序”的场景收益最大在“纯排序、无谓词”的场景收益有限。所以我自己上生产时会先跑spark.auron.stats.enabledtrue观察每个Stage的收益日志再决定哪些作业值得深度优化。5. 实战中踩过的坑与调优建议任何优化框架都有它的边界和暗坑。Auron用了大半年我踩过的坑不少挑几个最典型的分享出来希望你避开。5.1 表达式重写在什么情况下会静默失效最坑的不是报错而是重写规则不匹配但作业正常跑——你以为优化生效了实际没有。我遇到最多的几种失效场景类型精度不匹配。模式库里substring提取日期年月的规则要求字段类型是StringType且格式固定。如果数据里混入了几行脏数据比如2024-1-1这种少一位的格式规则会整体不匹配直接走原逻辑。这种情况不算错误但是会让你误以为优化失效了。复杂条件分支。when(col(dt) 2024-01-01, col(a)).otherwise(col(b))这类CaseWhen里嵌套日期解析模式库经常会命中不了因为解析路径被分支打散了。Auron对这种情况会退化为常量折叠收益小很多。自定义UDF的无注解替换。Auron的UDF替换需要业务方主动注册映射关系如果你没注册它不会猜测UDF语义。排查这个问题的办法就一句话看Auron日志。spark.auron.stats.logLevelINFO会把每个作业实际重写多少表达式打印出来。如果一个SQL明显属于模式库覆盖范围但重写计数为0大概率是类型或格式匹配出了问题。5.2 数据倾斜场景下Auron与AQE的配合数据倾斜是Spark世界永恒的痛Auron不能直接解决它但在某些情况下会让它更明显。举个例子Shuffle裁剪去掉了一次多余的Exchange但如果Join本身Key分布不均裁剪后倾斜直接体现在后续的HashAggregate上——因为没有中间的Exchange做一次数据打散。我建议在这样的作业里把Auron和Spark 3.x的AQE一起用起来。顺序上AQE是在运行阶段做动态调整Auron是在计划阶段做静态裁剪两者不冲突。实际操作中如果发现某个作业在启用Auron后倾斜加剧把spark.auron.rewrite.shuffle关掉针对这个作业跑一遍看看是不是裁剪导致的。还有一个更细的点Auron对聚合的“中间聚合”优化在面对极端倾斜的Key时会导致那个Key的处理Task压力更大。Auron的应对是增加一个倾斜检测如果部分聚合阶段某个分区的数据量超过中位数的5倍自动把该分区数据再拆成更细的粒度做两层预聚合。这个机制需要一点额外内存不必担心开销一般可接受。5.3 哪些场景不建议上Auron最后说点不该用的情况。我也不是所有作业都建议开Auron以下几种场景建议你保持默认关闭任务本身耗时很短秒级。Auron的优化分析本身有计算开销虽然通常小于0.5%但秒级任务可能占掉10%。大炮打蚊子不划算。依赖严格计算顺序的流式作业。Auron的谓词下推和投影裁剪会改变物理执行顺序虽然语义等价但流式场景下中间状态一致性验证成本较高风险大于收益。代码里大量使用动态类型和反射调用的UDF。重写规则踢不进去徒增分析开销。建议这类作业先把UDF改写成明确返回类型的函数再考虑接入。数据量只有几GB、跑在本地教学环境。说实话几GB的数据量优化感知不明显还容易掩盖基础操作的问题。先把代码写干净、把Shuffle规律摸清楚远比依赖框架重要。按我自己的经验Auron最适合的场景是大规模批处理、日级ETL、报表任务这类作业有明显可优化的执行计划空间优化一次能持续受益。而对于在线查询、交互式分析它作用没那么大不如把精力放在SQL写法本身。5.4 最后分享一个来自生产环境的小技巧如果你用的是Spark 3.3以上版本建议把Auron和spark.sql.adaptive.coalescePartitions.enabled一起打开。AQE会在Shuffle结束后自动合并小分区减少后续Stage的任务数而Auron裁剪掉多余的Shuffle后剩下的Shuffle数据量更集中合并效果会更好。我有个报表任务两者叠加后Executor的CPU利用率从31%提到了67%集群整体吞吐量上了一个台阶。这个配合思路是Auron文档里没写的属于我自己摸索出来的经验实测稳定推荐一试。
返回列表