前言
在第七篇文章 Standalone Client 模式提交流程(七) 中,我们分析了 Standalone Client 模式的完整提交流程。本文继续分析 Standalone Cluster 模式,其 Driver 运行在集群 Worker 节点上,而不是提交客户端。
Standalone Cluster 模式是 Spark 自带集群管理器的完整形态:Master 负责资源调度和 Driver 分发,Worker 负责启动 Driver 和 Executor。其流程比 Client 模式多了一层"提交网关"。
版本
Spark 3.2.3
Standalone Cluster 模式概述
Standalone Cluster 模式下,Driver 运行在集群中的某个 Worker 节点上。spark-submit 进程只负责提交请求,不运行用户程序。
注意:Master 和 Worker 需要提前启动(start-master.sh / start-worker.sh),spark-submit 时集群资源管理的基础设施已经就绪。
与 Standalone Client 的核心差异:
ClientAppRestSubmissionClientApp(包装器) | ||
--supervise |
流程对比总览:
Standalone Client:
spark-submit → 直接运行用户类 → SparkContext → 注册 Application → 启动 Executor
Standalone Cluster:
spark-submit → ClientApp → 向 Master 提交 Driver → Master 调度到 Worker
→ Worker: DriverWrapper → 反射调用用户类 → SparkContext → 注册 Application → 启动 Executor// 提交命令示例
spark-submit --master spark://host:7077 --deploy-mode cluster \
--classorg.apache.spark.examples.SparkPi\
$SPARK_HOME/examples/jars/spark-examples_2.12-3.2.3.jar0. RPC 通信方式汇总
本文涉及的 RPC 通信对照(与 Standalone Client 相同的部分不再重复):
masterRef.ask(RequestSubmitDriver) | ask | |
worker.endpoint.send(LaunchDriver) | send | |
WorkerWatcher 监控 | ||
StandaloneAppClient.registerWithMaster | send |
完整提交流程
1. SparkSubmit 准备阶段
1.1 prepareSubmitEnvironment
在 SparkSubmit.scala 的 prepareSubmitEnvironment 中,当 master.startsWith("spark://") 且 deployMode == "cluster" 时:
// SparkSubmit.scala — prepareSubmitEnvironment
// isStandaloneCluster: master=spark://xxx 且 deployMode=cluster
val isStandaloneCluster = clusterManager == STANDALONE && deployMode == CLUSTER首先进入 Fail fast 段落,排除了不支持的模式——Python 和 R 在 Standalone Cluster 模式下会直接报错终止:
// SparkSubmit.scala:prepareSubmitEnvironment — Fail fast 段落
// 不支持的模式直接 error 终止
(clusterManager, deployMode) match {
case (STANDALONE, CLUSTER) if args.isPython =>
error("Cluster deploy mode is currently not supported for python " +
"applications on standalone clusters.")
case (STANDALONE, CLUSTER) if args.isR =>
error("Cluster deploy mode is currently not supported for R " +
"applications on standalone clusters.")
case (LOCAL, CLUSTER) =>
error("Cluster deploy mode is not compatible with master \"local\"")
case (_, CLUSTER) if isShell(args.primaryResource) =>
error("Cluster deploy mode is not applicable to Spark shells.")
case _ =>
}通过 match 模式匹配在准备阶段就拦截,而不是等到运行时才报错。如果走 REST API 方式(useRest=true),StandaloneRestServer:buildDriverDescription 也会对 mainClass 做非空校验——Python 应用没有 mainClass,同样会提前失败:
// StandaloneRestServer.scala:buildDriverDescription
// REST 方式同样不支持 Python,mainClass 是必填字段
val mainClass = Option(request.mainClass).getOrElse {
thrownewSubmitRestMissingFieldException("Main class is missing.")
}拦截后,进入 Standalone Cluster 模式的分支,根据 useRest 区分 Legacy 和 REST 两种提交方式,设置 childMainClass:
// SparkSubmit.scala — prepareSubmitEnvironment 中的 Standalone Cluster 分支
if (args.isStandaloneCluster) {
if (args.useRest) {
// REST API 方式(需 spark.master.rest.enabled=true)
childMainClass = "org.apache.spark.deploy.rest.RestSubmissionClientApp"
childArgs += (args.primaryResource, args.mainClass)
} else {
// Legacy 方式(默认):通过 RPC 提交 Driver
childMainClass = "org.apache.spark.deploy.ClientApp"
if (args.supervise) { childArgs += "--supervise" }
Option(args.driverMemory).foreach { m => childArgs += ("--memory", m) }
Option(args.driverCores).foreach { c => childArgs += ("--cores", c) }
childArgs += "launch"
childArgs += (args.master, args.primaryResource, args.mainClass)
}
if (args.childArgs != null) { childArgs ++= args.childArgs }
}useRest 的默认值为 false,读取配置 spark.master.rest.enabled(默认为 "false")。因此默认走 Legacy 方式。
注意 childMainClass 不再是用户主类,而是 ClientApp 或 RestSubmissionClientApp——与 Standalone Client 模式(spark-submit 直接运行用户代码)不同,此模式下 spark-submit 进程变成了一个纯粹"提交客户端"。
childArgs 中有一行 childArgs += "launch"。这个 "launch" 最终会原封不动传给 ClientApp,参数传递链路如下:
// SparkSubmit.scala:prepareSubmitEnvironment — 拼装 childArgs
childArgs += "launch" // args(0) = "launch"
childArgs += (args.master, ...) // args(1..n) 为 master、jar、mainClass 等
// SparkSubmit.scala:runMain — 反射调用 ClientApp.main
ClientApp.main(childArgs.toArray) // 传入 ["launch", "spark://host:7077", "my.jar", "MyClass", ...]
// ClientArguments.scala — 解析参数
args(0) = "launch" → cmd = "launch" // 驱动 onStart 中的 case "launch" 分支2. 两种提交方式总览
Standalone Cluster 模式支持两种提交方式:
ClientApp | |||
RestSubmissionClientApp |
两种方式都只负责向 Master 提交 Driver 描述信息,然后轮询 Driver 状态。spark-submit 进程不运行用户代码。
spark-submit
├─ useRest=true → RestSubmissionClientApp → POST /v1/submissions/create → Master REST Server
└─ useRest=false → ClientApp → RPC RequestSubmitDriver → Master如果启用了 REST 方式但连接失败,会 fallback 到 Legacy 方式:
// SparkSubmit.scala submit()
if (args.isStandaloneCluster && args.useRest) {
try {
doRunMain() // 尝试 REST
} catch {
case e: SubmitRestConnectionException =>
args.useRest = false
submit(args, false) // 回退到 Legacy
}
} else {
doRunMain()
}3. Legacy 提交方式(默认,RPC)
3.1 ClientApp.start()
ClientApp 的 start 方法创建 RPC 环境,连接到 Master,注册 ClientEndpoint:
// Client.scala — ClientApp.start()
classClientAppextendsSparkApplication{
overridedefstart(args: Array[String], conf: SparkConf): Unit = {
val driverArgs = newClientArguments(args)
val rpcEnv = RpcEnv.create("driverClient", Utils.localHostName(), 0, conf,
newSecurityManager(conf))
val masterEndpoints = driverArgs.masters.map(RpcAddress.fromSparkURL)
.map(rpcEnv.setupEndpointRef(_, Master.ENDPOINT_NAME))
rpcEnv.setupEndpoint("client", newClientEndpoint(rpcEnv, driverArgs, masterEndpoints, conf))
rpcEnv.awaitTermination()
}
}3.2 ClientEndpoint.onStart()
ClientEndpoint.onStart() 是 Standalone Cluster 模式的关键,它构造 DriverDescription 并提交给 Master:
// Client.scala — ClientEndpoint.onStart()
overridedefonStart(): Unit = {
driverArgs.cmd match {
case"launch" =>
// Driver 的 mainClass 被设置为 DriverWrapper(不是用户类)
val mainClass = "org.apache.spark.deploy.worker.DriverWrapper"
// ...
val command = newCommand(mainClass,
Seq("{{WORKER_URL}}", "{{USER_JAR}}", driverArgs.mainClass) ++ driverArgs.driverOptions,
sys.env, classPathEntries, libraryPathEntries, javaOpts)
val driverDescription = newDriverDescription(
driverArgs.jarUrl, driverArgs.memory, driverArgs.cores,
driverArgs.supervise, command, driverResourceReqs)
// 向所有 Master 异步发送 RequestSubmitDriver
asyncSendToMasterAndForwardReply[SubmitDriverResponse](
RequestSubmitDriver(driverDescription))
}
}核心要点:
• mainClass被替换为DriverWrapper,而不是用户的 mainClass• 用户的 mainClass 被放入 command.arguments中,作为DriverWrapper的参数传入• {{WORKER_URL}}和{{USER_JAR}}是占位符,Worker 启动时替换为实际值• asyncSendToMasterAndForwardReply向所有 Master 发送请求,回复由receive处理
3.3 启动轮询
ClientEndpoint.onStart() 最后一步是启动定时任务,向 Master 轮询 Driver 状态(详细分析见第 7 节):
// Client.scala — ClientEndpoint.onStart()
forwardMessageThread.scheduleAtFixedRate(() => Utils.tryLogNonFatalError {
monitorDriverStatus()
}, 5000, REPORT_DRIVER_STATUS_INTERVAL, TimeUnit.MILLISECONDS)4. REST 提交方式(可选)
4.1 RestSubmissionClientApp
启用方式:spark-submit --conf spark.master.rest.enabled=true ...
RestSubmissionClientApp.start() 解析参数后调用 RestSubmissionClient:
// RestSubmissionClient.scala — RestSubmissionClientApp.start()
classRestSubmissionClientAppextendsSparkApplication{
overridedefstart(args: Array[String], conf: SparkConf): Unit = {
val client = newRestSubmissionClient(master)
val submitRequest = client.constructSubmitRequest(
appResource, mainClass, appArgs, sparkProperties, env)
client.createSubmission(submitRequest)
}
}4.2 REST API 通信
createSubmission 向 Master 发送 HTTP POST 请求:
POST http://master:6066/v1/submissions/create
Content-Type: application/json
{
"action": "CreateSubmissionRequest",
"appResource": "hdfs://...jar",
"mainClass": "org.apache.spark.examples.SparkPi",
"appArgs": ["10"],
"sparkProperties": {"spark.master": "spark://...", ...},
"environmentVariables": {...}
}请求流程:
// RestSubmissionClient.scala — createSubmission
for (m <- masters if !handled) {
validateMaster(m)
val url = getSubmitUrl(m)
response = postJson(url, request.toJson)
response match {
case s: CreateSubmissionResponseif s.success =>
reportSubmissionStatus(s) // 成功后轮询状态
handled = true
}
}提交成功后,同样会轮询 submission 状态(pollSubmissionStatus,最多 10 次,间隔 1s)。
5. Master 调度 Driver
无论是 Legacy 还是 REST 方式,最终 Master 都会收到提交 Driver 的请求。以 Legacy 方式为例:
5.1 Master 接收 RequestSubmitDriver
// Master.scala — receive 中的 RequestSubmitDriver 处理
caseRequestSubmitDriver(description) =>
val driver = createDriver(description) // 创建 DriverInfo
persistenceEngine.addDriver(driver) // 持久化(HA 支持)
waitingDrivers += driver // 加入等待队列
drivers.add(driver)
schedule() // 触发调度
context.reply(SubmitDriverResponse(self, true, Some(driver.id), ...))5.2 schedule() 调度 Driver
Master 的 schedule() 方法中,Driver 调度优先于 Executor 调度:
// Master.scala — schedule()
privatedefschedule(): Unit = {
// 优先调度 Driver
val shuffledAliveWorkers = Random.shuffle(workers.toSeq.filter(_.state == WorkerState.ALIVE))
for (driver <- waitingDrivers.toList) {
var launched = false
while (numWorkersVisited < numWorkersAlive && !launched) {
val worker = shuffledAliveWorkers(numWorkersVisited)
if (canLaunchDriver(worker, driver.desc)) {
launchDriver(worker, driver)
waitingDrivers -= driver
launched = true
}
numWorkersVisited += 1
}
}
// 然后调度 Executor
startExecutorsOnWorkers()
}canLaunchDriver 检查 Worker 是否有足够的空闲内存和 CPU 核心来运行 Driver。
launchDriver 向 Worker 发送 LaunchDriver 消息:
// Master.scala — launchDriver
worker.endpoint.send(LaunchDriver(driver.id, driver.desc, driver.resources))
driver.state = DriverState.RUNNING6. Worker fork 子进程运行 Driver
6.1 Worker 接收 LaunchDriver
// Worker.scala — receive
caseLaunchDriver(driverId, driverDesc, resources_) =>
val driver = newDriverRunner(conf, driverId, workDir, sparkHome, driverDesc, ...)
drivers(driverId) = driver
driver.start()DriverRunner.start() 在新线程中启动 Driver 进程:
// DriverRunner.scala — start()
private[worker] defstart() = {
newThread("DriverRunner for " + driverId) {
overridedefrun(): Unit = {
val exitCode = prepareAndRunDriver() // 准备并运行 Driver(详见 6.2)
// 根据退出码和是否被 kill 设置最终状态
finalState = if (exitCode == 0) Some(DriverState.FINISHED)
elseif (killed) Some(DriverState.KILLED)
elseSome(DriverState.FAILED)
}
// 通知 Worker:Driver 已结束
worker.send(DriverStateChanged(driverId, finalState.get, finalException))
}.start()
}start() 的核心逻辑:新线程运行 → prepareAndRunDriver() → ProcessBuilder.start() 创建全新 OS 进程(DriverWrapper)→ 根据退出码设置最终状态 → 通知 Worker。
DriverWrapper 和 Worker 是两个完全独立的 JVM 进程,只是运行在同一台机器上。Worker 通过 ProcessBuilder.start() fork 出子进程运行 DriverWrapper,并非在 Worker 进程内启动。
6.2 DriverRunner 准备并运行 Driver
DriverRunner.prepareAndRunDriver() 的完整源码和关键步骤:
// DriverRunner.scala — prepareAndRunDriver()
private[worker] defprepareAndRunDriver(): Int = {
val driverDir = createWorkingDirectory() // 1. 创建 Driver 工作目录
val localJarFilename = downloadUserJar(driverDir) // 2. 下载用户 JAR
defsubstituteVariables(argument: String): String = argument match {
case"{{WORKER_URL}}" => workerUrl // 3a. 替换 Worker URL
case"{{USER_JAR}}" => localJarFilename // 3b. 替换本地 JAR 路径
case other => other
}
val builder = CommandUtils.buildProcessBuilder(
driverDesc.command.copy(javaOpts = javaOpts),
securityManager, driverDesc.mem, sparkHome.getAbsolutePath, substituteVariables)
runDriver(builder, driverDir, driverDesc.supervise) // 4. 启动子进程
}对应之前的文字说明,这 4 步分别做了:
1. 创建 Driver 工作目录 2. 下载用户 JAR:通过 driverDesc.command中的jarUrl将用户 JAR 下载到本地(注意:Worker 通过Utils.fetchFile()下载,file:// 或裸路径走本地拷贝,需所有 Worker 该路径下有文件;也可用 HDFS/HTTP 等远程存储)3. 替换占位符: {{WORKER_URL}}替换为当前 Worker 的 RPC 地址,{{USER_JAR}}替换为下载后的本地路径4. 构建 ProcessBuilder:用 CommandUtils.buildProcessBuilder构造启动命令5. 启动子进程: ProcessBuilder.start()
实际启动的命令等价于:
java -cp <classpath> \
org.apache.spark.deploy.worker.DriverWrapper \
<workerUrl> <localJarPath> <userMainClass> [extraArgs...]6.3 DriverWrapper 入口
DriverWrapper 是 Worker fork 出的子进程的入口类。它接收 Worker URL、用户 JAR 路径和用户主类名:
// DriverWrapper.scala — main
case workerUrl :: userJar :: mainClass :: extraArgs =>
// 创建 RPC 环境,注册 WorkerWatcher 监控 Worker 连接
val rpcEnv = RpcEnv.create("Driver", host, port, conf, newSecurityManager(conf))
rpcEnv.setupEndpoint("workerWatcher", newWorkerWatcher(rpcEnv, workerUrl))
// 将用户 JAR 加入 ClassLoader
val loader = newMutableURLClassLoader(Array(userJarUrl), currentLoader)
Thread.currentThread.setContextClassLoader(loader)
setupDependencies(loader, userJar)
// 反射调用用户 mainClass 的 main 方法
val clazz = Utils.classForName(mainClass)
val mainMethod = clazz.getMethod("main", classOf[Array[String]])
mainMethod.invoke(null, extraArgs.toArray[String])WorkerWatcher 是一个关键组件——它监控与 Worker 的 RPC 连接。如果 Worker 进程挂了(网络断开),WorkerWatcher 会主动退出 Driver 进程,实现"命运共享"。
7. 轮询 Driver 状态
Client 向 Master 提交 Driver 后,spark-submit 进程需要知道:Driver 提交成功了吗?运行完了没有?我该不该退出? 由于 Master 的调度和 Worker 的启动都是异步的,Client 只能通过轮询来感知进展。
轮询结果决定了 spark-submit 的退出行为:
waitAppCompletion=false | exit(0),提交机使命完成 |
waitAppCompletion=true | exit(0) |
exit(-1) |
monitorDriverStatus 通过 asyncSendToMasterAndForwardReply 向 Master 发送 RequestDriverStatus,Master 回复后经 receive 路由到 reportDriverStatus:
// Client.scala — reportDriverStatus
defreportDriverStatus(
found: Boolean,
state: Option[DriverState],
workerId: Option[String],
workerHostPort: Option[String],
exception: Option[Exception]): Unit整体逻辑分三层:
reportDriverStatus(found, state, workerId, workerHostPort, exception)
if (found) → Master 找到了这个 driverId
├─ 打印 Driver 状态和所在 Worker(仅首次)
├─ exception match
│ ├─ Some(e) → 有异常,exit(-1)
│ └─ None → 按状态判断
│ ├─ FINISHED/FAILED/ERROR/KILLED → exit(0)
│ └─ 其他状态(RUNNING/SUBMITTED 等)
│ ├─ waitAppCompletion=false(默认)→ exit(0)
│ └─ waitAppCompletion=true → 继续轮询
└─ 来自 backup Master 的响应则忽略(HA 场景)
else → Master 不认识此 driverId → exit(-1)exception 为什么会出现——状态查询的回复怎么会带异常信息?看 Master 端如何回复 RequestDriverStatus:
// Master.scala — receive 中的 RequestDriverStatus 处理
caseRequestDriverStatus(driverId) =>
if (state != RecoveryState.ALIVE) {
// Master 处于非 ALIVE 状态(如 backup Master)
// exception 被设为异常消息,Client 端识别为 backup 后忽略
context.reply(DriverStatusResponse(found = false, None, None, None, Some(newException(msg))))
} else {
(drivers ++ completedDrivers).find(_.id == driverId) match {
caseSome(driver) =>
// 找到 driver,直接返回其状态和异常信息
// driver.exception 来自 DriverWrapper 上报的异常
context.reply(DriverStatusResponse(found = true, Some(driver.state),
driver.worker.map(_.id), driver.worker.map(_.hostPort), driver.exception))
caseNone =>
// 找不到此 driverId
context.reply(DriverStatusResponse(found = false, None, None, None, None))
}
}exception 在三个场景中的含义:
• Master 非 ALIVE(如 backup): found=false, exception=Some(msg),Client 识别为 backup 后静默忽略• Driver 运行出错: found=true, exception=driver.exception,进入Some(e)分支exit(-1)• 查无此 driverId: found=false, exception=None,进入最后的else分支exit(-1)
8. SparkContext 初始化
用户 mainClass 中的代码开始执行,例如 SparkPi.main() 中调用了 new SparkContext()。
这一步 与 Standalone Client 模式完全一致:
// SparkContext.scala — createTaskScheduler
caseSPARK_REGEX(sparkUrl) =>
val scheduler = newTaskSchedulerImpl(sc)
val masterUrls = sparkUrl.split(",").map("spark://" + _)
val backend = newStandaloneSchedulerBackend(scheduler, sc, masterUrls)
scheduler.initialize(backend)
(backend, scheduler)创建 TaskSchedulerImpl + StandaloneSchedulerBackend,没有分支区别。
8.1 StandaloneSchedulerBackend.start()
// StandaloneSchedulerBackend.scala — start()
overridedefstart(): Unit = {
super.start() // 创建 DriverEndpoint
// Cluster 模式不连接 LauncherBackend(Client 模式才连)
if (sc.deployMode == "client") {
launcherBackend.connect()
}
val command = Command("org.apache.spark.executor.CoarseGrainedExecutorBackend", ...)
client = newStandaloneAppClient(sc.env.rpcEnv, masters, appDesc, this, conf)
client.start() // 向 Master 注册 Application
waitForRegistration() // 信号量阻塞等待注册完成
}Cluster 模式跳过 launcherBackend.connect()——提交成功后 spark-submit 默认即退出,不需要接收 Driver 状态通知。
8.2 信号枪机制
与 Standalone Client 一样,使用 Semaphore(0) 阻塞等待 StandaloneAppClient 注册完成:
// StandaloneSchedulerBackend.scala
privateval registrationBarrier = newSemaphore(0)
privatedefwaitForRegistration() = {
registrationBarrier.acquire()
}当 StandaloneAppClient 收到 RegisteredApplication 回复后,回调 connected() 释放信号量:
overridedefconnected(msg: String): Unit = {
logInfo("StandaloneAppClient: " + msg)
notifyContext()
}
privatedefnotifyContext() = {
registrationBarrier.release()
}9. Application 注册与 Executor 启动
9.1 StandaloneAppClient 注册
client.start() → ClientEndpoint.onStart() → registerWithMaster(1) → tryRegisterAllMasters() → 向所有 Master 发送 RegisterApplication。
这一步与 Standalone Client 模式完全一致,包括重试机制(3 次,间隔 20 秒)。
9.2 Master 调度 Executor
Master 收到 RegisterApplication 后,创建 App 并触发 schedule():
// Driver 已经调度完成,现在调度 Executor
schedule()schedule() 中 Driver 部分已经跳过(waitingDrivers 为空),直接调用 startExecutorsOnWorkers(),其流程与 Standalone Client 完全一致:
1. scheduleExecutorsOnWorkers()— spreadOut 算法分配 Executor 到 Worker2. allocateWorkerResourceToExecutors()— 为每个 Worker 创建 ExecutorDesc3. launchExecutor()— 向 Worker 发送LaunchExecutor
9.3 Worker 启动 Executor
Worker 收到 LaunchExecutor,创建 ExecutorRunner,启动 CoarseGrainedExecutorBackend 子进程,替换占位符。与 Standalone Client 完全一致。
9.4 Executor 反向注册
CoarseGrainedExecutorBackend 进程启动后,通过 driverUrl 连接 DriverEndpoint,发送 RegisterExecutor,Driver 端创建 Executor 对象。
与 Standalone Client 完全一致。
10. 完整流程图
spark-submit (提交机)
│
├─ prepareSubmitEnvironment()
│ childMainClass = ClientApp(默认)或 RestSubmissionClientApp
│ childArgs = [launch, master, jar, mainClass, ...]
│
└─ runMain()
└─ [Legacy 方式] ClientApp.start()
├─ RpcEnv.create("driverClient", ...)
├─ setupEndpointRef(Master.ENDPOINT_NAME)
├─ setupEndpoint("client", new ClientEndpoint(...))
│
└─ ClientEndpoint.onStart()
├─ mainClass = DriverWrapper
├─ Command(DriverWrapper, [{{WORKER_URL}}, {{USER_JAR}}, userMainClass])
├─ asyncSendToMasterAndForwardReply(RequestSubmitDriver) → Master
└─ monitorDriverStatus() 定时轮询 Driver 状态
Master
├─ receive: RequestSubmitDriver
│ ├─ createDriver → waitingDrivers
│ ├─ schedule()
│ │ └─ launchDriver(worker, driver)
│ │ └─ worker.endpoint.send(LaunchDriver)
│ └─ reply: SubmitDriverResponse
│
└─ (后续) receive: RegisterApplication
└─ schedule() → launchExecutor → LaunchExecutor
Worker
├─ receive: LaunchDriver
│ └─ DriverRunner.start()
│ └─ prepareAndRunDriver()
│ ├─ 下载用户 JAR
│ ├─ 替换 {{WORKER_URL}} / {{USER_JAR}}
│ └─ Process: java DriverWrapper <workerUrl> <jar> <userMainClass>
Worker 子进程 - DriverWrapper
├─ RpcEnv.create("Driver", ...)
├─ WorkerWatcher 监控 Worker 连接
├─ 用户 JAR 加入 ClassLoader
└─ 反射调用 userMainClass.main(args)
用户代码 → new SparkContext()
├─ createTaskScheduler()
│ └─ TaskSchedulerImpl + StandaloneSchedulerBackend
│
└─ StandaloneSchedulerBackend.start()
├─ super.start() → DriverEndpoint
├─ StandaloneAppClient → registerWithMaster → RegisterApplication
├─ waitForRegistration() 信号量阻塞
│
└─ Master: schedule() → startExecutorsOnWorkers()
└─ Worker: ExecutorRunner → CoarseGrainedExecutorBackend
└─ 反向注册到 DriverEndpoint与 Standalone Client 的差异汇总
ClientAppRestSubmissionClientApp(包装器) | |||
--driver-memory | |||
--driver-cores | STANDALONE/MESOS/YARN/K8s, CLUSTER 生效,Client 模式不生效 | ||
--supervise | |||
| 完全一致 | |||
| 完全一致 | |||
Running Drivers 列表 | Running Drivers(只有 Running Applications) |
设计要点
1. 提交网关模式:Standalone Cluster 的 spark-submit 不运行用户代码,只做提交。 ClientApp和RestSubmissionClientApp是两种提交网关实现。2. DriverWrapper 反射层:Worker 通过 DriverWrapper启动用户 mainClass,中间加了一层 ClassLoader 隔离,保证用户依赖和 Spark 依赖不冲突。3. WorkerWatcher 命运共享:Driver 进程中的 WorkerWatcher监控 Worker 连接,Worker 挂则 Driver 自动退出,避免孤立进程。4. Legacy + REST 双通道:默认走 RPC 方式,可选 REST 方式(更稳定的跨版本协议),REST 失败会自动 fallback 到 Legacy。 5. 信号枪机制:与 Client 模式相同, StandaloneSchedulerBackend使用Semaphore(0)阻塞等待注册完成,区别在于此时代码运行在 Worker 节点的 DriverWrapper 子进程中,而非 spark-submit 进程。
🧐 分享、点赞、在看,给个3连击呗!👇
夜雨聆风