【Spark内核】Spark Driver 的整体架构
一、先给一个总判断
Spark Driver 本质上是一个应用级控制中心。
Spark Driver 的架构,可以记成四步:
接入用户逻辑 → 分析依赖并拆 Stage → 调度 Task 并下发给 Executor → 收集反馈并推进或恢复执行再细一点就是
SparkContext:初始化和统领整个应用DAGScheduler:按依赖拆 StageTaskScheduler:把 Stage 变成 Task 并调度SchedulerBackend:和集群通信,把任务送出去- 状态监听组件:记录执行过程、结果和失败
他们之间怎么协作
先由SparkContext把环境搭起来,再由DAGScheduler看清计算依赖,然后TaskScheduler把可执行任务安排出去,最后SchedulerBackend通过集群把任务发给 Executor,Executor 执行完(每个Stage)以后再把结果和状态回传给 Driver,Driver 再决定下一步。
Spark Driver 可以先理解成 Spark 应用的“控制中枢”。
它不直接承担大规模数据计算,而是负责把用户写下来的计算逻辑,逐层变成可以在集群里执行的任务,并在执行过程中持续跟踪状态、处理失败、推进后续计算。
所以,Driver 的核心不是“算”,而是“组织计算”。
二、Driver 里最重要的几个部分
1.SparkSession和SparkContext
这两层是 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 到来后,才真正开始执行
当用户调用count、collect、write这类 Action 时,Spark 才会把前面的计算描述变成真正的 Job。
也就是说,Driver 不是一开始就把所有东西都跑起来,而是等到“需要结果”的那一刻才正式调度。
3.DAGScheduler先拆 Stage
Driver 拿到 Job 之后,DAGScheduler会先看依赖关系。
如果某些计算之间是连续的,就可以放在同一个 Stage 里;如果中间遇到 Shuffle,就必须切开。
所以它做的第一件事是:
判断依赖 → 找 Shuffle 边界 → 拆出多个 Stage4.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 不是一次性发完任务就结束,而是边收反馈边推进。