架构设计 · 大数据运维

待环境验证

Spark 架构与执行过程:从 Driver 到 Task

通过有限分组样本理解 Driver、Executor、Stage、Shuffle 和 AQE,区分应用排队、计算与输出提交的耗时。

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

这篇知识解决什么问题

本篇用有限分组样本建立 Driver、Executor、Job、Stage 与 Task 的关系,再把应用慢拆成资源等待、计划、计算和输出提交。重点是读懂数据在哪里重分布,以及每一类失败应到哪一层寻找证据。

计算描述不等于已经执行

转换会建立计算依赖,action 才请求结果并推动实际任务。理解惰性执行后,才能区分计划成功、任务成功与输出提交成功,避免以 explain 成功代替运行验收。

Shuffle 需要同时考虑分区与资源

重新分布数据涉及网络、磁盘和任务调度。只有分区规模、热点和资源条件一起解释,才能判断增加 Executor 是否有用,不能将更多核数视为通用加速保证。

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

使用 Spark 作业运维图观察提交端、Driver、资源管理器、Executor 与存储的关系。图中的执行节点是逻辑聚合,不代表固定数量或当前分配;Driver 的实际位置、client/cluster 模式和日志入口需要从应用记录确认。

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

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 为基线,适用于理解批处理和 Spark SQL 应用的执行过程。YARN、Kubernetes 与 Standalone 的资源管理实现不同,本文介绍共同执行概念,并用 YARN 作为运维对照;Spark 4.x 和商业发行版的依赖、计划及配置应重新核对。示例未在实际集群运行,验证状态为 PENDING。

目标是解释 Driver、Executor、Job、Stage、Task 和 Shuffle 之间的关系,并把应用排队、计算变慢与结果写入失败分开。Spark 是计算引擎,数据可以来自 HDFS、对象存储或数据库;看到 Spark 进程并不能推断数据一定在内存中,也不能推断集群采用哪种存储。

Driver 与 Executor 各自做什么

Driver 运行应用主程序并协调任务,Executor 在集群节点中执行任务、存储缓存和 Shuffle 数据。资源管理器负责为应用提供执行资源,不负责替代 Driver 的 SQL 计划与任务调度。不同应用通常有各自的 Executor,不能把某应用缓存的 DataFrame 视为所有应用共享的数据缓存。

部署模式决定 Driver 位于提交客户端还是集群内部。client 模式下提交机器的网络和进程稳定性是应用依赖;cluster 模式下应去相应的集群日志系统找 Driver 输出。Executor 日志不能替代 Driver 异常记录。职责与部署模式依据 Spark 3.5.7 集群模式概览

从惰性计算理解 Job、Stage 与 Task

DataFrame 或 RDD 的转换通常先形成计算描述,真正需要结果的 action 才触发执行。一个应用可以产生多个 Job,每个 Job 由若干 Stage 组成,Stage 内的 Task 面向数据分区并行执行。应区分应用、SQL execution、Job 和 Stage 的标识,不能把同一个 SQL 中的多次 Job 误认为重复提交应用。

某些转换能在已有分区内继续计算,另一些操作需要重新分布数据。Shuffle 往往成为阶段边界并带来磁盘、网络和重试成本。多给 CPU 不一定能加速分区数量很少或数据严重倾斜的 Stage。惰性计算及 Shuffle 的概念见 RDD Programming Guide

用一千行样本观察计划与结果

下面是使用现有 SparkSession 的 PySpark 3.5.7 教学示例,只构造一千个整数,不读取业务表、不写入外部系统。执行 show 会提交实际计算,所以在学习环境中使用:

from pyspark.sql import functions as F

sample = spark.range(0, 1000, 1, 4)
summary = sample.groupBy((F.col("id") % 10).alias("bucket")).count()
summary.explain(mode="formatted")
summary.orderBy("bucket").show(10, truncate=False)

若示例成功,数学上应有十个 bucket,每组 count 为一百;这是预期,不是本文取得的实测输出。物理计划通常包含聚合与 Exchange 等节点,具体名字、层次和最终分区数会受版本和 AQE 配置影响,不应以计划文本逐字相同为验收标准。orderBy 也可能引入额外排序过程,因此不能据此例给复杂 SQL 估算固定 Stage 数。

SQL 计划与运行时计划要分开看

SQL 会经历解析、分析、优化和物理计划选择;字段能解析不代表数据一定可读,物理计划能生成也不代表运行一定成功。用 EXPLAIN FORMATTED 或 DataFrame explain 理解扫描、过滤、关联和聚合的位置,再对照运行中的 SQL 页面和实际指标。

Spark 3.5.7 的 AQE 可以利用运行时统计调整计划,因此提交前的解释结果不一定等于最终执行方案。若要证明倾斜优化或广播转换生效,查看已完成 SQL 的最终计划与指标。格式和模式定义见 EXPLAIN 文档,参数边界见 SQL 调优指南

用一条时间线定位应用慢在哪里

  • 提交到 Driver 启动:优先证据:调度器、AM 或 Pod 状态;典型候选原因:队列、权限、资源申请。
  • Driver 启动到首批 Task:优先证据:Driver 日志、文件发现与计划;典型候选原因:元数据访问、依赖加载、计划开销。
  • Task 执行期间:优先证据:Stage 分布、Executor 指标;典型候选原因:扫描、Shuffle、倾斜、GC、外部服务。
  • Task 完成到结果可见:优先证据:输出提交日志、目标系统状态;典型候选原因:提交协议、重试、目录或事务发布。

采样时保存 application ID、attempt、SQL execution ID 和时区。计算结束后的业务延迟可能来自下游目录刷新或索引,不能只按 Spark UI 的完成时间结案。YARN 中尚未取得资源的问题,接着使用 YARN 排队诊断

缓存、重算与失败重试的边界

缓存适合复用成本高、重复访问且资源允许的数据。首次访问仍需要填充缓存;资源不足或 Executor 丢失时,分区可能需要重新计算。缓存不等于备份,应用结束后的持久保存需要明确外部存储与提交结果。盲目缓存所有中间表会与执行内存竞争,应测量重复利用次数及缓存命中情况。

Task 重试意味着同一段用户代码可能再次执行,写数据库、调用接口等外部副作用需要单独保证幂等或事务语义。普通应用重启也不自动等于 Structured Streaming 从原 checkpoint 恢复;两者的输入位点与状态管理不同。流处理可继续阅读 Spark 流任务 checkpoint 与恢复

常见误区与观测准备

把所有错误都归因于 Executor 内存不足,会忽略 Driver 上的大规模 collect、文件列表和计划开销。collect 把结果带回 Driver,结果规模未知时不能拿它做生产数据预览。另一个误区是将更多 Executor 视为通用加速器;输入分区、倾斜、外部吞吐和资源队列都可能限制实际收益。

为事后分析配置受控的 event log 保存与 History Server,并确认日志路径权限、保留周期及敏感字段处理。History Server 依赖事件日志,不能保证每个未启用日志或日志缺失的应用都能完整回放。监控能力依据 Spark 3.5.7 Monitoring

验收、停止与回退边界

学习验收应能从样本计划指出一次数据重分布,解释预期分组结果,并把一个真实应用的等待、计算和提交阶段分开。运维验收还要保存数据窗口、有效配置、日志来源与阶段耗时,不把示例成功当作生产性能证明。后续可用 Spark SQL 性能诊断 形成一次可比较优化实验。

本文只读计划没有业务数据回退步骤;教学 show 只计算有限数据。若在真实应用调优中产生新输出,应写到独立候选目标再验收。取消或重跑应用之前,先检查已经提交的外部写入与幂等边界;进程停止不会撤销外部系统中已成功的写入。

参考资料

从现象到判断

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

  1. 提交端退出后应用停止,而集群中的部分 Executor 仍短暂存在。

    只读核对
    核对部署模式、Driver 所在位置、生命周期和首个异常日志。
    如何判读
    client 模式可能依赖提交机器;Executor 存在不能证明 Driver 仍能调度或应用可继续。
  2. SQL 已提交且计划生成,但首个 Task 启动前等待时间占大部分。

    只读核对
    分别查看资源分配、Driver 依赖初始化、文件发现和计划耗时。
    如何判读
    尚未执行算子时,任务内存或 Shuffle 参数未必针对当前瓶颈。
  3. 大多数 Task 已完成,少数任务长期运行,应用整体迟迟不能结束。

    只读核对
    比较同 Stage 的任务数据量、Shuffle、节点、重试和耗时分布。
    如何判读
    分区倾斜或节点差异可能形成长尾,需要进一步证据,不能仅看平均 CPU。
常见误区与判断边界 2 项

把缓存当作持久备份

缓存可能因资源回收或 Executor 丢失而重算,也不是不同应用的通用共享存储。业务需要长期保留的结果应有明确外部输出与提交验收。

通过 collect 预览未知大结果

collect 会把结果带回 Driver,未知结果规模可能改变原本低影响的诊断。使用预先固定的有限教学数据或明确边界的样本,并单独核对 Driver 资源。

交接时应留下的证据

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

  • 记录 Spark 与连接器版本、部署模式、Driver 位置、资源管理器和 application attempt。
  • 保留有限样本的数学预期、计划片段和实际学习环境结果,未运行的内容明确标为预期。
  • 保存提交、首任务、最慢 Stage 和输出可见的时间线,关联对应日志和执行标识。
  • 交接缓存、事件日志与外部写入的边界,注明取消或重跑不等于撤销已提交结果。

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

继续阅读与资料核对

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

返回原理导读

DOUYA OPS ECOSYSTEM

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

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