乐于分享
好东西不私藏

Spark 源码 | Standalone Cluster 模式提交流程(八)

Spark 源码 | Standalone Cluster 模式提交流程(八)

前言

在第七篇文章 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 的核心差异:

维度
Standalone Client
Standalone Cluster
Driver 位置
spark-submit 进程内(客户端)
Master 调度到某个 Worker 节点
spark-submit 主类
用户 mainClass
ClientApp
 或 RestSubmissionClientApp(包装器)
提交后 spark-submit
一直运行,直到应用退出
提交成功后默认即退出
LauncherBackend(状态通知通道)
建立(StandaloneSchedulerBackend.start 中)
不建立(提交后即退出,不需要)
Python/R 支持
支持
不支持
--supervise
 支持
不适用
支持(Driver 失败后自动重启)

流程对比总览:

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.jar

0. RPC 通信方式汇总

本文涉及的 RPC 通信对照(与 Standalone Client 相同的部分不再重复):

位置
方式
原因
3.1 Client → Master:masterRef.ask(RequestSubmitDriver)
ask
需要知道 Master 是否接受了 Driver 提交
4.1 Master → Worker:worker.endpoint.send(LaunchDriver)
send
通知 Worker fork 子进程运行 Driver
5.1 Worker → Driver:WorkerWatcher 监控
RPC 连接
监控 Worker 连接状态,Worker 挂则 Driver 退出
6.2 Driver → Master:StandaloneAppClient.registerWithMaster
send
通知 Master 注册 Application

完整提交流程

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 (STANDALONECLUSTERif args.isPython =>
    error("Cluster deploy mode is currently not supported for python " +
"applications on standalone clusters.")
case (STANDALONECLUSTERif args.isR =>
    error("Cluster deploy mode is currently not supported for R " +
"applications on standalone clusters.")
case (LOCALCLUSTER) =>
    error("Cluster deploy mode is not compatible with master \"local\"")
case (_, CLUSTERif 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 模式支持两种提交方式:

提交方式
启用条件
入口类
通信协议
Legacy 提交(默认)
spark.master.rest.enabled=false
ClientApp
RPC(Netty)
REST 提交
spark.master.rest.enabled=true
RestSubmissionClientApp
HTTP REST

两种方式都只负责向 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()
}, 5000REPORT_DRIVER_STATUS_INTERVALTimeUnit.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, trueSome(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.RUNNING

6. 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 == 0Some(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. 1. 创建 Driver 工作目录
  2. 2. 下载用户 JAR:通过 driverDesc.command 中的 jarUrl 将用户 JAR 下载到本地(注意:Worker 通过 Utils.fetchFile() 下载,file:// 或裸路径走本地拷贝,需所有 Worker 该路径下有文件;也可用 HDFS/HTTP 等远程存储)
  3. 3. 替换占位符{{WORKER_URL}} 替换为当前 Worker 的 RPC 地址,{{USER_JAR}} 替换为下载后的本地路径
  4. 4. 构建 ProcessBuilder:用 CommandUtils.buildProcessBuilder 构造启动命令
  5. 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 的退出行为:

场景
spark-submit 行为
提交成功(RUNNING),默认 waitAppCompletion=false
立即 exit(0),提交机使命完成
提交成功 + waitAppCompletion=true
继续轮询,等 Driver FINISHED/FAILED 后再 exit(0)
提交失败(exception / 找不到 driverId / backup Master)
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 = falseNoneNoneNoneSome(newException(msg))))
  } else {
    (drivers ++ completedDrivers).find(_.id == driverId) match {
caseSome(driver) =>
// 找到 driver,直接返回其状态和异常信息
// driver.exception 来自 DriverWrapper 上报的异常
        context.reply(DriverStatusResponse(found = trueSome(driver.state),
          driver.worker.map(_.id), driver.worker.map(_.hostPort), driver.exception))
caseNone =>
// 找不到此 driverId
        context.reply(DriverStatusResponse(found = falseNoneNoneNoneNone))
    }
  }

exception 在三个场景中的含义:

  • • Master 非 ALIVE(如 backup):found=false, exception=Some(msg),Client 识别为 backup 后静默忽略
  • • Driver 运行出错found=true, exception=driver.exception,进入 Some(e) 分支 exit(-1)
  • • 查无此 driverIdfound=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. 1. scheduleExecutorsOnWorkers() — spreadOut 算法分配 Executor 到 Worker
  2. 2. allocateWorkerResourceToExecutors() — 为每个 Worker 创建 ExecutorDesc
  3. 3. 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 的差异汇总

编号
差异点
Cluster 模式
Client 模式
1
spark-submit 主类
ClientApp
/RestSubmissionClientApp(包装器)
用户 mainClass
2
Driver 位置
Worker 节点(子进程)
spark-submit 进程内
3
是否使用 DriverWrapper
是,Worker fork 子进程运行 DriverWrapper 反射调用用户类
4
Python/R 支持
不支持,直接报错
支持
5
spark-submit 进程与 Driver 状态通信
不建立状态通知通道(提交后即退出,不需要)
通过 LauncherBackend 建立 Socket 通道,接收 Driver 状态变化通知
6
--driver-memory
生效(用于 Worker 分配资源)
所有集群管理器 + 所有部署模式均生效
7
--driver-cores
生效(用于 Worker 分配资源)
仅 STANDALONE/MESOS/YARN/K8s, CLUSTER 生效,Client 模式不生效
8
--supervise
支持
不适用
9
SparkContext 初始化
TaskSchedulerImpl + StandaloneSchedulerBackend
完全一致
10
信号枪机制
Semaphore(0)
完全一致
11
spark-submit 进程
提交成功后默认即退出
一直运行到应用退出
12
用户 JAR 访问
需所有 Worker 节点在相同路径下都有该 JAR,或使用远程存储(HDFS/HTTP)
仅提交机本地即可,无需提前分发
13
Master Web UI 显示
有 Running Drivers 列表
无 Running Drivers(只有 Running Applications

设计要点

  1. 1. 提交网关模式:Standalone Cluster 的 spark-submit 不运行用户代码,只做提交。ClientApp 和 RestSubmissionClientApp 是两种提交网关实现。
  2. 2. DriverWrapper 反射层:Worker 通过 DriverWrapper 启动用户 mainClass,中间加了一层 ClassLoader 隔离,保证用户依赖和 Spark 依赖不冲突。
  3. 3. WorkerWatcher 命运共享:Driver 进程中的 WorkerWatcher 监控 Worker 连接,Worker 挂则 Driver 自动退出,避免孤立进程。
  4. 4. Legacy + REST 双通道:默认走 RPC 方式,可选 REST 方式(更稳定的跨版本协议),REST 失败会自动 fallback 到 Legacy。
  5. 5. 信号枪机制:与 Client 模式相同,StandaloneSchedulerBackend 使用 Semaphore(0) 阻塞等待注册完成,区别在于此时代码运行在 Worker 节点的 DriverWrapper 子进程中,而非 spark-submit 进程。

🧐 分享、点赞、在看,给个3连击👇