前些天发现了一个巨牛的人工智能学习网站,通俗易懂,风趣幽默,忍不住给大家分享一下。点击跳转到网站:https://www.captainai.net/dongkelun
前言
在第五篇文章 Yarn Client 模式提交流程(五) 中,我们分析了 Yarn Client 模式的完整提交流程。本文将继续分析 Spark 的另一种常用部署模式——Standalone Client 模式的提交流程。
Standalone 模式是 Spark 自带的集群管理器,不依赖 YARN、Mesos 等外部资源调度系统。其架构包含 Master(集群管理器)和 Worker(工作节点)两个角色,Master 负责资源调度,Worker 负责执行计算任务。
关于 Standalone 集群的搭建与配置,可参考之前的文章 Spark Standalone 集群配置。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 | // 提交命令示例 |
0. RPC 通信方式汇总
本文多处涉及 Spark RPC 通信,这里做一个统一说明。完整的 Spark RPC 框架总结见 Spark RPC 学习总结:
rpcEnv.setupEndpoint(name, endpoint):注册本地RpcEndpoint,注册完成后自动触发endpoint.onStart()。之后可通过返回的RpcEndpointRef收发消息。rpcEnv.setupEndpointRef(address, name):获取远程RpcEndpoint的引用,用于向远程端点发送消息。ref.send(msg):异步发送消息,无返回值。对应RpcEndpoint.receive。用于”通知”类消息。ref.ask[T](msg)/ref.askSync[T](msg):发送消息并等待返回值。对应RpcEndpoint.receiveAndReply。用于”请求-响应”类消息。
本文涉及的 RPC 通信对照:
| 位置 | 方式 | 原因 |
|---|---|---|
5.2 Client → Master 注册:masterRef.send(RegisterApplication) |
send |
只需通知,不阻塞等待,Master 成功后通过回调回复 |
6.1 Master → Client 回复:driver.send(RegisteredApplication) |
send |
同上,单向通知 |
6.5 Master → Worker:worker.endpoint.send(LaunchExecutor) |
send |
通知 Worker 启动 Executor |
6.5 Worker → Driver:exec.application.driver.send(ExecutorAdded) |
send |
通知 Driver 有新的 Executor 被分配 |
8.2 Executor 获取配置:driverRef.askSync(RetrieveSparkAppConfig) |
askSync |
需要获取配置内容,必须等结果 |
8.2 Executor 反向注册:ref.ask[Boolean](RegisterExecutor(...)) |
ask |
需要知道 Driver 是否确认注册成功 |
本文涉及的 RPC Endpoint 名称及注册位置:
| 角色 | RPC Endpoint 名称 | 注册代码 | 对应类 |
|---|---|---|---|
| Master | "Master" |
Master.scala setupEndpoint("Master", ...) |
Master |
| Worker | "Worker" |
Worker.scala setupEndpoint("Worker", ...) |
Worker |
| Driver 端 | "CoarseGrainedScheduler"(即 CoarseGrainedSchedulerBackend.ENDPOINT_NAME) |
CoarseGrainedSchedulerBackend.scala setupEndpoint(ENDPOINT_NAME, ...) |
DriverEndpoint |
| Executor 端 | "Executor" |
CoarseGrainedExecutorBackend.scala setupEndpoint("Executor", backend) |
CoarseGrainedExecutorBackend |
| AppClient | "AppClient" |
StandaloneAppClient.scala setupEndpoint("AppClient", ...) |
ClientEndpoint |
--driver-url spark://CoarseGrainedScheduler@driverHost:driverPort中的CoarseGrainedScheduler就是这个 RPC Endpoint 名称,不是类名- Executor 启动后以
"Executor"名称注册 RPC Endpoint,与"CoarseGrainedScheduler"对称
完整提交流程
1. SparkSubmit 准备阶段
在第二篇文章 SparkSubmit 提交流程分析(二) 中,我们分析了 prepareSubmitEnvironment 方法。该方法根据 --master 和 --deploy-mode 设置 childMainClass,决定 main 方法在哪里执行。
对于 Standalone Client 模式,关键代码在 SparkSubmit.scala 的 prepareSubmitEnvironment 方法中。当 deployMode == CLIENT 时,设置 childMainClass:
1 | // SparkSubmit.scala — prepareSubmitEnvironment 方法 |
关键点: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
createTaskScheduler 方法
createTaskScheduler 是 SparkContext 内部方法,根据 master URL 创建对应的 TaskScheduler 和 SchedulerBackend:
1 | // SparkContext.scala |
关键点:Standalone 模式走的是
SPARK_REGEX分支,直接在createTaskScheduler内部创建TaskSchedulerImpl+StandaloneSchedulerBackend。而 YARN 模式走的是masterUrl兜底分支,通过 Java SPI 机制加载ExternalClusterManager实现类来创建,具体分析见第四篇文章。
创建完成后,SparkContext 初始化过程中会调用 _taskScheduler.start() → TaskSchedulerImpl.start() 内部调用 backend.start() 即 StandaloneSchedulerBackend.start(),正式进入提交流程的核心。
4. StandaloneSchedulerBackend.start() — 核心入口
StandaloneSchedulerBackend 继承自 CoarseGrainedSchedulerBackend,实现了 StandaloneAppClientListener 接口。它的 start() 方法是整个 Standalone Client 提交流程的核心入口:
1 | // StandaloneSchedulerBackend.scala |
4.1 启动命令模板与 CoarseGrainedExecutorBackend
从代码中可以看到,Standalone 模式下的 Executor 启动入口类为:
1 | org.apache.spark.executor.CoarseGrainedExecutorBackend |
与 YARN 模式有所不同:
| 模式 | Executor 入口类 |
|---|---|
| Standalone | CoarseGrainedExecutorBackend |
| Yarn Client | YarnCoarseGrainedExecutorBackend(继承自 CoarseGrainedExecutorBackend) |
CoarseGrainedExecutorBackend 的启动命令分两阶段构建:
阶段一(Driver 端):构建命令模板 — StandaloneSchedulerBackend.start()
1 | // StandaloneSchedulerBackend.scala — start() 方法 |
这个 ApplicationDescription 通过 RegisterApplication 消息发送给 Master,Master 调度时传给 Worker。
阶段二(Worker 端):替换占位符拼出最终命令 — ExecutorRunner.fetchAndRunExecutor()
1 | // ExecutorRunner.scala — fetchAndRunExecutor() 方法 |
其中占位符替换的核心方法:
1 | // ExecutorRunner.scala — 占位符替换 |
为什么
--driver-url不是占位符?因为 Driver 地址在StandaloneSchedulerBackend.start()时已经确定(driverHost+driverPort),不需要等到 Worker 端再填。而executorId、hostname等要等到 Worker 真正分配资源时才知道实际值。
替换后的最终命令见第 7.2 节。
4.2 ApplicationDescription
ApplicationDescription 封装了应用的描述信息,包括应用名称、需要的 CPU/内存、Executor 启动命令等,Master 将根据这些信息进行资源调度:
1 | val appDesc = ApplicationDescription( |
5. StandaloneAppClient 注册流程
RPC 基础知识(
setupEndpoint/send/ask等)见第 0 章「RPC 通信方式汇总」。
StandaloneAppClient 是 Driver 端与 Master 通信的客户端。创建后,调用 client.start() 触发注册流程:
1 | // StandaloneAppClient.scala |
5.1 ClientEndpoint.onStart()
setupEndpoint 注册 RpcEndpoint 时自动触发 onStart()(见 0 RPC 基础),开始注册流程:
1 | // StandaloneAppClient.scala — ClientEndpoint 内部类 |
5.2 registerWithMaster() 与 tryRegisterAllMasters()
registerWithMaster() 调用 tryRegisterAllMasters() 向所有 Master 地址发送注册请求,并设置定时器处理重试:
1 | // StandaloneAppClient.scala — ClientEndpoint |
重试机制总结:
- 最大重试次数:3 次
- 重试间隔:20 秒
- 总超时:最多 60 秒(3 × 20s)
- 并发:同时向所有 Master 地址发起注册
为什么要向多个 Master 注册?
核心原因:Spark Standalone Master HA(高可用)机制
虽然单机测试时只有一个 Master,但生产环境中 Master 是单点故障(SPOF)。Spark 通过以下方式支持 Master HA:
- 配置多个 Master 地址:
spark.master参数支持逗号分隔多个地址,如:典型 HA 部署为一主一备两个 Master,但数组不做数量限制,理论上可多于两个(实际极少需要)1
spark.master=spark://master1:7077,master2:7077,master3:7077
- ZooKeeper 协调:HA 模式下多个 Master 节点通过 ZooKeeper 选举一个 Active Master,其余为 Standby
masterUrls的解析:每个地址都会作为1
2// 支持逗号分隔的多个 Master 地址,用于 Standalone Master HA 场景,典型一主一备
val masterUrls = sparkUrl.split(",").map("spark://" + _)StandaloneAppClient.masterRpcAddresses中的一个元素- “广撒网”策略:
- Client/Worker 启动时不知道哪个 Master 是 Active
- 并行向所有配置的 Master 发送注册请求
- 只有 Active Master 会响应
RegisteredApplication - Standby Master 会忽略或拒绝注册
- 注册成功后,后续通信只与 Active Master 进行
Worker 端也采用完全相同的策略向所有 Master 注册,原理一致。
tryRegisterAllMasters() 向所有 Master 地址发送 RegisterApplication 消息:
1 | // StandaloneAppClient.scala — ClientEndpoint |
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 |
Round-Robin 轮询机制详解:
curPos 是跨 Driver 持续累加的,不重置回 0:
- Driver A 尝试 Worker[0](资源不足)→ Worker[1](成功启动,curPos=2)
- Driver B 从 Worker[2] 开始尝试 → Worker[3](成功启动,curPos=4)
这种”从上次停下的位置继续”的方式,保证多个 Driver 间公平分配 Worker,不会出现所有 Driver 都挤在同一个 Worker 上的情况。
为什么 Client 模式不走 Driver 调度?
关键在 schedule() 遍历的是 waitingDrivers 队列,而 Client 模式发送的是 RegisterApplication 消息,Master 处理时创建 ApplicationInfo 并加入 waitingApps。Driver 进入 waitingDrivers 的唯一途径是 Master 收到提交 Driver 的请求(Cluster 模式的 SubmitDriverRequest)。
1 | Client 模式:RegisterApplication → ApplicationInfo → waitingApps → startExecutorsOnWorkers() |
所以 Client 模式下 waitingDrivers 始终为空,Driver 调度循环直接跳过,走到 Executor 调度。
- 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 ... |
这些参数的含义和构建过程见 4.1 节”阶段二(Worker 端)”。
8. Executor 注册到 Driver
8.1 CoarseGrainedExecutorBackend 启动
CoarseGrainedExecutorBackend 进程启动后,执行 main() 方法,解析命令行参数并调用 run():
1 | // CoarseGrainedExecutorBackend.scala |
run() 方法执行 Executor 启动的完整流程:
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.receiveAndReply |
8.4 创建 Executor 实例
CoarseGrainedExecutorBackend 收到 RegisteredExecutor 消息后,创建 Executor 实例:
1 | // CoarseGrainedExecutorBackend.scala — receive 方法 |
至此,CoarseGrainedExecutorBackend 进程 完成了注册(8.2),进程内部创建了 Executor 实例(executor = new Executor(...)),可以开始接收并执行 Task 了。
关于 Executor 这个名称的歧义说明
本文及 Spark 文档中多处出现”Executor”,实际指两个不同层次的概念:
| 概念 | 本质 | 对应代码 |
|---|---|---|
| Executor 进程 | 独立 JVM 进程 | CoarseGrainedExecutorBackend.main() 启动的进程 |
| Executor 实例 | 进程内部的 Task 执行与 RDD 缓存对象 | new Executor(executorId, ...) |
- 业界通用的 “Spark Executor”、UI 界面上显示的 Executor 条目、本文第 7~8 章讨论的注册过程,都是指 Executor 进程
- 第 8.4 节
new Executor(...)特指进程内部的 Executor 对象,是实现细节
不同模式下 Executor 进程的入口类有所不同:
| 模式 | Executor 进程入口类 | 说明 |
|---|---|---|
| Standalone | CoarseGrainedExecutorBackend |
Spark 自带调度系统,直接启动子进程 |
| Yarn | YarnCoarseGrainedExecutorBackend |
继承自 CoarseGrainedExecutorBackend,适配 YARN Container 生命周期 |
| Kubernetes | CoarseGrainedExecutorBackend |
与 Standalone 相同入口类,通过 K8s 启动 Pod |
但无论哪种模式,Executor 进程内部最终都通过 new Executor(...) 创建实例,Task 执行逻辑完全一致。
关系:Worker/Container 启动 Executor 进程 → 进程向 Driver 注册 → 注册成功后进程内部创建 Executor 实例 → 实例执行 Task
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(即 Executor 进程) - CoarseGrainedExecutorBackend 进程 启动 →
run()先通过临时 RpcEnv 连接 Driver 获取 SparkConf 配置 → 创建正式 SparkEnv → 注册 RPC Endpoint 触发onStart()→ 发送RegisterExecutor→ Driver 确认注册 → 进程内部创建 Executor 实例(new 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 启动、资源调度等方面的差异。