适用范围与问题定义
本文以 Apache Spark 3.5.7 批处理 SQL 为基线,适用于扫描缓慢、Join 倾斜、Shuffle 开销和资源配置诊断。不同数据源连接器、表格式及发行版可能改变下推和执行行为;流任务恢复另见 Spark checkpoint 与恢复。所有示例待实际环境验证,本站状态为 PENDING。
先将“变慢”定义为同样数据窗口、相同结果语义与可比较资源条件下的耗时变化。保留 SQL 修订、应用标识、输入分区与大小、配置和运行时段。数据量翻倍、缓存失效、资源队列拥堵和代码回归是不同假设,不能仅以两次总耗时判断优化优劣。
先定位扫描、重分布还是长尾
从应用时间线确认慢发生在实际任务执行期间,再进入 SQL 与 Stage 页面。检查输入行数和字节、扫描文件数、Shuffle 读写、spill、GC、失败重试以及 Task 时长分布。优先看最慢阶段占比和最慢任务与中位任务的差距,不把全部累计 CPU 时间解释成用户等待时间。
如果大量任务都相近地变慢,考虑整体输入或共享资源;若只有少数任务持续拖尾,核对分区规模、键分布、节点和重试。不同版本或连接器暴露的指标可能不同,缺少某列要保留“未采集”,不要以零值替代。观测入口见 Spark 3.5.7 监控指南。
用固定分区 SQL 阅读物理计划
下面只解释一份已确认的日分区查询,执行前将表名与日期改成环境中的小范围样本。EXPLAIN 会规划查询及访问相关元数据,不应对未知超大文件清单无限重复:
EXPLAIN FORMATTED
SELECT customer_id, SUM(amount) AS total_amount
FROM REPLACE_WITH_CATALOG.REPLACE_WITH_DB.REPLACE_WITH_TABLE
WHERE business_date = DATE '2026-01-01'
GROUP BY customer_id;检查 Scan 的读取列、分区过滤及数据源过滤信息,再寻找 Exchange、聚合和 Join。日期条件写在 SQL 中并不证明已经分区裁剪;若列经过函数转换、类型不匹配或目录未按该列分区,实际扫描可能仍然很大。LIMIT 只限制返回行数,在聚合或排序中不能被当作扫描数据的硬上限。计划格式见 EXPLAIN。
先减少不必要输入与文件开销
让筛选尽早作用于可裁剪的分区和数据源列,避免不需要的 SELECT *,确认列类型与业务谓词一致。是否真正下推要以连接器文档和计划为准,不能只凭 SQL 外观判断。统计信息应与数据变更窗口一致;若估计值严重过时,应通过数据平台既定流程评估更新,而不是在排障时顺便全库收集统计。
小文件会增加发现和调度开销,合并文件应在独立候选数据上验证结果与查询收益。单纯减少输出文件数并不证明扫描量减少。HDFS 输入的治理流程见 Hadoop 小文件治理,Parquet 的分区发现与读写边界见 Spark Parquet 文档。
Join 选择需要数据量和语义证据
先确认关联条件、Join 类型、两侧过滤后大小及键分布。广播适合可安全放入相关进程内存的小侧,但文件压缩大小不是运行时对象的内存上界,强制广播可能造成超时或内存压力。Hint 是计划建议,还受策略是否适用于当前 Join 影响,不是无条件的执行承诺。
对长尾 Join,抽取已经批准的有限键分布统计,识别空值、默认键或热点业务实体。热点拆分必须保持内外连接、去重与聚合语义;随意增加随机盐值可能改变结果。先用原始查询和候选查询的有限样本做语义对照,再讨论吞吐收益。Join Hint 与 AQE 的能力范围见 SQL 调优文档。
正确理解 AQE 与分区参数
Spark 自 3.2.0 起默认启用 AQE,但现场可以覆盖;在 3.5.7 中,它能够基于运行时统计合并 Shuffle 分区、调整部分 Join 策略和处理受支持的倾斜 Join。应从当前会话和应用 Environment 页面确认有效配置,不从版本推断功能已生效。
SET spark.sql.adaptive.enabled;
SET spark.sql.shuffle.partitions;
SET spark.sql.autoBroadcastJoinThreshold;以上为读取配置,不修改设置。返回结果表示当前会话值,不证明其他应用使用相同配置。分区数太少可能形成大任务,太多则增加调度成本;AQE 的建议分区大小也不是内存使用上限。调参时一次只改一个假设所需的项,并保留最终计划,不能只展示提交前 explain 截图。
把 spill、GC 与内存错误分开
spill 说明部分执行数据落到磁盘,需要结合总量、磁盘速度与任务耗时判断成本;出现 spill 并不自动意味着作业失败或必须扩大内存。长 GC 可能与对象、缓存或分区规模有关。Executor 丢失还要查看容器退出原因、资源限制和节点日志,不能由最后一行 FetchFailed 倒推唯一根因。
先定位是 Driver、Executor 还是 Python 进程的资源压力。大规模 collect、广播构建和复杂计划可能影响 Driver;增加 Executor 内存对此不一定有效。对同节点重复失败应关联主机与磁盘证据,对各节点相同键失败则检查数据和算子,保留支持与反证。
组织一次可比较的优化实验
固定输入快照或关闭分区,记录基线与候选 SQL、有效配置、资源分配和运行窗口。先做结果一致性检查,再比较总耗时、最慢阶段、长尾、Shuffle、spill 和资源使用。业务允许的浮点误差需预先说明;排序、空值、重复记录及外连接结果都应覆盖。
缓存命中、并发负载、资源动态变化会影响对比,所以应将冷暖状态与环境波动写入报告。只有出现波动或新假设时才扩大实验次数,不能用大量重跑凑出最好的一次。写入型 SQL 的候选输出应进入独立目标,避免优化实验重复覆盖正式表。
常见误区、验收与回退
常见误区是先调大所有资源、给所有 Join 加广播 Hint,以及为了缩减文件统一使用单分区。它们可能把瓶颈移到 Driver、单个 Task 或共享队列,甚至改变其他应用可用资源。另一误区是只验收耗时,忽略结果行数、主键和业务汇总发生变化。
验收需要结果一致、代表性窗口改善、资源成本可接受且无新失败。若结果语义变化、长尾扩大或共享集群受影响,停止候选并恢复原 SQL 与配置。只读诊断无需数据回滚;候选已写出的数据和外部副作用必须单独处理,取消应用不会撤销它们。执行概念可回到 Spark 架构与执行过程。