ARTICLE DETAIL

建站实战干货

来自一线的建站与推广经验沉淀,每一条都经过真实交付验证。

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

2026/8/16 4:21:07 拓冰建站 浏览量
【Spark内核】Spark Driver 的整体架构 一、先给一个总判断Spark Driver 本质上是一个应用级控制中心。Spark Driver 的架构可以记成四步接入用户逻辑 → 分析依赖并拆 Stage → 调度 Task 并下发给 Executor → 收集反馈并推进或恢复执行再细一点就是SparkContext初始化和统领整个应用DAGScheduler按依赖拆 StageTaskScheduler把 Stage 变成 Task 并调度SchedulerBackend和集群通信把任务送出去状态监听组件记录执行过程、结果和失败他们之间怎么协作先由SparkContext把环境搭起来再由DAGScheduler看清计算依赖然后TaskScheduler把可执行任务安排出去最后SchedulerBackend通过集群把任务发给 ExecutorExecutor 执行完每个Stage以后再把结果和状态回传给 DriverDriver 再决定下一步。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 应该拆成几个 StageStage 之间的先后顺序是什么所以它管的是“阶段怎么切”不是“任务发给谁”。你可以把它理解成负责把一条完整的计算链拆成可执行的阶段链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 执行完以后还要把结果和状态再回传给 DriverDriver 再根据反馈决定下一步怎么走。所以真正的关系其实是一个闭环Driver 生成计划 → 拆分阶段 → 下发任务 → Executor 执行 → 状态回传 → Driver 推进后续计算2. 再看职责边界更准确地说Driver 里的各个部分分工是这样的SparkContext负责“启动和统领”。DAGScheduler负责“按依赖拆 Stage”。TaskScheduler负责“把 Stage 变成 Task并安排执行”。SchedulerBackend负责“把 Task 送到集群里”。监听和状态组件负责“记录过程让 Driver 知道发生了什么”。这几个部分不是并列堆在一起的而是层层衔接的。四、一次完整执行时它们怎么配合1. 用户先写计算逻辑用户写 DataFrame 或 SQL 的时候Spark 通常不会马上执行。比如valresultspark.read.parquet(/orders).filter($statusPAID).groupBy($user_id).sum(amount)这时 Driver 先接住的是“计算描述”不是最终结果。2. Action 到来后才真正开始执行当用户调用count、collect、write这类 Action 时Spark 才会把前面的计算描述变成真正的 Job。也就是说Driver 不是一开始就把所有东西都跑起来而是等到“需要结果”的那一刻才正式调度。3.DAGScheduler先拆 StageDriver 拿到 Job 之后DAGScheduler会先看依赖关系。如果某些计算之间是连续的就可以放在同一个 Stage 里如果中间遇到 Shuffle就必须切开。所以它做的第一件事是判断依赖 → 找 Shuffle 边界 → 拆出多个 Stage4.TaskScheduler再把 Stage 拆成 TaskStage 一旦确定Driver 就会把它拆成多个 Task。一个分区通常对应一个 Task。然后TaskScheduler负责把这些 Task 安排给合适的 Executor。所以这里的关系是Stage ↓ Task ↓ Executor 执行5.SchedulerBackend负责真正下发TaskScheduler决定“谁来跑”SchedulerBackend决定“怎么发出去”。它把任务通过集群通信机制送到 ExecutorExecutor 再真正开始干活。6. Executor 回报Driver 继续推进Executor 执行完成后会把结果、错误信息、Shuffle 输出位置等信息返回 Driver。Driver 收到反馈以后会做两类判断如果当前 Stage 成功了就推进下一个 Stage如果失败了就按失败类型决定重试还是重算所以 Driver 不是一次性发完任务就结束而是边收反馈边推进。