【Spark内核】Spark Driver 的整体架构

【Spark内核】Spark Driver 的整体架构

一、先给一个总判断

Spark Driver 本质上是一个应用级控制中心。

Spark Driver 的架构,可以记成四步:

接入用户逻辑 → 分析依赖并拆 Stage → 调度 Task 并下发给 Executor → 收集反馈并推进或恢复执行

再细一点就是

  • SparkContext:初始化和统领整个应用
  • DAGScheduler:按依赖拆 Stage
  • TaskScheduler:把 Stage 变成 Task 并调度
  • SchedulerBackend:和集群通信,把任务送出去
  • 状态监听组件:记录执行过程、结果和失败


他们之间怎么协作

先由SparkContext把环境搭起来,再由DAGScheduler看清计算依赖,然后TaskScheduler把可执行任务安排出去,最后SchedulerBackend通过集群把任务发给 Executor,Executor 执行完(每个Stage)以后再把结果和状态回传给 Driver,Driver 再决定下一步。

Spark Driver 可以先理解成 Spark 应用的“控制中枢”。

它不直接承担大规模数据计算,而是负责把用户写下来的计算逻辑,逐层变成可以在集群里执行的任务,并在执行过程中持续跟踪状态、处理失败、推进后续计算。

所以,Driver 的核心不是“算”,而是“组织计算”。


二、Driver 里最重要的几个部分

1.SparkSessionSparkContext

这两层是 Driver 的入口。

用户写 Spark 程序,最先接触的一般是SparkSession,而SparkSession背后真正连接 Spark 核心运行时的是SparkContext

它们主要负责:

  • 接住用户提交的 Spark 应用
  • 初始化 Driver 的运行环境
  • 建立后续调度所需的核心对象
  • 作为用户代码和 Spark 内核之间的入口

可以把它理解成:

用户代码 ↓ SparkSession ↓ SparkContext ↓ Driver 内部调度体系

如果没有这一层,后面的计划生成、任务调度、失败恢复都无从谈起。


2.DAGScheduler

这是 Driver 里最核心的“依赖分析器”。

它的作用是把用户的计算逻辑,按照数据依赖关系拆成多个 Stage。

它主要回答的是:

  • 哪些计算可以连续做
  • 哪些地方因为 Shuffle 必须切开
  • 一个 Job 应该拆成几个 Stage
  • Stage 之间的先后顺序是什么

所以它管的是“阶段怎么切”,不是“任务发给谁”。

你可以把它理解成:

负责把一条完整的计算链,拆成可执行的阶段链

3.TaskScheduler

这是 Driver 里的“任务调度器”。

如果说DAGScheduler管的是 Stage,那么TaskScheduler管的就是 Stage 里具体的 Task。

它的作用是:

  • 接收某个 Stage 生成的一批 Task
  • 决定哪些 Task 先跑
  • 决定 Task 发到哪个 Executor 上
  • 处理任务失败后的重试
  • 根据资源、本地性和调度策略做分配

你可以把它理解成:

负责把一个阶段里的具体任务安排出去

它不负责分析依赖,只负责把可以执行的任务真正调度起来。


4.SchedulerBackend

这是 Driver 和集群资源之间的连接层。

它的作用是:

  • 向集群管理器申请 Executor
  • 接收 Executor 注册
  • 维护 Driver 和 Executor 的通信
  • 把 Task 真正发送到 Executor
  • 接收资源变化和执行状态反馈

它本身不决定业务逻辑,也不负责拆 Stage,它更像一个“适配器”

不同部署环境下,比如 YARN、Kubernetes、Standalone,底层实现不一样,但 Driver 看见的是统一的调度接口。


5. 状态和监听组件

Driver 里还有一组容易被忽略,但非常重要的组件,它们负责记录和展示执行过程。

典型的有:

  • 事件监听
  • Spark UI 状态更新
  • Job、Stage、Task 的状态跟踪
  • 指标和日志收集
  • Shuffle 输出和任务结果的元数据维护

这些组件不直接参与“怎么计算”,但它们决定了 Driver 能不能知道:

  • 现在执行到哪一步了
  • 哪个 Stage 成功了
  • 哪个 Task 失败了
  • 为什么失败
  • 后面该不该重试

它们负责的是“看见”和“记住”。


三、这些部分之间怎么串起来

1. 先看最核心的链路

Driver 内部最核心的关系,大致是这样的:

用户代码 ↓ SparkSession / SparkContext ↓ DAGScheduler ↓ TaskScheduler ↓ SchedulerBackend ↓ Executor

但这不是一条单向流水线。

因为 Executor 执行完以后,还要把结果和状态再回传给 Driver,Driver 再根据反馈决定下一步怎么走。

所以真正的关系其实是一个闭环:

Driver 生成计划 → 拆分阶段 → 下发任务 → Executor 执行 → 状态回传 → Driver 推进后续计算

2. 再看职责边界

更准确地说,Driver 里的各个部分分工是这样的:

SparkContext

负责“启动和统领”。

DAGScheduler

负责“按依赖拆 Stage”。

TaskScheduler

负责“把 Stage 变成 Task,并安排执行”。

SchedulerBackend

负责“把 Task 送到集群里”。

监听和状态组件

负责“记录过程,让 Driver 知道发生了什么”。

这几个部分不是并列堆在一起的,而是层层衔接的。


四、一次完整执行时,它们怎么配合

1. 用户先写计算逻辑

用户写 DataFrame 或 SQL 的时候,Spark 通常不会马上执行。

比如:

valresult=spark.read.parquet("/orders").filter($"status"==="PAID").groupBy($"user_id").sum("amount")

这时 Driver 先接住的是“计算描述”,不是最终结果。


2. Action 到来后,才真正开始执行

当用户调用countcollectwrite这类 Action 时,Spark 才会把前面的计算描述变成真正的 Job。

也就是说,Driver 不是一开始就把所有东西都跑起来,而是等到“需要结果”的那一刻才正式调度。


3.DAGScheduler先拆 Stage

Driver 拿到 Job 之后,DAGScheduler会先看依赖关系。

如果某些计算之间是连续的,就可以放在同一个 Stage 里;如果中间遇到 Shuffle,就必须切开。

所以它做的第一件事是:

判断依赖 → 找 Shuffle 边界 → 拆出多个 Stage

4.TaskScheduler再把 Stage 拆成 Task

Stage 一旦确定,Driver 就会把它拆成多个 Task。

一个分区通常对应一个 Task。

然后TaskScheduler负责把这些 Task 安排给合适的 Executor。

所以这里的关系是:

Stage ↓ Task ↓ Executor 执行

5.SchedulerBackend负责真正下发

TaskScheduler决定“谁来跑”,SchedulerBackend决定“怎么发出去”。

它把任务通过集群通信机制送到 Executor,Executor 再真正开始干活。


6. Executor 回报,Driver 继续推进

Executor 执行完成后,会把结果、错误信息、Shuffle 输出位置等信息返回 Driver。

Driver 收到反馈以后,会做两类判断:

  • 如果当前 Stage 成功了,就推进下一个 Stage
  • 如果失败了,就按失败类型决定重试还是重算

所以 Driver 不是一次性发完任务就结束,而是边收反馈边推进。