别再熬夜翻DAG图了!Spark任务优化,其实可以交给AI Agent

1.如何通过Spark Web UI定位任务问题?

在之前的文章中,我们提到了一些通过 Spark Web UI 定位任务问题的方式:

下面,我们再进一步详细讲解一下这部分内容,这对后面讲解 Spark 任务优化至关重要。

1.1 进入Spark Web UI界面

在日常的开发工作中,我们总会遇到 Spark 应用运行失败、或是执行效率未达预期的情况。对于这些问题都可以通过 Spark UI 来获取最直接、最直观的线索,在全面地审查 Spark 应用的同时,迅速定位问题所在。如果我们把失败的、或是执行低效的 Spark 程序看作是“病人”的话,那么 Spark UI 中关于应用的众多度量指标(Metrics),就是这个病人的“体检报告”。结合多样的 Metrics,身为“大夫”的开发者即可结合经验来迅速地定位“病灶”。

官网:https://archive.apache.org/dist/spark/docs/3.4.4/web-ui.html#sql-tab

打开 Spark UI,最上面的导航条,这里罗列着 Spark UI 所有的一级入口(备注:如果是非SparkSQL的程序,将不会有SQL一级入口),如下图所示。

1.2 点击Stage,查看Stage整体的执行情况

我们知道,每一个作业可能包含多个Stage,在 Stages 页面,Spark UI 罗列了应用中涉及的所有 Stages,这些 Stages 分属于不同的作业。要想查看哪些 Stages 隶属于哪个 Job,还需要从 Jobs 的 Descriptions 二级入口进入查看

Stages 页面,更多地是一种预览,要想查看每一个 Stage 的详情,同样需要从“Description”进入 Stage 详情页。

  • Input:指真正读取的文件大小,如果表是分区表,则代表读取的分区文件大小。如果数据表有10个字段,只select了3个字段并发生了列裁剪,则Input表明是3个字段的存储大小。
  • Output:输出到HDFS上的文件大小,如果结果数据是压缩的,则代表压缩后的大小。
  • Shuffle Write:为了Shuffle所准备的数据,未来会有其他的Stage来读取,该部分数据会写到磁盘上。
  • Shuffle Read:Shuffle阶段读取的数据大小,既包含Executor本地的数据,也包含从远程Executor读取的数据。

某些Stage除了会显示总的Task数,执行成功Task数之外,还会显示failed task数。failed task数量就代表该Stage中执行失败的Task数量。是因为Spark有Task级别的重试来保证容错。

spark.task.maxFailures代表一个task连续执行失败几次会被中止,默认设置为4

这个时候,我们需要在Stage页面找出执行时间异常的Stage,去进一步定位问题。

1.3 查看单个异常Stage执行情况

点击Stage对应的Description,进入到详情页。我们先来看看Stage的详情页包含哪些信息?(详细说明可以去

https://archive.apache.org/dist/spark/docs/3.4.4/web-ui.html#sql-tab 看Stage detail)

官网介绍

重点关注

这里需要关注两个核心指标:

  • Shuffle Read Size / Records:如果某个 Stage 读取的数据量(Shuffle Read)远大于其他 Stage,说明上游 Stage 产生了数据膨胀,可能存在 数据倾斜。
  • Tasks 指标:关注 Tasks 列表中的 Duration 列。
  • 现象:大部分 Task 在几秒内完成,但有少数几个 Task 耗时极长(几分钟甚至几小时)。
  • 结论:典型的 数据倾斜。
  • 定位:点击该 Stage 进入详情,查看 Shuffle Read Size 列,通常耗时长的 Task 读取的数据量是其他 Task 的几十倍甚至几百倍。可以记录下 Host 地址,结合 Executors 页面查看该节点是否异常。

1.4 查看Executor的运行情况

Executor选项卡介绍

“Executor”选项卡显示了为应用程序创建的执行器的摘要信息,包括内存和磁盘使用情况以及任务和 shuffle 信息。“存储内存”列显示了用于缓存数据的已用和保留内存量。

“Executor”选项卡不仅提供资源信息(每个执行器使用的内存、磁盘和核心数量),还提供性能信息(GC 时间和 shuffle 信息)。

单击Executor 0 的 ‘stderr’ 链接,可在其控制台中查看详细的标准错误日志。

Executor问题定位

重点关注

  1. 失败与死亡节点
  • 关注点:Dead 列表。如果有 Executor 挂掉(Dead),任务就会在另一个节点重试。
  • 定位:如果任务一直失败或极慢,发现有 Executor 频繁死亡,点击 Logs 链接查看 stderr 日志。常见原因:OOM(内存溢出)、FetchFailedException(网络或磁盘问题)。
  1. 活跃节点负载不均

    0 结论:可能该节点所在的物理机资源被抢占,或者存在 数据本地化 问题(数据在远端,拉取耗时过长)。

  • 关注点:Tasks 列(任务数)和 Duration 列(总执行时间)。
  • 现象:某个 Executor 执行的 Task 数量特别少,或者总耗时特别长。
  1. 内存与 GC
  • 关注点:Storage Memory(存储内存)和 Shuffle Write/Read。
  • 现象:如果某个 Executor 的 GC Time 异常高(例如超过总任务时间的 10%),说明内存压力过大,导致频繁垃圾回收,严重拖慢速度。

1.5 定位异常SQL

SQL详情页介绍

先介绍一下SQL详情页,以下面SQL为例:

SELECT count(DISTINCT if(server_id = 1, user_id, null)) server_1 ,count(DISTINCT if(server_id = 2, user_id, null)) server_2 ,count(DISTINCT if(server_id = 3, user_id, null)) server_3 ,count(DISTINCT if(server_id = 4, user_id, null)) server_4FROM ods_game_dev.ods_user_login

在 SQL Tab 一级入口,我们看到有 1个条目

点击图中的“Decription”,即可进入到该作业的执行计划页面,如下图所示。

每个方块都代表了一种算子,鼠标在算子的色块上悬停,下方会显示该节点的详细信息:

  • Duration:该节点总耗时(毫秒)。
  • Records:输入/输出行数。
  • Data Size:输入/输出数据量(字节)。
  • Shuffle Read/Write:对于 Exchange 节点,显示 Shuffle 读写的记录数和数据量。
  • Peak Memory、Spill 等:用于判断内存压力。

SQL逻辑定位

第一步:获取 Stage 的唯一标识

  • 在 Stages 页面,每个 Stage 都有一个 Stage ID(例如 stage 5)。
  • 记下这个 Stage ID,以及它所归属的 Job ID(如果需要)。

第二步:进入 SQL 页面,找到对应的查询

  1. 点击顶部导航栏的 SQL 标签。
  2. 页面会列出所有已执行的 SQL 查询(包括 DataFrame 操作),每个查询都有 Description 和 Duration。
  3. 根据 Stage 归属的 Job 时间或查询描述,找到最可能包含该 Stage 的查询,点击其 Description 进入详情页。
  4. 提示:如果 Stage 归属的 Job 执行了多个 SQL,可以在 Jobs 页面查看 Job 的 SQL 列表(通过 Job 详情中的 SQL ID)。

第三步:在 DAG 可视化图中定位 Stage(★★★★★)

SQL 详情页的 DAG Visualization 展示了该 SQL 的物理执行计划。

  1. 通过节点标注查找
  • 算子类型(如 Scan、Exchange、HashAggregate、SortMergeJoin)

  • Stage ID(例如 Stage 5)

  • DAG 图中的每个矩形节点通常包含:

  • 在 Spark 3.x 中,节点上方会直接显示 Stage Id。如果未直接显示,可以将鼠标悬停或点击节点,在弹出的详情框中会显示该节点所属的 Stage ID。

  1. 通过节点列表查找
  • 详情页右侧或下方有一个 节点列表,按执行顺序列出所有物理算子。
  • 每个算子条目也会标注 Stage ID,可以快速定位目标 Stage。
  1. 确认节点类型
  • 定位到目标 Stage 后,观察该节点的算子类型。常见的物理算子与 SQL 逻辑的对应关系如下:

第四步:关联到具体的 SQL 代码片段

  1. 利用节点详情中的表达式
  • 点击 DAG 中的节点,下方会显示该节点的 详细信息,包括输入/输出表达式、过滤条件、聚合函数等。
  • 例如,一个 HashAggregate 节点会列出聚合函数和分组字段,比如 keys: [user_id],说明 SQL 中存在 GROUP BY user_id。
  1. 查看物理计划文本
  • 在 SQL 详情页,可以找到 Details 或 Physical Plan 按钮,点击后显示完整的物理计划文本。
  • 在物理计划中搜索 Stage ID,可以找到对应的算子以及它包含的表达式,这些表达式直接反映了 SQL 中的逻辑。
  1. 结合 SQL 原始文本
  • 如果 SQL 是纯文本执行的,可以在 SQL 页面的查询描述中看到原始 SQL(可能被截断)。
  • 将物理计划中的表达式与 SQL 文本对照,即可确定 Stage 对应的是哪部分逻辑。
  • 例如:
  • 物理计划中出现 Exchange 且下游是 SortMergeJoin → 对应 SQL 中的 JOIN 操作。
  • 出现 HashAggregate 且分组字段是 date → 对应 SQL 中的 GROUP BY date。
  • 出现 BroadcastExchange → 对应 SQL 中触发了广播 join 的表。

示例:

1.6 常见问题场景与对应 UI 特征

2.异常任务优化思路

Spark 任务的异常优化是一个从资源到代码逻辑的逐层深入过程。当任务出现慢、失败或不稳定时,建议按照以下四个维度依次排查与优化:资源问题 → 并发配置 → 数据倾斜 → 异常问题(如 HDFS Shuffle 慢节点)

资源问题

很多时候,任务产出慢,可能是由资源问题导致的!

  • 一般来说,资源的层级是这样看的:公司集群规模 => 部门可用集群 => 资源队列额度 => 任务优先级
  • 队列资源打满的情况下,即使任务优先级很高,也可能导致产出延迟;任务优先级很低的情况下,即使队列资源没有打满,任务也可能执行的很慢
  • 另外,不同的任务类型所提交的队列是不同的,并且每个队列的资源都是有限的:比如数据查询、线上任务、补数据(回溯数据)任务所在的队列就是不同的,配额也不同
  • 一般来说:线上任务队列资源 > 数据查询队列资源 >= 补数据队列资源
  • 队列资源其实就是可供分配的内存+CPU核数

贴一个思路导览:

如果资源不紧张,但是任务上仍然存在资源问题,可以通过增加 spark.executor.memory,spark.executor.cores,来让任务在执行时申请到更多的 Executor 资源。

怎么定位任务资源上存在问题?

常见表现

  • Task 频繁 GC,GC Time 占比超过 10%–20%。
  • Executor 频繁 OOM 或被 YARN/K8s 杀掉。
  • 单个 Executor 处理数据量远超其内存,导致大量溢写磁盘(Spill)。
  • 任务整体吞吐量低,CPU 利用率不足。

排查方法

  • 在 Spark Web UI 的 Executors 页面查看:
  • Storage Memory:是否远小于配置的内存。
  • Shuffle Write/Read 与内存对比,是否存在大量溢写。
  • GC Time 与任务总时间的比例。
  • Dead Executors 及对应的日志,查找 OOM 或 FetchFailed 异常。

并发配置

常见表现

  • Stage 中 Task 数量极少(例如几十个),但每个 Task 处理数据量极大,执行时间很长。
  • 或者 Task 数量过多(数万甚至数十万),每个 Task 处理数据量极小(几 KB),调度开销巨大。
  • CPU 使用率低,但任务长时间处于“Pending”状态。

排查方法

  • 在 Stages 页面查看 Stage 的 Number of Tasks 以及每个 Task 的 Shuffle Read/Write 数据量。

并发问题的解决思路通常是通过参数配置来解决,这块后面单独出专题讲解一下Spark的参数配置。

3.现阶段,如何利用 AI Agent 提效 Spark 的任务优化?

3.1 数仓同学优化任务的现实困境

认知困境:看到问题,抓不住关键

当一个任务跑崩时,Spark UI 已经给出了大量信息 — 上百个 Stage、数千个 Task、密密麻麻的 DAG 图,数仓同学需要:

  • 在几百个 Stage 中翻页筛选
  • 在密集的 DAG 图中追踪数据流向
  • 区分哪些是正常节点、哪些是冗余计算

结果:信息过载导致“看得到问题,却抓不住关键”。即使是有经验的工程师,也要花费大量时间才能从海量信息中定位到真正的瓶颈。

时间困境:有时间时没需求,有需求时没时间

复杂 SQL 的完整优化闭环包括:理解业务逻辑 → 分析执行计划 → 定位瓶颈 → 改写 SQL → 验证等价性 → 上线观察。

这个闭环通常需要 1 到 2 天的整块专注时间。

但现实是:

  • 业务需求排满日程,性能优化永远被挤到“有空再说”
  • 当任务真正跑崩、业务方催促时,压力最大,反而最没有时间从容优化
  • 对于维护几十上百个定时任务的同学来说,每个任务都做一次深度优化,成本根本不可接受

结果:优化成了“救火式”的被动响应,而非主动治理。

信任困境:知道问题在哪,不敢动手去改

即使定位到问题,并构思出改写方案,验证环节同样令人却步:

  • 需要在小数据量上反复验证结果等价性
  • 数据规模大,试跑成本高
  • 手工校验几乎不可行,稍有不慎就可能引入数据质量问题

结果:很多优化想法停留在“想改但不敢改”的状态。最终只能选择加资源、加并发、加超时,用“堆机器”的方式绕过问题,而不是真正解决问题

这也解释了为什么越来越多的团队开始探索AI Agent 介入优化流程:不是要取代人,而是要把人从“翻 DAG 图、对比执行计划、手工改写 SQL”的低效重复劳动中解放出来,让人聚焦于业务判断和策略选择。

3.2 AI Agent 怎么解决这个问题?怎么为任务优化提效?

Spark 任务优化 Agent工作流概览:

环节一:任务发现与元数据采集

Agent 做什么

  • 通过调度平台 API 定时拉取团队成员的所有 HSQL 任务,筛选出运行时长超过设定阈值的高耗时任务
  • 对每个筛选出的任务,调用 Spark History Server API 获取最近一次执行的 Stage 级指标(executorRunTime、shuffle 读写量、Task 数量)以及物理执行计划文本

提效点

  • 自动覆盖所有任务,无需人工逐一翻阅调度平台
  • 将分散在 UI 各处的指标统一整理为结构化数据,为后续分析提供高质量输入

环节二:异常识别与瓶颈分析(AI核心价值点,核心提效点)

Agent 做什么

  • 按 executorRunTime 对 Stage 排序,自动锁定 Top N 瓶颈 Stage
  • 结合物理执行计划,将每个 Stage 映射到具体的 SQL 操作(全表扫描、Join、聚合),识别数据量大但过滤率高、重复扫描、低效 Join 顺序等反模式
  • 计算每个 Stage 内 Task 执行时间的 p95 与 p50 比值,检测是否存在数据倾斜,区分“数据量大”与“数据倾斜”两类瓶颈

提效点

  • 从上百个 Stage 中自动定位真正的瓶颈,不再依赖人工凭经验猜测
  • 提供根因分析,为后续方案选择提供准确依据

环节三:方案生成(参数调优 / SQL 改写)(AI核心价值点,核心提效点)

Agent 做什么

  • 根据瓶颈类型,从优化知识库中匹配解决方案(可以人工维护一些优化技巧或者思路或者参数配置参考,便于AI学习),并针对当前任务生成具体建议
  • 参数调优:如调整广播阈值、shuffle 分区数、开启自适应查询执行等
  • SQL 改写:如将重复扫描统一物化、拆分大表 Join 为两阶段、优化 Join 顺序、合并多次独立扫描

提效点

  • 输出不再只是“问题定位”,而是“具体怎么改”的可执行方案
  • 量化预期收益,帮助用户优先落地高价值优化

环节四:自动测试与数据验证

Agent 做什么

  • 准备测试环境:自动创建测试表或在隔离队列中准备测试数据
  • 执行基线:运行原始 SQL,记录执行指标和结果集
  • 执行优化方案:依次运行每个优化后的 SQL,记录相同维度的执行指标
  • 数据验证:采用多层递进验证(行数对比、关键字段哈希对比、全字段聚合对比、随机抽样对比),确保优化前后结果集完全一致,差异方案自动标记为不可用

提效点

  • 彻底消除手工跑数验证的低效与高风险
  • 通过多层校验确保数据准确性,为上线提供信心

环节五:对比报告与用户确认

Agent 做什么

  • 生成结构化测试对比报告,包含:
  • 任务基本信息与瓶颈摘要
  • 各优化方案的具体修改点
  • 数据验证结果(通过/不通过)
  • 性能对比(耗时、扫描量、Stage 数、资源消耗等)
  • 推荐操作及上线前注意事项
  • 通过消息渠道将报告推送给任务负责人,并附带确认入口(采纳/拒绝/稍后处理)

提效点

  • 自动生成专业报告,减少人工整理数据的时间
  • 将决策与执行分离,提升协作效率

环节六:自动化部署上线

Agent 做什么

  • 根据用户确认的方案,自动执行上线操作:
  • 参数调整:通过调度平台 API 更新任务配置
  • SQL 替换:提交代码到 Git 仓库或直接更新调度平台中的 SQL 内容
  • 上线验证:触发一次正式环境运行,监控任务状态,对比执行结果与测试报告预期
  • 异常处理:若任务失败或出现异常,自动回滚并通知用户
  • 记录上线结果,用于后续优化效果追踪和知识库积累

提效点

  • 实现从优化建议到生产生效的全流程自动化
  • 自动回滚机制降低变更风险,保障生产稳定性

3.3 人工 VS AI Agent

总结对比

学AI大模型的正确顺序,千万不要搞错了

🤔2026年AI风口已来!各行各业的AI渗透肉眼可见,超多公司要么转型做AI相关产品,要么高薪挖AI技术人才,机遇直接摆在眼前!

有往AI方向发展,或者本身有后端编程基础的朋友,直接冲AI大模型应用开发转岗超合适!

就算暂时不打算转岗,了解大模型、RAG、Prompt、Agent这些热门概念,能上手做简单项目,也绝对是求职加分王🔋

📝给大家整理了超全最新的AI大模型应用开发学习清单和资料,手把手帮你快速入门!👇👇

学习路线:

✅大模型基础认知—大模型核心原理、发展历程、主流模型(GPT、文心一言等)特点解析
✅核心技术模块—RAG检索增强生成、Prompt工程实战、Agent智能体开发逻辑
✅开发基础能力—Python进阶、API接口调用、大模型开发框架(LangChain等)实操
✅应用场景开发—智能问答系统、企业知识库、AIGC内容生成工具、行业定制化大模型应用
✅项目落地流程—需求拆解、技术选型、模型调优、测试上线、运维迭代
✅面试求职冲刺—岗位JD解析、简历AI项目包装、高频面试题汇总、模拟面经

以上6大模块,看似清晰好上手,实则每个部分都有扎实的核心内容需要吃透!

我把大模型的学习全流程已经整理📚好了!抓住AI时代风口,轻松解锁职业新可能,希望大家都能把握机遇,实现薪资/职业跃迁~

这份完整版的大模型 AI 学习资料已经上传CSDN,朋友们如果需要可以微信扫描下方CSDN官方认证二维码免费领取【保证100%免费