前些天发现了一个巨牛的人工智能学习网站,通俗易懂,风趣幽默,忍不住给大家分享一下。点击跳转到网站:https://www.captainai.net/dongkelun
前言
在第五篇文章 Yarn Client 模式提交流程(五) 中,我们分析了 Yarn Client 模式的完整提交流程。本文将继续分析 Spark 的另一种常用部署模式——Standalone Client 模式的提交流程。
Standalone 模式是 Spark 自带的集群管理器,不依赖 YARN、Mesos 等外部资源调度系统。其架构包含 Master(集群管理器)和 Worker(工作节点)两个角色,Master 负责资源调度,Worker 负责执行计算任务。
版本
Spark 3.2.3
Standalone Client 模式概述
Standalone Client 模式下,Driver 运行在提交应用的客户端机器上,向 Master 注册 Application,Master 调度 Worker 启动 Executor。
与 Yarn Client 的直观对比:
Yarn Client 模式中的 ResourceManager + NodeManager 对应到 Standalone 模式中的 Master + Worker。区别在于 Standalone 模式是 Spark 自带的简易调度系统,功能比 YARN 精简,但启动和调试更加方便。
| 模式 | Master 地址 | childMainClass | Driver 位置 |
|---|---|---|---|
| Standalone Client | spark://host:port |
args.mainClass(用户主类) |
客户端 |
| Standalone Cluster | spark://host:port |
SparkApplication 或 Wrapper |
集群 |
| Yarn Client | yarn |
args.mainClass(用户主类) |
客户端 |
1 | // 提交命令示例 |
完整提交流程
1. SparkSubmit 准备阶段
在第二篇文章 SparkSubmit 提交流程分析(二) 中,我们分析了 prepareSubmitEnvironment 方法。该方法根据 --master 和 --deploy-mode 设置 childMainClass,决定 main 方法在哪里执行。
对于 Standalone Client 模式,关键代码在 SparkSubmit.scala 的 prepareSubmitEnvironment 方法中:
1 | // SparkSubmit.scala — prepareSubmitEnvironment 方法 |
集群管理器识别后,进入 Standalone 模式的配置赋值逻辑。当 deployMode == CLIENT 时,设置 childMainClass:
1 | // SparkSubmit.scala — prepareSubmitEnvironment 方法(第 777-783 行) |
关键点:Standalone Client 和 Yarn Client 在这步完全一致——
childMainClass都设为args.mainClass,也都是在客户端通过runMain反射执行用户主类的 main 方法。Driver 在提交客户端运行。
整理后有以下对应关系:
| 模式 | childMainClass | main 方法执行位置 | Driver 位置 |
|---|---|---|---|
| Standalone Client | args.mainClass(用户主类) |
runMain 在客户端反射调用 |
客户端 |
| Standalone Cluster | SparkApplication 或 Wrapper |
runDriver 在 Worker 内调用 |
集群 |
2. runMain 执行用户主类
prepareSubmitEnvironment 返回后,runMain 方法通过反射加载并执行用户主类的 main 方法:
1 | // SparkSubmit.scala — runMain 方法 |
JavaMainApplication.start 内部通过反射调用用户主类的 main 方法:
1 | // SparkSubmit.scala |
用户主类 main 方法中的典型调用链:
1 | 用户主类.main |
3. SparkContext 初始化与 TaskScheduler 创建
SparkContext 初始化时,会创建 TaskScheduler 和 SchedulerBackend。在第四篇文章 SparkContext 初始化与 TaskScheduler 创建流程(四) 中已经详细分析了该流程。
当 master = "spark://..." 时,通过 createTaskScheduler 方法创建的是:
- TaskScheduler:
TaskSchedulerImpl - SchedulerBackend:
StandaloneSchedulerBackend
1 | // SparkContext.scala — createTaskScheduler 方法(简化) |
注意:这个模式匹配只是一个简化示意。实际源码使用了
ExternalClusterManager的 SPI 机制来创建,但最终效果如上所述。
创建完成后,SparkContext 初始化过程中会调用 _taskScheduler.start(),该调用会委托给 StandaloneSchedulerBackend.start(),正式进入提交流程的核心。
4. StandaloneSchedulerBackend.start() — 核心入口
StandaloneSchedulerBackend 继承自 CoarseGrainedSchedulerBackend,实现了 StandaloneAppClientListener 接口。它的 start() 方法是整个 Standalone Client 提交流程的核心入口:
1 | // StandaloneSchedulerBackend.scala |
4.1 super.start() — CoarseGrainedSchedulerBackend.start()
CoarseGrainedSchedulerBackend.start() 创建 DriverEndpoint,是 Driver 端与 Executor 通信的 RPC 端点:
1 | // CoarseGrainedSchedulerBackend.scala |
DriverEndpoint 负责处理 Executor 的注册(RegisterExecutor)、任务分发(LaunchTask)、状态更新(StatusUpdate)等核心 RPC 消息。
4.2 启动命令模板与 CoarseGrainedExecutorBackend
从代码中可以看到,Standalone 模式下的 Executor 启动入口类为:
1 | org.apache.spark.executor.CoarseGrainedExecutorBackend |
与 YARN 模式有所不同:
| 模式 | Executor 入口类 |
|---|---|
| Standalone | CoarseGrainedExecutorBackend |
| Yarn Client | YarnCoarseGrainedExecutorBackend(继承自 CoarseGrainedExecutorBackend) |
CoarseGrainedExecutorBackend 的启动参数使用了占位符,Worker 在真正启动 Executor 进程时会替换为实际值:
1 | val args = Seq( |
4.3 ApplicationDescription
ApplicationDescription 封装了应用的描述信息,包括应用名称、需要的 CPU/内存、Executor 启动命令等,Master 将根据这些信息进行资源调度:
1 | val appDesc = ApplicationDescription( |
5. StandaloneAppClient 注册流程
StandaloneAppClient 是 Driver 端与 Master 通信的客户端。创建后,调用 client.start() 触发注册流程:
1 | // StandaloneAppClient.scala |
5.1 ClientEndpoint.onStart()
setupEndpoint 注册 RpcEndpoint 时会自动触发 onStart() 方法,开始注册流程:
1 | // StandaloneAppClient.scala — ClientEndpoint 内部类 |
5.2 registerWithMaster() 与 tryRegisterAllMasters()
registerWithMaster() 调用 tryRegisterAllMasters() 向所有 Master 地址发送注册请求,并设置定时器处理重试:
1 | // StandaloneAppClient.scala — ClientEndpoint |
重试机制总结:
- 最大重试次数:3 次
- 重试间隔:20 秒
- 总超时:最多 60 秒(3 × 20s)
- 并发:同时向所有 Master 地址发起注册
tryRegisterAllMasters() 向所有 Master 地址发送 RegisterApplication 消息:
1 | // StandaloneAppClient.scala — ClientEndpoint |
注意:使用
send(异步、无需响应)而不是ask(同步、需响应)。Master 注册成功后,会通过RegisteredApplication消息回复。
5.3 注册成功回调
当 Master 回复 RegisteredApplication 消息后,ClientEndpoint.receive 处理该消息:
1 | // StandaloneAppClient.scala — ClientEndpoint.receive |
listener.connected() 对应 StandaloneSchedulerBackend.connected():
1 | // StandaloneSchedulerBackend.scala |
至此,StandaloneSchedulerBackend.start() 中的 waitForRegistration() 返回,注册完成。Driver 端继续执行后续逻辑。
注册流程时序:
1
2
3
4
5
6
7
8
9
10
11
12
13
14 StandaloneSchedulerBackend.start()
→ client.start() → setupEndpoint("AppClient", ClientEndpoint)
→ ClientEndpoint.onStart()
→ registerWithMaster(1)
→ tryRegisterAllMasters()
→ masterRef.send(RegisterApplication(appDesc, self))
[Master 处理注册,回复 RegisteredApplication]
→ receive: RegisteredApplication
→ listener.connected(appId)
→ registrationBarrier.release()
waitForRegistration() 返回 → 注册完成
6. Master 注册与调度
6.1 Master.receive — RegisterApplication
Master 收到 RegisterApplication 消息后,处理流程如下:
1 | // Master.scala — receive 方法 |
createApplication 方法创建 ApplicationInfo:
1 | // Master.scala |
registerApplication 将 Application 加入管理集合:
1 | // Master.scala |
6.2 schedule() 调度入口
schedule() 是 Master 的资源调度入口。每次有新应用注册、新 Worker 加入、资源变化时都会调用:
1 | // Master.scala |
调度优先级:
- Drivers 优先调度(Round-Robin 方式遍历 Worker)
- Executors 随后调度(通过
startExecutorsOnWorkers())
6.3 startExecutorsOnWorkers() — Executor 调度
这是 Standalone 模式下 Executor 资源分配的核心方法:
1 | // Master.scala |
核心步骤:
- 按 FIFO 顺序遍历
waitingApps队列 - 过滤可用 Worker:存活状态、资源充足
- 按空闲核心数降序排序:优先使用资源充足的 Worker
- 计算分配方案:通过
scheduleExecutorsOnWorkers()决定每个 Worker 分配多少核心 - 执行分配:通过
allocateWorkerResourceToExecutors()在 Worker 上启动 Executor
6.4 scheduleExecutorsOnWorkers() — 核心分配算法
这个方法实现具体的资源分配逻辑,支持两种模式:分散(spreadOut) 和 集中(非 spreadOut):
1 | // Master.scala |
spreadOut 算法说明:
- spreadOut = true(默认):将 Executor 尽可能分散到多个 Worker 上,有利于数据本地性
- spreadOut = false:尽可能让 Executor 集中到尽可能少的 Worker 上
- 可通过
spark.deploy.spreadOut配置(默认 true)
6.5 allocateWorkerResourceToExecutors() 与 launchExecutor()
分配确定后,allocateWorkerResourceToExecutors() 在具体 Worker 上启动 Executor:
1 | // Master.scala |
关键动作:
- Master 向 Worker 发送
LaunchExecutor消息,让 Worker 启动 Executor 进程- Master 向 AppClient 发送
ExecutorAdded通知,告知 Driver 有新的 Executor 被分配
7. Worker 启动 Executor
Worker 收到 Master 的 LaunchExecutor 消息后,创建 ExecutorRunner 并启动:
7.1 Worker.receive — LaunchExecutor
1 | // Worker.scala — receive 方法 |
7.2 ExecutorRunner.start() 与 fetchAndRunExecutor()
ExecutorRunner.start() 在新线程中启动 fetchAndRunExecutor(),真正执行 CoarseGrainedExecutorBackend 进程:
1 | // ExecutorRunner.scala |
fetchAndRunExecutor() 替换占位符并启动 CoarseGrainedExecutorBackend 进程:
1 | // ExecutorRunner.scala |
占位符替换方法:
1 | // ExecutorRunner.scala |
启动后的进程命令大致如下:
1 | /bin/java -server ... |
8. Executor 注册到 Driver
8.1 CoarseGrainedExecutorBackend 启动
CoarseGrainedExecutorBackend 进程启动后,执行 main() 方法:
1 | // CoarseGrainedExecutorBackend.scala |
8.2 onStart() — 反向注册
setupEndpoint 注册 CoarseGrainedExecutorBackend 到 RpcEnv 后,自动触发 onStart() 方法,开始向 Driver 反向注册:
1 | // CoarseGrainedExecutorBackend.scala |
8.3 DriverEndpoint 处理 RegisterExecutor
Driver 端的 DriverEndpoint 收到 RegisterExecutor 消息后:
1 | // CoarseGrainedSchedulerBackend.scala — DriverEndpoint |
8.4 创建 Executor 计算引擎
CoarseGrainedExecutorBackend 收到 RegisteredExecutor 消息后,创建真正的 Executor 对象:
1 | // CoarseGrainedExecutorBackend.scala — receive 方法 |
至此,CoarseGrainedExecutorBackend(进程)完成了注册,内部创建了 Executor(计算对象),可以开始接收任务了。
9. 总结
流程概述
Standalone Client 模式的完整提交流程:
- SparkSubmit 提交程序 →
prepareSubmitEnvironment设置childMainClass = args.mainClass→runMain在客户端反射执行用户主类的 main 方法,Driver 在客户端启动 - SparkContext 初始化 → 创建
TaskSchedulerImpl+StandaloneSchedulerBackend - StandaloneSchedulerBackend.start() → 创建
DriverEndpoint→ 构建ApplicationDescription→ 创建StandaloneAppClient - StandaloneAppClient 向 Master 发送
RegisterApplication(重试 3 次,间隔 20 秒) - Master 处理注册 → 创建
ApplicationInfo→ 回复RegisteredApplication→ 调用schedule()调度资源 - Master.schedule() →
startExecutorsOnWorkers()→ FIFO 遍历waitingApps→ 过滤可用 Worker →allocateWorkerResourceToExecutors()→launchExecutor()发送LaunchExecutor到 Worker - Worker 收到
LaunchExecutor→ 创建ExecutorRunner→fetchAndRunExecutor()→ 启动子进程CoarseGrainedExecutorBackend - CoarseGrainedExecutorBackend 启动 → 连接
DriverEndpoint→ 发送RegisterExecutor→ 注册成功后创建Executor对象 - Driver 继续执行用户业务代码 → 创建 RDD/DataFrame → Action 算子触发 Job → DAGScheduler/TaskScheduler 调度 Task 到已注册的 Executor 上执行
Standalone Client vs Yarn Client 对比
| 特性 | Standalone Client | Yarn Client |
|---|---|---|
| Driver 位置 | 客户端 | 客户端 |
| 集群管理器 | Master + Worker | ResourceManager + NodeManager |
| 资源调度 | Master.schedule() 自研调度器 | YARN 框架调度 |
| Executor 入口类 | CoarseGrainedExecutorBackend |
YarnCoarseGrainedExecutorBackend |
| Executor 启动方式 | Worker 直接启动子进程 | NMClient 启动 Container |
| AM/中间角色 | 无(Client 直连 Master) | ExecutorLauncher(AM) |
| Executor 资源申请 | Master 主动分配 | YarnAllocator 向 RM 申请 |
| 注册流程 | StandaloneAppClient → Master → Worker | YarnClientSchedulerBackend → RM → AM → NM |
| 重试机制 | 3 次,每次 20 秒 | YARN 框架重试 |
| 适用场景 | Spark 自建集群,轻量级 | 已有 Hadoop/YARN 集群,企业生产 |
核心调用链
1 | SparkSubmit.main |
完整的流程图
1 | ┌─────────────────────────────────────────────────────────────────────────┐ |
括号中的 [x.x] 对应本文各章节编号,便于快速回溯源码分析。
下一次我们将分析 Standalone Cluster 模式的详细提交流程,以及它和 Client 模式在 Driver 启动、资源调度等方面的差异。