性能优化 · 大数据运维

待环境验证

Spark SQL 执行计划与性能诊断

结合扫描、Join、Shuffle、Task 长尾和 AQE 最终计划定位 SQL 瓶颈,用固定输入验证结果一致性与资源收益。

Apache Spark大数据
阅读导引 · 理解后再操作

这篇知识解决什么问题

本篇以固定输入和结果语义建立性能基线,从 Scan、Join、Shuffle 与 Task 长尾收窄瓶颈,再用 AQE 最终计划验证优化是否生效。每次实验都应解释改善来自哪里,并保留资源成本和结果一致性的证据。

先比较相同输入再评价速度

SQL 修订、数据分区、缓存和资源窗口都会影响耗时。基线与候选应保持可解释的比较条件,先确认结果一致再评价性能,不能把最好一次运行作为唯一结论。

计划选择需要运行指标支持

提交前 explain 有助于理解候选执行路径,AQE 可能在运行时改写它。证明广播或倾斜优化生效,应查看已完成查询的最终计划和实际数据量。

进入文章正文
关联架构图解6 个组件 · 点击展开

从 Spark 作业图中跟踪数据源、Executor 计算、Shuffle 和结果写入,识别耗时所在阶段。图不展示查询的真实物理计划、分区数量或动态资源变化,需要将具体 SQL execution 的最终计划与指标叠加理解。

查看场景架构与实施步骤
参考架构 · 非实时拓扑

Spark 作业诊断与流式恢复运维架构

从 Driver、Executor、Stage 与 SQL 计划定位 Spark 作业瓶颈,结合事件日志和 checkpoint 核验批处理与流式恢复边界。

  • 数据 / 请求
  • 控制 / 管理
  • 观测 / 查询

点击组件,在图下方查看职责;连线编号对应流向解读。小屏可横向滚动,或直接展开文字说明。

Spark 作业诊断与流式恢复运维架构:组件关系图从 Driver、Executor、Stage 与 SQL 计划定位 Spark 作业瓶颈,结合事件日志和 checkpoint 核验批处理与流式恢复边界。 Driver → 资源管理器:资源申请;资源管理器 → Executor:资源分配;Driver → Executor:Task 调度;Executor → 输入与输出存储:读取与提交;Driver → 事件日志与历史服务:事件持久化与读取;Driver → 流式持久恢复点:流式进度协调;Executor → 流式持久恢复点:有状态计算持久化。箭头说明见下方流向解读。
逻辑参考图,待环境验证。采用 Spark 3.5.7 的批处理与 Structured Streaming 微批模式。YARN、Kubernetes 和 Standalone 的资源管理不同,部署平台状态与 Spark 作业状态需分别检查。

Driver

控制 / 治理

构造作业与阶段并调度任务,需要独立观察 Driver 内存。

查看关联工具
全部组件职责 6 个组件
资源管理器
按所选部署平台分配 Executor 资源,三者不是同时必需。
Driver
构造作业与阶段并调度任务,需要独立观察 Driver 内存。工具介绍 Driver
Executor
执行计算任务,任务重试和丢失 Executor 应与 Stage 时间线对应。工具介绍 Executor
输入与输出存储
保存业务输入输出;外部提交语义与算子重算独立核验。
事件日志与历史服务
event log 提供已结束作业的可观测来源,History Server 不调度任务。工具介绍 事件日志与历史服务
流式持久恢复点
仅 Structured Streaming 分支使用,每个查询采用独立持久目录;包含进度及有状态查询所需状态,不能与业务输出或 event log 混用。
流向解读 7 条连接
  1. 1

    Driver 资源管理器

    控制 / 管理 · 资源申请

    Driver 与部署层协作获得执行资源,具体生命周期取决于运行模式。

  2. 2

    资源管理器 Executor

    控制 / 管理 · 资源分配

    资源管理器提供运行 Executor 的资源。

  3. 3

    Driver Executor

    控制 / 管理 · Task 调度

    Driver 调度 task 并跟踪执行与重试。

  4. 4

    Executor 输入与输出存储

    数据 / 请求 · 读取与提交

    任务访问业务输入输出;重算安全性需看目标提交协议。

  5. 5

    Driver 事件日志与历史服务

    观测 / 查询 · 事件持久化与读取

    Driver 写 event log,History Server 读取;此节点合并表示日志与历史查看服务。

  6. 6

    Driver 流式持久恢复点

    控制 / 管理 · 流式进度协调

    Driver 协调微批进度和提交记录,恢复时需与同一查询及输入输出语义相容。

  7. 7

    Executor 流式持久恢复点

    数据 / 请求 · 有状态计算持久化

    执行侧保存有状态查询所需状态;状态后端与检查点存储插件须支持当前部署。

故障域与操作边界

适用版本与部署模式

采用 Spark 3.5.7 的批处理与 Structured Streaming 微批模式。YARN、Kubernetes 和 Standalone 的资源管理不同,部署平台状态与 Spark 作业状态需分别检查。

数据与变更边界

不通过反复重提作业、清空 checkpoint、全量 collect 或扩大 Driver 内存代替诊断。批任务重跑与流式恢复必须考虑已提交输出和外部副作用。

架构依据与版本核对 2 篇官方资料

图解是本站基于官方资料整理的逻辑参考;实施前仍需核对实际部署版本、组件支持范围与变更审批。

适用范围与问题定义

本文以 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 架构与执行过程

参考资料

从现象到判断

先收集证据,再缩小范围。以下是判读路径,不代表已经确认根因或获准变更。

  1. SQL 带有日期条件,但扫描文件与输入字节仍远超预计的日分区规模。

    只读核对
    检查计划中的分区过滤、列类型、函数转换与数据源读取路径。
    如何判读
    谓词存在不等于已裁剪;优先确认实际输入范围,不能先增大并行度掩盖扫描。
  2. Join Stage 只有少量任务持续拖尾,整体资源利用率逐渐下降。

    只读核对
    对照任务读写量、热点键样本和最终 Join 计划,区分节点异常与数据分布。
    如何判读
    长尾可能来自倾斜;需要保持 Join 语义的验证,不可随意盐化或删去热点键。
  3. 优化后出现更多 spill 或 Executor 丢失,最终日志显示 FetchFailed。

    只读核对
    保存首次容器退出原因、GC、分区规模、磁盘与节点证据,区分前因和连锁错误。
    如何判读
    FetchFailed 可能是上游丢失的后果;单条最后错误不能证明唯一内存根因。
常见误区与判断边界 2 项

将 LIMIT 当作扫描预算

聚合和排序可能需要读取更多输入后才能返回少量结果。诊断查询应通过已知输入分区限制范围,不能仅凭 LIMIT 声称开销有界。

给所有关联强制广播

压缩文件大小不代表广播对象内存上界,Join 类型和策略支持也有限制。应先核实过滤后规模与运行计划,强制 Hint 不是通用性能修复。

交接时应留下的证据

作为记录提纲使用,不是自动检查结果;未取得的证据应标记缺口,并注明负责人。

  • 保存 SQL 与配置修订、输入窗口、文件与行数、资源分配和缓存状态,形成可比较基线。
  • 归档 Scan、Join、Exchange 与 AQE 最终计划,以及最慢阶段和 Task 长尾证据。
  • 比较记录数、重复键、空值、业务汇总和允许误差,再记录耗时、Shuffle、spill 与资源成本。
  • 交接候选停止条件、原 SQL 与配置恢复方式,以及已写出结果需要单独处理的范围。

记录需包含环境、版本、时间与时区;分享前脱敏,不附访问令牌、密码或完整业务敏感数据。

继续阅读与资料核对

补充相关主题,再结合当前环境的实施记录形成结论。

返回原理导读

DOUYA OPS ECOSYSTEM

贡献你的经验,帮助更多运维人

把故障复盘、标准流程和最佳实践沉淀为可检索、可复用的知识内容。