ARTICLE DETAIL

资讯详情

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

Spark 3.0升级实战:AQE、DPP与性能优化深度解析

Spark 3.0升级实战:AQE、DPP与性能优化深度解析 先说结论Spark 3.0 是我用过的几个大版本里升级性价比最高的一次。如果你还在 2.4.x 上纠结要不要动我的建议很直接——如果你每天被数据倾斜、小文件、手动调参折磨3.0 的 AQE自适应查询执行就是为你准备的如果你被生产环境的资源浪费和诡异的 OOM 搞得焦头烂额那这次升级更是值得认真对待。这篇文章我尽量不写官腔把我拆解 Spark 3.0 的源码、跑测试、上生产后沉淀下来的东西一次性讲清楚覆盖性能提升和功能增强两大主线适合正在评估升级的架构师也适合被 Spark 调优折磨的开发同学。我始终记得第一次在 Spark 2.4 集群上手动估算 shuffle 分区数时的心情——那感觉就像闭着眼睛开手动挡换错档就得顿挫一下。3.0 最大的贡献不是某一个具体功能而是把“调参”这件脏活累活从人身上搬走了一部分。这个思路贯穿了 AQE、动态分区裁剪、动态优化等几乎所有核心改进。下面我从这几个维度一个个拆开说。1. 性能提升的核心引擎自适应查询执行AQE1.1 AQE 到底在解决什么实际问题先讲一个我自己的故事。去年维护过一套基于 Spark 2.4 的离线数仓任务每天几百个顺跑任务其中大约有两成会在下午的某个时段翻车。翻车原因高度集中在三类Shuffle 过程中数据倾斜导致单个 task 跑几个小时某个 reducer 因为数据量远超预期直接 OOM以及下游小文件爆炸把 HDFS NameNode 内存吃出一身冷汗。问题的根源听着很朴素Spark 在生成执行计划时并不知道真实数据的分布。它只能按照你给的参数、统计信息、或者干脆靠猜提前定好 reducer 数量和 join 策略。而真实数据一旦跟预估偏差过大执行计划就废了。Spark 3.0 引入的 AQE本质上是把“先定计划再执行”这件事改成了“边执行边修正”。它会在 shuffle 阶段完成后把真实的 map 输出统计信息拿回来重新审视当前的执行计划自动调整后续阶段的分区数、join 策略、甚至给倾斜的 join 分配额外的 task。这样一来很多原本需要人工反复试参数才能解决的问题3.0 在运行时就替你处理掉了。1.2 AQE 的三个核心优化点逐个拆解我把 AQE 拆成三个子能力这仨在生产中能解决的问题完全不同别混为一谈。动态合并 Shuffle 分区这个优化针对的是“分区数量定错了”的场景。举个例子你按照经验给某个任务设置了 2000 个 shuffle 分区结果上游数据实际只有 200GB 级别每个分区处理的数据才 100MB根本用不着 2000 个 task。在这种情况下2.4 只能按 2000 个分区老老实实跑浪费大量调度开销和 task 启动时间。而 3.0 在 shuffle 写完之后会获取所有 map 输出的分区大小如果发现相邻分区都比较小就会自动把它们合并起来减少下游 task 的数量。合并规则默认是通过spark.sql.adaptive.coalescePartitions.minPartitionNum和spark.sql.adaptive.coalescePartitions.initialPartitionNum控制的。你可以理解为AQE 会计算一个“目标分区大小”比如默认 64MB然后从前往后依次累加相邻分区直到超过这个阈值就划成一个新分区。注意它是顺序合并而不是重分布所以不会引入额外的 shuffle。我在生产环境见过最夸张的一个案例某个任务设置 5000 个分区上游实际只有 30GB 数据开 AQE 动态合并后实际跑起来只有 900 多个分区任务耗时直接砍半。对于跑 3 个小时以上的大任务来说这种省下来的时间非常可观。动态切换 Join 策略这个优化针对的是“join 策略选错了”的场景。最常见的问题是小表广播阈值设得不够导致 Spark 没有选择 BroadcastHashJoin而是选了 SortMergeJoin。SortMergeJoin 要 shuffle 两边的全量数据对磁盘 IO 和网络的压力远高于广播 join。在 2.4 里如果你预估的“小表”实际上在某个时间点数据量暴涨那只能等着任务被拖死。3.0 的 AQE 会在 shuffle 执行完后重新精确评估每个 join 输入的统计信息如果发现其中一边的数据量已经明显小于spark.sql.autoBroadcastJoinThreshold就会自动把 SortMergeJoin 改成 BroadcastHashJoin。这个动作不需要重新提交任务只是在执行计划层面替换算子。这点对于生产环境的价值非常大。大促、月末结算、数据回刷这些场景表的数据量经常是波动的你没法为每一种情况设置完美的阈值。AQE 相当于给执行计划上了一层保险。动态优化倾斜 Join倾斜是 Spark 任务里最恶心的坑没有之一。以前处理倾斜全靠手工要么找出倾斜 key 加盐打散要么拆分成两个任务分别跑代码改起来又丑又脆弱。3.0 的 AQE 针对 join 的倾斜做了自动化处理当它检测到某个分区处理的数据量远超其他分区时会自动把倾斜的分区分成多个小分区并为这些小分区分别启动 task。也就是说原本一个 10 小时的 task会被拆成 10 个 1 小时的 task 并发跑整体进度就拉回来了。触发倾斜检测的条件是通过spark.sql.adaptive.skewJoin.enabled控制的生产上我一般还会配合调整spark.sql.adaptive.skewJoin.skewedPartitionFactor和spark.sql.adaptive.skewJoin.skewedPartitionThresholdInBytes。简单说如果你发现某个 key 的数据量是正常分区的 5 倍以上同时绝对数据量超过某个阈值就会被识别为倾斜分区。这两个参数需要根据你的数据实际情况来调默认值偏保守但千万不要盲目调低阈值——拆得太细会导致小 task 过多调度开销反而盖过收益。1.3 AQE 的开关与版本差异说完功能提个关键的落地细节Spark 3.0 刚发布时AQE 并不是默认开启的需要显式设置spark.sql.adaptive.enabledtrue。如果你用的 3.0.0/3.0.1不开的话以上所有优化都不会生效。真正把 AQE 变成默认选项是从 Spark 3.2.0 开始的。我在迁移时踩过一个坑把 2.4 的作业直接切到 3.0没开 AQE跑出来的性能不但没提升某些场景反而变慢了因为 3.0 默认动态分区裁剪是开启的它会在某些情况下引入额外的计划优化开销。建议你在切换大版本后先用小范围内的一批核心任务做对比测试再决定全局配置的基线。2. 动态分区裁剪DPP让 SQL 少读数据的黑科技2.1 DPP 的原理与适用场景动态分区裁剪Dynamic Partition PruningDPP在 Spark 3.0 里从实验特性转正成为默认开启的功能。它解决的核心问题是星型模型里大事实表 JOIN 小维度表时能不能提前把事实表里根本不会被命中的分区裁剪掉。举个例子一张订单事实表按日期分区另一张小维度表只有最近 30 天的有效订单状态。在 2.4 里Spark 只会用静态分区裁剪也就是 WHERE 条件里明确写的分区过滤。如果 WHERE 条件没写日期那么 join 时事实表的所有分区都要读一遍。DPP 的思路是在 join 的另一个输入上做一次预计算收集维度表中真正会被用到的日期集合然后用这个集合去裁剪事实表的分区从而大大减少扫描的数据量。在 3.0 里默认spark.sql.optimizer.dynamicPartitionPruning.enabledtrue聚合之后的小表会自动被用于裁剪大表的分区不需要你改 SQL。这个特性对分区数量多、单分区数据量大的表效果极其显著。我在一个 ODS 层的任务上实测过事实表 1800 个分区join 之后只读了 47 个分区扫描量降了 97%任务从 40 分钟压到 11 分钟代价仅仅是多了一次对维度表的小规模聚合。2.2 DPP 的边界和限制需要注意DPP 听着很美但它不是万能的有几个在生产里容易踩的点。第一DPP 依赖统计信息所以要求表必须执行过ANALYZE TABLE。如果你没有收集统计信息的习惯DPP 很可能不会生效。我们当时就把spark.sql.statistics.histogram.enabledtrue打开并且每周定时对核心表跑一次 analyze。第二DPP 在以下这些场景是无效的非等值 join、用 OR 连接的过滤条件、join 条件不是分区键而是分区键上的表达式、子查询里的动态裁剪。如果你发现 DPP 没生效先看看执行计划里是否出现了DynamicPruningSubquery这个节点。没有的话基本就是条件不满足。第三小表如果本身很大预聚合的代价会抵消裁剪收益。当维度表超过广播阈值时DPP 可能会退化成 SortMergeJoin这时你要留意任务是否因此变成了两个大 shuffle。2.3 实操中要不要人工干预我的经验是DPP 的收益通常来源于“查询本身没有写分区过滤条件”的场景比如报表系统里用户选了日期范围但日期字段在另一张表里。遇到这种查询开 DPP 后基本不需要人工干预执行计划自己会处理。但要注意DPP 在 3.0 里对某些很复杂的 join 条件依然无能为力比如事实表的分区键在 join 条件里被函数包装了date_sub(a.dt, 1) b.dt这种情况下构建动态裁剪条件太复杂会被优化器直接放弃。如果你有这种 SQL建议要么提前预计算要么手动加上分区过滤条件别把希望全压在 DPP 上。另外提醒一句DPP 开启后偶发会出现“动态裁剪子查询”读取小表时额外多跑了一个 stage这在调度资源紧张、并发高的场景下可能引入额外的队列等待。你可以通过spark.sql.optimizer.dynamicPartitionPruning.reuseBroadcastOnly控制是否只复用已经广播的表数据避免重复计算。3. 功能增强一ANSI SQL 兼容与新的 SQL 语义3.1 为什么 ANSI 兼容对生产这么重要Spark 3.0 的一个重要更新是引入了 ANSI SQL 合规模式这直接影响了很多线上 SQL 的跑批行为。以前 Spark SQL 在类型转换和算术运算上非常“宽容”甚至可以说是放纵整数除法直接截断小数、字符串转数字失败返回 NULL、超长字符串截断也不报错。这带来一个隐蔽问题数据质量在跑批过程中悄悄丢了但任务本身是绿的没人发现。开启 ANSI 模式spark.sql.ansi.enabledtrue后这些操作会变成运行时错误任务直接失败而不是静默返回 NULL。初看这似乎是个“坏消息”因为线上可能会有不少旧 SQL 因为这种行为差异而失败。但从数据工程的角度这是天大的好事——它逼着你尽早暴露数据问题而不是等报表上线或者数仓数据已经跑歪了才发现。我们迁移过程中把 ANSI 模式开放到了一个单独的测试环境用一个月的全量任务做回归结果揪出了 6 个以前被静默吞掉的类型转换问题。这些问题都是真实数据格式异常导致的只不过以前被 Spark 的宽松策略掩盖了。3.2 CAST 行为变更与隐式转换的坑这次升级里最容易让老任务“莫名其妙挂掉”的就是 CAST 的语义变化。在 3.0 里默认情况下类似CAST(abc AS INT)仍然会返回 NULL但当你开启 ANSI 模式后它会直接抛出异常。这个变化对 ETL 脚本的影响是致命的——如果你的清洗逻辑依赖“转换失败就得到 NULL然后在下游用 COALESCE 处理”那么 ANSI 模式下任务会直接中断。举个实际例子我们的埋点日志表里有个字段duration历史数据里有少量脏值形如1.5秒以前CAST(duration AS INT)返回 NULL下游再用IFNULL兜底。开了 ANSI 模式后第一步 CAST 就报错任务直接红。我的建议是对这类清洗逻辑建议在开启 ANSI 模式之前先显式用try_cast3.0 提供了try_cast函数转换失败返回 NULL替代裸 CAST保证兜底逻辑仍然可控。还有一个隐蔽的变化是隐式类型转换。3.0 收紧了一些隐式转换规则比如 INT 到 LONG 的隐式转换在某些表达式里会失败。你如果写的 SQL 里大量依赖“字符串和数字比较”以前 Spark 会把字符串隐式转成数字现在可能会直接报错要求你显式 CAST。这类问题比较隐蔽建议通过EXPLAIN和ANALYZE两个阶段来排查。3.3 新的 SQL 函数和语法功能功能增强部分3.0 给 SQL 函数库补了一批新函数我挑几个在生产中管用的说。try_cast上面提到了转换失败返回 NULL 而不是抛异常适合做数据清洗兜底。bit_count、bit_and、bit_or、bit_xor位运算函数处理权限标记、开关状态之类的数据非常方便省得你写 UDF。date_part从日期时间字段里提取特定部分语法和 PostgreSQL 对齐做时间维度拆解很顺手。any_value从分组里随机取一个值以前要用first或max兜底现在直接有原生函数了注意它的语义是“任意值”不要用它来替代有序取值。另外3.0 对GROUP BY 别名、HAVING中使用别名这些语法的支持更完整了写复杂汇总逻辑时能省不少事。我个人还比较喜欢的是DEFAULT列值支持建表时指定默认值这在实时写入的场景可以少写很多判断逻辑。4. 功能增强二UDAF 新接口与 DataSource API v2 演进4.1 UDAF 新接口从继承到组合的转变Spark 3.0 对自定义聚合函数UDAF的 API 做了重构引入了新的Aggregator接口同时把旧的UserDefinedAggregateFunction标记为 deprecated。如果你在 2.x 时代写过 UDAF那么这里需要留意新接口的使用方式完全不同。新接口的逻辑是这样的你定义一个输入类型、一个缓冲类型、一个输出类型然后实现reduce更新缓冲区、merge合并两个缓冲区、finish输出最终结果三个方法。缓冲区的序列化由 Spark 自动处理你只需要定义好BufferEncoder。我举个例子统计一天内每条用户路径的访问次数并去重。旧接口你要继承一个类覆写一堆initialize/update/evaluate方法还要自己管理中间状态的 schema。新接口只需声明输入是某个 case class缓冲区是另一个 case class剩下交给框架。这个改进对复杂统计逻辑的可维护性提升非常明显代码量减少一半逻辑也更清晰。有一点要专门提醒新接口对缓冲类型的序列化是有要求的必须是可编码Encodable的类型。如果你在缓冲区里放了一个自定义的复杂 Java 对象编译能过运行时大概率在 shuffle 阶段报not supported错误。解决方法是把复杂对象拆分成基本类型和集合或者手动指定表达式编码器ExpressionEncoder。4.2 DataSource API v2为流批一体铺路DataSource API v2 在 3.0 中被进一步推进和稳定化。这个改动对普通业务开发者来说感知不强但对数据源提供商和自研存储系统接入的人来说是一个里程碑。v2 API 最大的变化是数据源不再只是一个简单的RelationProvider而变成了有着完整生命周期管理、支持分区裁剪、支持列裁剪、支持谓词下推、能感知流式写入的完整数据源抽象。它把“读”和“写”拆成了独立的接口并且引入了Table/TableProvider/ScanBuilder/WriteBuilder这套工厂模式。说句大白话以前你想把自研的存储系统接进 Spark得 hack 一堆内部接口实现很粗糙很多优化没法生效。现在按照 v2 接口规范实现一套Table和ScanSpark 就能自动支持谓词下推和列裁剪查询性能会有一个质的提升。我在实际项目里将一个内部自研的列式存储接入 Spark 3.0 时最明显的感受是代码结构清晰了很多。v2 的读写 API 是完全分离的数据源的提交commit机制更规范支持分阶段写入这让“精确一次”的语义在流批一体场景下成为可能。如果你所在团队有计划做实时数仓或者在调研自研存储的 Spark 接入3.0 的 DataSource v2 是值得投入的方向。4.3 流处理能力增强从微批到连续处理3.0 还增强了 Structured Streaming引入了Continuous Processing模式实验性。跟传统的微批模式不同连续处理模式通过spark.sql.streaming.continuous.enabledtrue配合format(continuous)使用理论上可以把端到端延迟降到毫秒级。4.2 性能细节谁在默默帮你省钱3.0 还做了很多不显眼但影响深远的性能优化比如shuffle 性能的改进引入更快的序列化器和更高效的 IO 路径。ZSTD 压缩3.0 默认支持 ZSTD 压缩格式对高压缩比场景比如日志类数据非常友好。动态分区写入优化写动态分区时3.0 对分区裁剪和文件提交的路径做了优化小文件问题有所缓解。同时3.0 在S3对象存储的适配、Kubernetes的原生支持spark.kubernetes.allocation.driver.self.scheduler等上也做了大量工作。如果你已经在容器化跑 Spark 作业3.0 的 Kubernetes 支持会明显改善扩展性和资源利用。5. 生产级监控与运维增强不只看日志5.1 新的监控指标与 Prometheus 支持3.0 的监控体系有了明显的补强Prometheus 支持是其中我最看重的一块。以前 Spark 的监控主要靠 Spark UI 和 Spark History Server但生产环境往往需要一套统一的监控告警体系UI 里的指标拉不出来运维就只能靠日志猜测。3.0 增加了 Prometheus 监控端点spark.ui.prometheus.enabledtrue可以直接暴露 executor 的 JVM 指标、shuffle 读写指标、task 执行指标等接入 Prometheus Grafana 之后集群的整体健康度和作业性能变化一目了然。我把这个端点加到现有的监控体系里之后那次“某任务 Shuffle 读耗时明显上涨”的问题就是靠指标曲线的拐点提前发现磁盘故障的。5.2 结构化日志告别 grep 大海捞针除了监控指标3.0 还引入了结构化日志Structured Logging。这个功能把日志从纯文本变成了 JSON 格式的 key-value 事件官方框架里的核心日志都带上了timestamp、level、source、message这些结构化字段。生产环境里这个功能的价值在于给日志加上了上下文信息排查问题的时候不再是grep一堆原始字符串然后靠肉眼对照而是可以按任务 ID、阶段 ID 去聚合查询。你可以通过spark.sql.adaptive.logLevel控制 AQE 运行日志的级别默认是info后台跑完一批任务后用日志聚合平台把AdaptiveSparkPlan相关日志捞出来基本就能还原整个执行计划的调整过程。5.3 关于升级的一个忠告最后说一个必须强调的坑Spark 3.0 清理了一批在 2.x 时代就已经废弃的 API 和配置项。你从 2.4 升级到 3.0要么直接用兼容模式启动spark.sql.legacy.*这些参数可以保留旧行为要么就提前规划一个“兼容性改造缓冲期”。比如spark.sql.hive.convertMetastoreParquet、spark.sql.sources.partitionOverwriteMode这些参数默认行为变了一旦你跑“整表覆盖写”或者“动态分区覆盖写”的任务不仔细核对很容易发生数据污染事故。我建议在测试环境跑一遍全量任务清单用 3.0 的执行计划对比 2.4 的执行计划重点看 Shuffle 次数、Join 类型、分区裁剪是否生效再决定全局升级的节奏。我个人在实际操作中的体会是Spark 3.0 的升级最忌讳“一刀切”。先把 AQE 打开用一周时间观察核心任务的时间和资源曲线再切换 ANSI 模式用回归测试兜底最后再把监控和告警体系配上。按这个节奏走大概率不会出大问题。
返回列表