大数据运维 · 高级

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

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

Apache Spark大数据监控告警

场景目标

从 Driver、Executor、Stage 与 SQL 计划定位 Spark 作业瓶颈,结合事件日志和 checkpoint 核验批处理与流式恢复边界。 所有输出应绑定真实采样窗口和业务样本,交接未关闭问题与下次检查时间。

环境要求

准备 Application ID、attempt、制品和提交参数、event log 与 History Server 权限,记录数据快照和输出路径。REST 示例要求 SPARK_HISTORY_URL 指向已认证的只读 HTTPS 网关,SPARK_APP_ID 为具体应用;多 attempt 应在 API 路径追加已确认 attempt。

参考架构 · 非实时拓扑

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 流式持久恢复点

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

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

从架构到实施

  1. 01

    核对应用和资源分配

    应用 ID、提交时间、资源申请和实际分配可对应,明确批或流式作业。

  2. 02

    区分内存与存储压力

    失败能落到具体进程和资源维度,数据量及分区分布可解释。

  3. 03

    隔离重跑并对账输出

    业务键、总量、批次边界和失败重试输出无未解释差异。

故障域与操作边界

适用版本与部署模式

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

数据与变更边界

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

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

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

方案说明

适用架构

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

不适用边界

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

全局验收

以实际拓扑、业务样本与持续进度共同验收,记录时间窗口、变更前后证据及恢复缺口。本文是待环境验证的参考流程,预计时间只覆盖首轮诊断与小范围试点,不包括大规模重建或真实生产故障演练。

配套知识

官方参考

工具编排

1 个关联工具
  1. Apache Spark运行状态与诊断来源使用该产品原生状态、日志与计划获取诊断证据,具体方式见各步骤。

实施步骤

共 8 步
  1. 01

    核对应用和资源分配

    固定应用、attempt、Driver 运行位置和调度平台;区分提交未获资源、Driver 失败与任务执行缓慢。Spark UI 和 History Server 是观测入口,不负责给应用分配容器。

    curl --fail --silent --show-error --max-time 10 "${SPARK_HISTORY_URL:?set SPARK_HISTORY_URL}/api/v1/applications/${SPARK_APP_ID:?set SPARK_APP_ID}"
    验证标准

    应用 ID、提交时间、资源申请和实际分配可对应,明确批或流式作业。

    停止与回退

    应用身份或 attempt 不明确时暂停重提,先核对调度器记录和日志目录。

    返回步骤起点
  2. 02

    沿 Stage 找到长尾任务

    读取已确认 attempt 的 Stage 信息,对比 task 分位数、输入、shuffle、spill、GC 和失败重试;保留最慢任务与正常任务样本。Stage 重试不能简单累加为独立业务运行。

    curl --fail --silent --show-error --max-time 10 "${SPARK_HISTORY_URL:?set SPARK_HISTORY_URL}/api/v1/applications/${SPARK_APP_ID:?set SPARK_APP_ID}/stages"

    上面的列表用于定位 Stage;分位数需继续进入 SQL/Stage UI,或对已确认的 STAGE_ID 与 STAGE_ATTEMPT_ID 请求 taskSummary。下例沿用已核对的应用路径,多 application attempt 时追加正确 attempt。

    curl --fail --silent --show-error --max-time 10 "${SPARK_HISTORY_URL:?set SPARK_HISTORY_URL}/api/v1/applications/${SPARK_APP_ID:?set SPARK_APP_ID}/stages/${STAGE_ID:?set STAGE_ID}/${STAGE_ATTEMPT_ID:?set STAGE_ATTEMPT_ID}/taskSummary?quantiles=0.5,0.95,0.99"
    验证标准

    能区分普遍资源不足与少量倾斜 task,首个失败阶段及原因可追溯。

    停止与回退

    事件日志缺失则标记证据不全,不为补日志直接重新运行有外部副作用的任务。

    返回步骤起点
  3. 03

    结合 SQL 实际计划

    使用 EXPLAIN 与 SQL UI 分析裁剪、Join 策略、统计信息和实际运行的 AQE 计划。检查过滤是否下推、shuffle 是否过多和广播侧是否真正足够小;计划与运行数据规模必须一致。

    验证标准

    计划中的扫描、交换和 Join 与 Stage 指标相互验证,优化目标明确。

    停止与回退

    发现模型或结果语义变化时停止改写,保留原 SQL 与业务结果基线。

    返回步骤起点
  4. 04

    区分内存与存储压力

    分开 Driver、Executor 堆、off-heap、容器额外内存和 shuffle 临时磁盘。对照 GC、spill 和丢失 Executor 记录,检查 collect、缓存和分区尺寸;单纯增加内存可能延长 GC 或降低并发。

    验证标准

    失败能落到具体进程和资源维度,数据量及分区分布可解释。

    停止与回退

    资源边界未确认前不提高全局额度,临时盘满时停止扩量并保留 shuffle 故障证据。

    返回步骤起点
  5. 05

    在固定输入上单项优化

    选一个已定位的改动,如过滤顺序、分区策略、Join 或缓存生命周期;固定输入快照、并发和输出目标,进行有限样本对照。不用多个参数同时变化掩盖原因。

    验证标准

    数据结果一致,端到端时间、资源成本与失败率都有前后记录。

    停止与回退

    结果或资源恶化则恢复旧 SQL、提交参数和缓存策略,保留试验输出隔离。

    返回步骤起点
  6. 06

    检查流式进度与恢复点

    对流式作业采集 query progress 的 batch ID、输入和处理速率、水位线及状态行数,记录唯一 checkpoint 路径。输出提交、offset 与状态共同决定恢复行为,startingOffsets 不覆盖已有恢复进度。

    验证标准

    checkpoint 归属唯一,源保留期覆盖恢复点,慢批次可解释。

    停止与回退

    禁止两个查询共用 checkpoint 或删除目录强制从头;状态无法读取时停止升级。

    返回步骤起点
  7. 07

    隔离重跑并对账输出

    批任务使用独立输出路径验证重跑;流任务先核对查询结构、状态 schema 和 sink 兼容性再恢复。foreachBatch 需业务幂等或批次去重,不因恢复成功就声称外部写 exactly-once。

    验证标准

    业务键、总量、批次边界和失败重试输出无未解释差异。

    停止与回退

    目标结果不符时保留原输出和 checkpoint,停止生产切换,按副作用单独制定补偿。

    返回步骤起点
  8. 08

    归档基线与交接

    归档脱敏 event log、SQL 计划、输入快照、配置差异和业务对账。将排队时间与执行时间分开,记录下一次高峰检查和恢复预算;没有事件日志的运行不编造算子证据。

    验证标准

    接班人员能按应用 ID 还原瓶颈、改动与恢复范围,长期观察项有负责人。

    停止与回退

    验收不通过时冻结推广,恢复最近的可逆参数,不覆盖已经验证的输入或输出。

    返回步骤起点

DOUYA OPS ECOSYSTEM

体验豆芽自研工具与场景能力

部分场景提供体验环境,用于功能验证、测试和技术交流。