Spark 核心之 Standalone-Client 模式:原理、流程与源码级深度拆解

摘要:Spark 的部署模式是面试和实际生产中最高频的话题之一,而 Standalone-Client 模式更是 90% 的开发者最常用却对其底层一知半解的模式。本文从整体架构、组件职责、启动流程、消息交互、任务调度五个维度,配合 3 张原创深色架构图和完整代码示例,带你彻底搞懂 Standalone-Client 模式的每一个细节。

关键词:Spark Standalone, Client 模式, Driver, Master, Executor, DAGScheduler, TaskScheduler, RPC, CoarseGrainedExecutorBackend


在这里插入图片描述

一、为什么你需要深入理解 Standalone-Client 模式?

如果你用过 Spark,你大概率执行过这样的命令:

spark-submit \
  --master spark://master:7077 \
  --deploy-mode client \
  --executor-memory 2G \
  --total-executor-cores 4 \
  my-app.jar

但你是否深入思考过以下问题:

  • Driver 到底运行在哪里? 是在 Master 节点还是在提交客户端?
  • Client 模式和 Cluster 模式的核心区别是什么? 为什么 Client 模式适合开发调试?
  • 从 spark-submit 到 Task 运行,中间经历了哪些步骤? 谁向谁注册?谁启动谁?
  • 如果 Driver 挂了,正在运行的 Executor 会怎样?
  • Master 是如何做资源调度的?Executor 是怎么被分配和启动的?

本文将逐一解答这些问题,从整体架构 → 启动流程 → 消息交互 → 任务调度四个层面,用 3 张原创深色架构图 + 源码级分析,带你彻底吃透 Standalone-Client 模式。


二、Spark Standalone 集群整体架构

在深入 Client 模式之前,先理解 Standalone 集群的整体架构。Spark Standalone 采用经典的 Master-Worker 主从架构,这是理解一切部署模式的基础。

2.1 四类核心角色

角色职责生命周期运行位置
Driver解析用户代码、生成 DAG、调度 Task、收集结果随应用启动/结束Client 模式:提交客户端;Cluster 模式:Worker 节点
Master集群资源管理、Worker 注册/心跳、Executor 分配调度集群级常驻进程独立 Master 节点
Worker管理本节点资源 (CPU/内存)、启动/停止 Executor集群级常驻进程各 Worker 节点
Executor运行 Task、缓存 RDD、提供 Shuffle 服务随应用启动/结束Worker 节点上的 JVM 进程

2.2 整体架构图

图 1:Spark Standalone 集群整体架构与组件交互关系

(架构图:01-architecture.png — 见上方生成的深色 SVG 架构图)

关键设计要点:

  1. Master 是"资源中介"而非"计算调度者":Master 只负责将资源(CPU 核心、内存)分配给应用,具体的 Task 调度由 Driver 内部的 DAGScheduler 和 TaskScheduler 完成。

  2. Executor 反向注册到 Driver:Worker 只是 fork 出 Executor 进程,Executor 启动后需要主动向 Driver 注册,之后 Driver 才能直接向 Executor 发送 Task。

  3. Driver 与 Executor 直连通信:注册完成后,Driver 与 Executor 之间通过 Netty RPC 直接通信,不再经过 Master。

// 源码:CoarseGrainedExecutorBackend.scala
// Executor 启动后反向注册到 Driver
override def onStart(): Unit = {
  logInfo("Connecting to driver: " + driverUrl)
  rpcEnv.asyncSetupEndpointRefByURI(driverUrl).flatMap { ref =>
    driver = Some(ref)
    // 向 Driver 发送 RegisterExecutor 消息
    ref.ask[Boolean](RegisterExecutor(executorId, self, hostname, cores, extractLogUrls))
  }(ThreadUtils.sameThread).onComplete { ... }
}

三、Client 模式 vs Cluster 模式:一张表彻底搞清

这是面试中的高频考点,也是实际选型的关键依据。

对比维度Client 模式Cluster 模式
Driver 运行位置提交客户端 JVMWorker 节点上的某个 JVM 进程
spark-submit 进程生命周期贯穿整个应用运行周期提交后即可退出
日志输出客户端控制台直接可见需通过 spark-submit --status 或 Web UI 查看
网络要求Driver 与各 Worker/Executor 之间网络必须互通Driver 在集群内部,无此问题
适用场景开发调试、交互式查询 (spark-shell)生产环境、定时任务
故障影响提交客户端宕机 → 整个应用失败客户端宕机不影响应用运行
资源占用Driver 占用客户端资源Driver 占用集群 Worker 资源

3.1 为什么 Client 模式适合开发调试?

┌─────────────────────────────────────────────────────┐
│          你的笔记本电脑 (Client 模式 Driver)          │
│                                                       │
│  spark-shell / IDEA 本地运行          Spark 集群       │
│  ┌──────────────┐      RPC     ┌──────────────────┐  │
│  │   Driver ❤️   │ ◄──────────► │  Master + Worker │  │
│  │  (你的 JVM)   │   日志/结果   │  + Executors     │  │
│  └──────────────┘              └──────────────────┘  │
│                                                       │
│  ✅ 日志直接打印在控制台                               │
│  ✅ 可以在代码中打断点 (breakpoint)                    │
│  ✅ 异常堆栈即时可见                                  │
│  ✅ 适合 interactive 分析                            │
└─────────────────────────────────────────────────────┘

但代价也很明确:你的网络必须能直连集群中每一个 Worker 节点和 Executor。如果 Executor 在云上的 VPC 内而你本地在办公网,Client 模式就无法工作——这正是 Cluster 模式的用武之地。


四、Client 模式启动流程:消息时序深度拆解

这是全文最核心的章节。我们将从 spark-submit 敲下回车的那一刻开始,逐步骤拆解整个启动链路。

4.1 四阶段总览

阶段核心动作关键消息
Phase 1:提交与注册Driver 向 Master 注册,申请资源RegisterApplicationRegisteredApplication
Phase 2:Executor 启动Master 调度 Worker fork Executor,反向注册LaunchExecutorRegisterExecutorRegisteredExecutor
Phase 3:任务执行DAGScheduler 切分 Stage,TaskScheduler 分发 TaskLaunchTaskStatusUpdate
Phase 4:应用结束Executor 销毁,Driver 注销KillExecutorsUnregisterApplication

(架构图:02-startup-sequence.png — 见上方消息时序图)

4.2 Phase 1:提交与注册

Step 1 — spark-submit 启动 Driver

spark-submit \
  --master spark://master:7077 \
  --deploy-mode client \
  --executor-memory 2G \
  --total-executor-cores 4 \
  my-app.jar

此时发生的事情:

  1. SparkSubmit 类解析参数,判断 deployMode == client
  2. 当前 JVM 进程中创建 SparkContext
  3. SparkContext 内部初始化三大核心组件:
    • DAGScheduler:负责 Stage 切分
    • TaskScheduler:负责 Task 分发(实现类 TaskSchedulerImpl
    • SchedulerBackend:负责与集群通信(实现类 StandaloneSchedulerBackend
// 源码:SparkContext.scala (简化)
// 根据 master URL 创建对应的 SchedulerBackend + TaskScheduler
private def createTaskScheduler(sc: SparkContext, master: String): (SchedulerBackend, TaskScheduler) = {
  master match {
    case SPARK_REGEX(sparkUrl) =>
      val scheduler = new TaskSchedulerImpl(sc)
      val backend = new StandaloneSchedulerBackend(scheduler, sc, sparkUrl)
      scheduler.initialize(backend)
      (backend, scheduler)
    // ...
  }
}

Step 2 — RegisterApplication

StandaloneSchedulerBackend 启动后,内部的 ClientEndpoint 向 Master 发送 RegisterApplication 消息:

// 源码:StandaloneAppClient.scala
class ClientEndpoint(override val rpcEnv: RpcEnv, ...) extends ThreadSafeRpcEndpoint {
  override def onStart(): Unit = {
    // 向 Master 发送注册请求
    try {
      registerMasterFutures = tryRegisterAllMasters()
    } catch { ... }
  }
  
  private def tryRegisterAllMasters(): Array[JFuture[_]] = {
    masterRpcAddresses.map { masterAddress =>
      // 异步向每个 Master 注册
      rpcEnv.setupEndpointRef(masterAddress, Master.ENDPOINT_NAME)
        .ask[RegisterApplicationResponse](RegisterApplication(appDescription, self))
    }
  }
}

Master 收到请求后:

  • 创建 ApplicationInfo 对象,分配 appId
  • 将 App 加入等待调度队列
  • 返回 RegisteredApplication(appId)

Step 3 — 请求 Executor 资源

Driver 收到 RegisteredApplication 后,根据用户配置的 --total-executor-cores--executor-memory,向 Master 请求启动 Executor:

// StandaloneAppClient.scala
case RegisteredApplication(appId, masterRef) =>
  // 注册成功后请求 Executor
  // numExecutors = total-executor-cores / executor-cores(或 spark.cores.max)
  appId_ = appId
  registered = true
  // 调用 SchedulerBackend 的 doRequestTotalExecutors
  ...

4.3 Phase 2:Executor 启动与反向注册

这是最容易让人困惑的环节。关键理解:Executor 不是由 Driver 直接启动的,而是由 Master 指令 Worker 启动的

Step 4 — Master 调度并指令 Worker

Master 的调度逻辑:

// 源码:Master.scala - schedule() 方法 (简化)
private def schedule(): Unit = {
  // 遍历等待运行的应用
  for (app <- waitingApps) {
    // 1. 筛选可用 Worker(满足 CPU/内存需求)
    val usableWorkers = workers.filter(_.state == WorkerState.ALIVE)
      .filter(canLaunchExecutor(_, app.desc))
    // 2. 在每个 Worker 上启动 Executor
    for (worker <- usableWorkers if app.coresLeft > 0) {
      // 3. 分配 cores 并发送 LaunchExecutor 消息
      val coresToUse = math.min(worker.coresFree, app.coresLeft)
      launchExecutor(worker, app, coresToUse)
    }
  }
}

调度策略核心要点:

  1. SpreadOut 算法(默认):尽量将 Executor 分散到不同 Worker,提高数据本地性
  2. FIFO / FAIR 调度:多应用场景下的资源分配策略
  3. Worker 筛选条件:内存 ≥ executor-memory,可用核心 > 0

Step 5 — Worker fork Executor 进程

Worker 收到 LaunchExecutor 消息后,fork 一个新的 JVM 进程运行 CoarseGrainedExecutorBackend

# Worker 内部执行的等价命令
java -cp spark-assembly.jar \
  org.apache.spark.executor.CoarseGrainedExecutorBackend \
  --driver-url spark://CoarseGrainedScheduler@client:PORT \
  --executor-id 0 \
  --cores 2 \
  --memory 2G \
  --hostname worker-node-1

Step 6 — Executor 反向注册到 Driver

这是最关键的环节。Executor 进程启动后,主动向 Driver 发起反向注册

// 源码:CoarseGrainedExecutorBackend.scala
override def onStart(): Unit = {
  logInfo("Connecting to driver: " + driverUrl)
  rpcEnv.asyncSetupEndpointRefByURI(driverUrl).flatMap { ref =>
    driver = Some(ref)
    // 🔑 关键:向 Driver 发送 RegisterExecutor
    ref.ask[Boolean](RegisterExecutor(executorId, self, hostname, cores, extractLogUrls))
  }(ThreadUtils.sameThread).onComplete {
    case Success(_) =>
      // 注册成功,Driver 返回 RegisteredExecutor
      self.send(RegisteredExecutor)
    case Failure(e) =>
      exitExecutor(1, s"Cannot register with driver: ...")
  }
}

为什么是"反向注册"?

因为 Worker fork 出来的只是一个 JVM 进程,它只知道 Driver URL(通过命令行参数传入),不知道如何被 Driver 发现。所以 Executor 必须主动向 Driver “报到”。

Step 7 — Driver 确认注册

Driver 收到 RegisterExecutor 后:

// 源码:CoarseGrainedSchedulerBackend.scala
case RegisterExecutor(executorId, executorRef, hostname, cores, logUrls) =>
  if (executorDataMap.contains(executorId)) {
    // 已注册 → 返回错误
    executorRef.send(RegisterExecutorFailed("Duplicate executor ID"))
  } else {
    // 新 Executor → 记录并确认
    val data = new ExecutorData(executorRef, executorRef.address, hostname, cores, ...)
    executorDataMap.put(executorId, data)
    executorRef.send(RegisteredExecutor)
    // 通知 TaskScheduler 有新的计算资源
    ...
  }

至此,Executor 与 Driver 之间的直连通信通道建立完毕。

4.4 Phase 3:任务执行

Driver 拿到 Executor 列表后,即可开始执行用户代码中的 Action 算子。详细流程见第五章。

4.5 Phase 4:应用结束

当用户代码执行完毕(或 SparkContext.stop() 被调用):

  1. Driver 向所有 Executor 发送 StopExecutor 消息
  2. Executor 进程退出
  3. Driver 向 Master 发送 UnregisterApplication
  4. Master 清理 App 记录和资源占用

五、任务调度与执行流水线

理解了 Client 模式的启动流程后,我们来看任务实际上是怎么被调度和执行的。
在这里插入图片描述

5.1 从用户代码到 Task 的转化链

用户代码 (RDD Transformation)
        │
        ▼
  RDD DAG (惰性求值,构建 Lineage)
        │
        ▼  Action 触发
  DAGScheduler 切分 Stage
        │
        ▼  Wide Dependency = Stage 边界
  TaskSet (每个 Partition 一个 Task)
        │
        ▼  TaskScheduler 分发
  Executor 执行 Task

5.2 DAGScheduler:Stage 切分的核心逻辑

// 简化版概念代码
def handleJobSubmitted(jobId: Int, finalRDD: RDD[_], ...): Unit = {
  // 1. 从 finalRDD 回溯依赖链,创建 ResultStage
  val finalStage = createResultStage(finalRDD, ...)
  
  // 2. 提交 Stage
  submitStage(finalStage)
}

private def submitStage(stage: Stage): Unit = {
  // 3. 递归提交父 Stage(ShuffleMapStage 必须先完成)
  val missing = getMissingParentStages(stage).sortBy(_.id)
  if (missing.isEmpty) {
    // 没有缺失的父 Stage → 提交当前 Stage 的 TaskSet
    submitMissingTasks(stage, jobId.get)
  } else {
    // 先提交父 Stage
    for (parent <- missing) {
      submitStage(parent)
    }
  }
}

Stage 切分规则:遇到 Shuffle Dependency (Wide Dependency) 即切分。

举例:

val rdd = sc.textFile("hdfs://...")     // HadoopRDD
  .flatMap(_.split(" "))                // MapPartitionsRDD (Narrow)
  .map((_, 1))                           // MapPartitionsRDD (Narrow)
  .reduceByKey(_ + _)                    // ShuffledRDD (Wide!)
  .filter(_._2 > 10)                     // MapPartitionsRDD (Narrow)
  .collect()                             // Action 触发

对应的 Stage 划分:

Stage 0 (ShuffleMapStage):  textFile → flatMap → map → reduceByKey 的 Map 端
Stage 1 (ResultStage):      reduceByKey 的 Reduce 端 → filter → collect

5.3 TaskScheduler:数据本地性与任务分发

TaskSchedulerImpl 拿到 TaskSet 后,通过 TaskSetManager 管理每个 Task 的分发:

// 数据本地性优先级(从高到低)
// PROCESS_LOCAL  >  NODE_LOCAL  >  RACK_LOCAL  >  ANY
// 同进程内缓存     同节点磁盘      同机架          任意

本地性等待策略

// TaskSetManager 源码
private def getAllowedLocalityLevel(curTime: Long): TaskLocality.TaskLocality = {
  // 逐步降低本地性要求
  // PROCESS_LOCAL 等待 spark.locality.wait.process (默认 3s)
  // NODE_LOCAL    等待 spark.locality.wait.node   (默认 3s)
  // RACK_LOCAL    等待 spark.locality.wait.rack   (默认 3s)
  ...
}

5.4 Executor 内部执行模型

Executor 是一个多线程的 JVM 进程,核心执行逻辑:

// 源码:Executor.scala - TaskRunner.run()
class TaskRunner(..., taskDescription: TaskDescription) extends Runnable {
  override def run(): Unit = {
    // 1. 反序列化 Task
    val ser = SparkEnv.get.closureSerializer.newInstance()
    val task = ser.deserialize[Task[Any]](taskDescription.serializedTask, ...)
    
    // 2. 设置 TaskMemoryManager,管理内存
    val taskMemoryManager = new TaskMemoryManager(env.memoryManager, taskId)
    
    // 3. 执行 Task
    val res = task.run(
      taskAttemptId = taskId,
      taskMemoryManager = taskMemoryManager, ...
    )
    
    // 4. 序列化结果并发送给 Driver
    val result = new DirectTaskResult[Any](ser.serialize(res), ...)
    execBackend.statusUpdate(taskId, TaskState.FINISHED, result.serialized())
  }
}

Executor 内部的并发模型

  • 一个 Executor = 一个 JVM 进程
  • 线程池大小 = spark.executor.cores(默认 1)
  • 每个 core 同一时间运行一个 Task
  • Task 之间共享 Executor 的 JVM 堆内存

六、Client 模式与 Cluster 模式的核心源码差异

理解了两者的宏观差异后,从源码层面看最关键的入口:

// 源码:SparkSubmit.scala
private def runMain(args: SparkSubmitArguments, ...): Unit = {
  if (args.isStandaloneCluster) {
    // ===== Cluster 模式 =====
    // 1. 在本地 fork 出一个子进程(ClientApp),负责向 Master 注册
    // 2. Master 会在某个 Worker 上启动 DriverWrapper
    // 3. 当前 spark-submit 进程可以退出
    runMainRestOrChild(args, childArgs, childClasspath, ...)
  } else {
    // ===== Client 模式 =====
    // 直接在当前 JVM 中反射调用用户 main 方法
    // Driver 就是当前进程
    runMainChild(args, childArgs, childClasspath, ...)
  }
}

Cluster 模式额外做的事

  1. ClientApp 向 Master 注册,告知需要启动 Driver
  2. Master 选择一个 Worker,fork DriverWrapper 进程
  3. DriverWrapper 内部再创建 SparkContext = 实际 Driver
  4. 所以 Cluster 模式下,spark-submit 进程的生命周期很短(提交完即退出)

而在 Client 模式 下,spark-submit 进程本身就是 Driver,必须一直存活到 Job 执行完毕。


七、生产环境最佳实践与避坑指南

7.1 何时使用 Client 模式?

✅ 适用场景❌ 不适用场景
本地开发与调试 (spark-shell / IDE)生产环境定时任务 (Cron Job)
交互式数据探索 (Notebook)提交客户端与集群网络不通
需要实时查看日志和堆栈提交后希望立即退出 (CI/CD 流水线)
需要使用断点调试大规模长时间运行的应用

7.2 常见配置调优

spark-submit \
  --master spark://master:7077 \
  --deploy-mode client \
  # Executor 配置
  --executor-memory 4G \              # Executor JVM 堆内存
  --executor-cores 2 \                 # 每个 Executor 的 core 数
  --total-executor-cores 8 \           # 总 core 数
  # Driver 配置(Client 模式下影响本地 JVM!)
  --driver-memory 2G \                 # Driver 内存
  --driver-cores 2 \                   # Driver 使用的 cores
  # 并行度配置
  --conf spark.default.parallelism=16 \  # RDD 默认分区数
  --conf spark.sql.shuffle.partitions=200 \  # SQL Shuffle 分区数
  # 调度配置
  --conf spark.scheduler.mode=FIFO \      # 调度模式
  --conf spark.locality.wait=3s \         # 数据本地性等待时间
  my-app.jar

7.3 常见故障排查

问题 1:Executor 无法反向注册到 Driver
ERROR CoarseGrainedExecutorBackend: Cannot register with driver

根因:Executor 所在节点无法访问 Driver 的 RPC 端口(通常是随机端口)。

解决

# 固定 Driver RPC 端口(Client 模式下 Driver 在提交客户端)
--conf spark.driver.port=7078
# 或显式指定 Driver Host(如果是多网卡机器)
--conf spark.driver.host=192.168.1.100
# 确保 Worker 节点能访问 Driver 主机的该端口
问题 2:Initial job has not accepted any resources
WARN TaskSchedulerImpl: Initial job has not accepted any resources

根因:集群资源不足,或 Worker 未正常注册。

排查步骤

  1. 检查 Master Web UI (http://master:8080):Worker 是否 Alive?
  2. 检查所需的内存和 CPU 是否超过集群总资源
  3. 检查 Worker 日志:是否有 LaunchExecutor 失败的错误
问题 3:Client 模式 Driver OOM

Client 模式下 Driver 运行在提交客户端的 JVM,如果执行 .collect() 拉取大量数据到 Driver,容易 OOM。

解决

# 增大 Driver 内存
--driver-memory 8G
# 避免 collect() 大数据量 → 改用 saveAsTextFile 或 take(n)
# 监控 Driver 内存使用(GC 日志)
--conf spark.driver.extraJavaOptions="-XX:+PrintGCDetails"

八、总结

核心知识点速记

要点一句话总结
Driver 位置Client 模式在提交客户端 JVM,Cluster 模式在 Worker 节点
四步启动注册 → 申请资源 → Master 调度 Worker → Executor 反向注册
反向注册Executor 是 Worker fork 的子进程,主动连接 Driver 注册
Task 调度Driver 的 DAGScheduler + TaskScheduler 负责,Master 只分配资源
适用场景Client 模式适合调试,Cluster 模式适合生产
关键区别spark-submit 进程生命周期的不同

三步走:从理解到实战

  1. 画一遍启动时序图:把 Client 模式启动的 7 个核心步骤自己走一遍
  2. 看一遍源码入口:从 SparkSubmit.scalaStandaloneSchedulerBackendCoarseGrainedExecutorBackend
  3. 搭一套本地环境:在本机起一个 Standalone 集群(Master + Worker),分别用 Client / Cluster 模式提交任务,观察 Web UI 的变化

当你真正理解了 Client 和 Cluster 模式的区别,理解了 Executor 反向注册的设计,你就真正掌握了 Spark 部署模式的精髓。


作者:starzy | AI Data Engineer / 大数据技术实践者
博客:blog.starzy.cn | GitHub:starzy1990.github.io
专注 AI Agent · LangGraph · RAG · 大数据架构 · 数据工程实践

更多推荐