前些天发现了一个巨牛的人工智能学习网站,通俗易懂,风趣幽默,忍不住给大家分享一下。点击跳转到网站:https://www.captainai.net/dongkelun
前言
在第五篇文章 Yarn Client 模式提交流程(五) 中,我们详细分析了 Yarn Client 模式的提交流程。本文继续分析 Yarn Cluster 模式的提交流程,重点介绍与 Client 模式的关键差异——Driver 如何在 AM 内启动。
版本
Spark 3.2.3
Yarn Cluster 模式概述
Yarn Cluster 模式下,Driver 运行在 YARN 集群的 ApplicationMaster 容器中,客户端提交后即可断开,作业在集群中独立运行。
1 | // 提交命令示例 |
与 Yarn Client 模式的核心区别
| 特性 | Yarn Client | Yarn Cluster |
|---|---|---|
| Driver 位置 | 客户端 | AM 内部 |
| childMainClass | args.mainClass(用户类) | YarnClusterApplication |
| isClusterMode 判断 | args.userClass == null | args.userClass != null |
| AM 入口方法 | runExecutorLauncher() | runDriver() |
| 资源申请时机 | AM 启动后直接申请 Executor(Driver 已在客户端运行) | AM 先启动 Driver 线程,等 SparkContext 初始化完毕后再申请 Executor |
| 资源申请顺序 | Driver → AM → Executor | AM → Driver → Executor |
| Driver 启动方式与位置 | 客户端先启动 Driver(app.start 直接反射调用用户 main 方法),Driver 初始化 SparkContext 时,在 YarnClientSchedulerBackend.start() 中提交 Application 启动 AM |
客户端先启动 AM(app.start → YarnClusterApplication.start() → Client.run() → submitApplication() 提交 Application),AM 的 runDriver() 中再通过 startUserApplication() 反射调用用户 main 方法启动 Driver |
| 跨线程同步机制 | 无。Driver 在客户端主线程运行,AM 与 Driver 是两个独立进程,不需要跨线程同步 | Promise + wait/notify。Driver(用户线程)与 AM 主线程在同一进程不同线程,SparkContext 初始化完后通过 YarnClusterScheduler.postStartHook() → sparkContextInitialized() 填充 Promise,通知 AM 主线程 awaitResult 返回;用户线程通过 wait() 挂起,等 AM 注册和申请完资源后通过 resumeDriver() → notify() 唤醒 |
Yarn Cluster 模式完整提交流程
阶段一:客户端提交(SparkSubmit → YarnClusterApplication)
1.1 SparkSubmit 准备
回顾第二篇文章 SparkSubmit 提交流程分析(二),prepareSubmitEnvironment 方法中,Yarn Cluster 模式的配置如下:
1 | // Yarn Cluster 模式 |
关键差异:
- Client 模式:
childMainClass = args.mainClass,直接在客户端反射调用用户主类的 main 方法 - Cluster 模式:
childMainClass = "org.apache.spark.deploy.yarn.YarnClusterApplication",用户主类名和 jar 路径作为--class和--jar参数传入
1.2 runMain 加载 YarnClusterApplication
根据第二篇文章的分析,runMain 方法会:
- 加载
childMainClass(即YarnClusterApplication) YarnClusterApplication实现了SparkApplication接口,直接实例化- 调用
app.start(childArgs.toArray, sparkConf)
1 | // runMain 内部,加载 childMainClass 后启动应用 |
要点:YarnClusterApplication 实现了 SparkApplication,所以走的是直接实例化分支,而不是 JavaMainApplication 包装。
1.3 YarnClusterApplication.start()
1 | // Yarn Cluster 模式入口类,实现了 SparkApplication 接口 |
YarnClusterApplication.start() 创建 Client 实例并调用 Client.run()。注意这里的 Client 构造参数第三个是 null,与 Client 模式不同(Client 模式传入的是 sc.env.rpcEnv),因为此时 RPC 环境还没建立。
1.4 Client.run() → submitApplication()
调用位置对比:
- Cluster 模式:
YarnClusterApplication.start()→Client.run(),直接在客户端提交入口中调用 - Client 模式:
YarnClientSchedulerBackend.start()→Client.submitApplication()(详见第五篇文章),Driver 已在客户端运行后,在 SparkContext 初始化过程中触发提交
两种模式最终调用的 submitApplication() 方法流程相同,差异在于 createContainerLaunchContext 中构造的 AM 启动命令不同。
1 | def run(): Unit = { |
run() 方法的执行逻辑:
- 先提交:调用
submitApplication()向 YARN RM 提交 Application - 再等待:根据
spark.yarn.submit.waitAppCompletion配置决定行为:- 默认等待(
waitAppCompletion=true):进入monitorApplication循环,定期轮询 RM 获取应用状态,直到应用结束 - fireAndForget 模式(
waitAppCompletion=false):提交后立即返回,只打印一次状态(Cluster 模式下fireAndForget默认为 true,因为客户端提交后可断开)
- 默认等待(
submitApplication() 流程与第五篇分析的 Yarn Client 模式相同,但需要注意:两种模式下的 createContainerLaunchContext() 构造的 AM 启动命令不同(见下一节)。以下是完整源码注释:
1 | def submitApplication(): ApplicationId = { |
与 Client 模式的关键差异隐藏在 createContainerLaunchContext 中。
1.5 createContainerLaunchContext 模式差异
createContainerLaunchContext 方法中根据 isClusterMode 决定 AM 启动类:
1 | val amClass = |
以下是在 createContainerLaunchContext 中构建 AM 启动命令的完整源码分析。核心分为三步:JVM 参数 → AM 参数 → 拼接最终命令。
① JVM 参数(javaOpts)
1 | val javaOpts = ListBuffer[String]() |
② AM 参数(amArgs)
1 | // 启动类:Cluster → ApplicationMaster,Client → ExecutorLauncher |
除了 --jar 参数外,jar 文件还通过 YARN LocalResource 机制分发给 AM。createContainerLaunchContext 前半部分调用了 prepareLocalResources,将 jar 上传到 HDFS staging 目录:
1 | // createContainerLaunchContext 内部,准备 AM 所需的本地资源 |
prepareLocalResources 中对 args.userJar 的处理:
1 | // prepareLocalResources 内部 |
因此 jar 通过两个途径到达 AM:
--jar参数:args.userJar原始路径作为命令行参数传入,ApplicationMasterArguments解析后使用- YARN LocalResource:由
prepareLocalResources上传到 HDFS staging 目录,AM 启动时 YARN 自动本地化到容器文件系统中
1 | // 拼接完整 AM 参数 |
③ 最终命令拼接
1 | val commands = prefixEnv ++ |
因此 Cluster 模式 AM 的实际启动命令大致为:
1 | JAVA_HOME/bin/java -server -Xmx<amMemory>m \ |
这里 --jar 的值是 args.userJar(即 SparkSubmit 中传入的 args.primaryResource),是用户在 spark-submit 命令中指定的 jar 路径。该路径由 Client.prepareLocalResources() 作为 YARN LocalResource 上传到 HDFS staging 目录,AM 容器启动时 YARN 会将其本地化到容器的本地文件系统中。
Client 模式对比:Client 模式 AM 命令不包含 --class 和 --jar(userClass 和 userJar 均为 Nil),AM 启动类为 ExecutorLauncher。这是判断 isClusterMode 的关键——args.userClass != null。
阶段二:AM 启动(ApplicationMaster → runDriver)
2.1 ApplicationMaster 入口
NodeManager 根据 ContainerLaunchContext 中的命令启动 ApplicationMaster.main():
1 | def main(args: Array[String]): Unit = { |
ApplicationMasterArguments 解析 AM 命令行参数,支持:
| 参数 | 说明 |
|---|---|
--jar JAR_PATH |
用户 JAR 路径 |
--class CLASS_NAME |
用户主类名 |
--arg ARG |
用户程序参数(可多次使用) |
--properties-file FILE |
Spark 配置文件路径 |
--dist-cache-conf |
分布式缓存配置文件路径 |
2.2 ApplicationMaster.run()
1 | // ApplicationMaster 主入口方法,根据模式分支执行 |
Cluster 模式独有的初始化:
- 设置 Web UI 端口为随机端口(避免与集群中其他 Spark 进程冲突)
- 设置
spark.master = yarn - 设置
spark.submit.deployMode = cluster - 保存
spark.yarn.app.id
2.3 runDriver() — Cluster 模式核心方法
1 | private def runDriver(): Unit = { |
执行流程(时序角度):
1 | runDriver() 用户线程 (startUserApplication) |
阶段三:用户线程与 SparkContext 初始化
3.1 startUserApplication() 启动用户主类
1 | // 在独立线程中启动用户主类,返回用户线程引用 |
关键点:
- 用户主类(如
SparkPi)在一个名为"Driver"的独立线程中执行 - 通过反射调用用户主类的
main方法 - 异常时通过
sparkContextPromise.tryFailure通知runDriver()中awaitResult等待的 AM 主线程 - 正常结束或异常结束时,
finally块中调用sparkContextPromise.trySuccess(null)确保runDriver()不会永久阻塞
“Driver” 线程的命名含义:这个线程名直接反映了 Yarn Cluster 模式下 Driver 的角色——用户主类的 main 方法在哪里执行,Driver 就在哪里。这里 main 方法运行在 AM 内部的独立线程中,因此 Driver 就在 AM 内部。
3.2 SparkContext 初始化与 YarnClusterScheduler
用户主类的 main 方法中创建 SparkContext 时,在 createTaskScheduler 阶段(详见第四篇文章),由于 master = "yarn" 且 deployMode = "cluster",YarnClusterManager 会创建:
1 | // YarnClusterManager |
YarnCluster 模式创建的对象:
| 组件 | 类名 |
|---|---|
| TaskScheduler | YarnClusterScheduler(继承自 YarnScheduler → TaskSchedulerImpl) |
| SchedulerBackend | YarnClusterSchedulerBackend(继承自 YarnSchedulerBackend → CoarseGrainedSchedulerBackend) |
3.3 YarnClusterScheduler 的关键作用
YarnClusterScheduler 的核心作用:在 SparkContext 初始化完成后,通知 AM 主线程(runDriver())继续执行。这是 Yarn Cluster 模式下 Driver 和 AM 协同的”信号枪”。
为什么这里需要跨线程等待?因为 runDriver() 和用户主类的 main 方法运行在两个不同线程中:
- AM 主线程:执行
runDriver(),需要获取 SparkContext 对象中的rpcEnv、DRIVER_PORT等信息 - 用户线程(Driver):执行用户
main方法,内部创建 SparkContext
SparkContext 对象是在用户线程中创建的,AM 主线程无法直接访问。因此需要一种跨线程通信机制:用户线程创建完 SparkContext 后通过 Promise 通知 AM 主线程,主线程拿到 sc 对象后才能执行 registerAM 和 createAllocator。
1 | // YarnClusterScheduler 是 Yarn Cluster 模式的 TaskScheduler |
ApplicationMaster.sparkContextInitialized(sc) 的完整调用链(伴生对象委托给实例方法):
1 | // 伴生对象静态方法(YarnClusterScheduler.postStartHook 调用此方法) |
调用链:
SparkContext构造函数中顺序执行:先调用_taskScheduler.start(),再调用_taskScheduler.postStartHook()_taskScheduler.start()内部调用backend.start()完成 DriverEndpoint 创建_taskScheduler.postStartHook()由于多态,实际调用的是YarnClusterScheduler.postStartHook()YarnClusterScheduler.postStartHook()内部调用ApplicationMaster.sparkContextInitialized(sc)sparkContextPromise.success(sc)填充 Promise,用户线程中这一步执行后,AM 主线程中阻塞的awaitResult(sparkContextPromise.future)感知到 Promise 被填充,返回并拿到 sc 对象
注:
_taskScheduler.postStartHook()在 SparkContext 构造函数中是固定调用的,Client 模式也会调用。但 Client 模式下YarnScheduler未重写该方法,走的是TaskSchedulerImpl的默认实现——调用waitBackendReady()等待 backend 就绪。YarnClusterScheduler重写了该方法,在waitBackendReady()之前插入了sparkContextInitialized通知逻辑。
3.4 YarnClusterSchedulerBackend.start()
1 | // YarnClusterSchedulerBackend 是 Yarn Cluster 模式的 SchedulerBackend |
与 Yarn Client 模式的 YarnClientSchedulerBackend.start()(向 RM 提交 Application)不同,Cluster 模式下 AM 已经运行在集群中,所以 start() 只是做绑定和初始化工作:
- 获取当前 AM 的 Attempt ID
bindToYarn()绑定 Application ID 和 Attempt IDsuper.start()调用CoarseGrainedSchedulerBackend.start()创建DriverEndpoint(RPC 端点,等待 Executor 注册)- 设置预期 Executor 总数
startBindings()完成绑定
3.5 sparkContextInitialized wait/notify 机制
ApplicationMaster.sparkContextInitialized(sc) 的实现:
1 | // 私有实例方法(通过伴生对象委托调用) |
wait/notify 配对流程:
1 | 用户线程 (Driver) AM 主线程 (runDriver) |
这种 wait/notify 机制确保:
- 用户线程在 SparkContext 初始化后暂停,等待 AM 完成注册和创建 Allocator
- 避免用户线程过早执行业务代码,而 Executor 尚未分配
- AM 初始化完后通过
resumeDriver()唤醒用户线程
阶段四:AM 注册与资源申请
4.1 registerAM()
runDriver() 中获取 SparkContext 后,调用 registerAM() 向 YARN RM 注册:
1 | // ApplicationMaster.registerAM() |
YarnRMClient.register() 内部调用 Hadoop YARN 的 AMRMClient.registerApplicationMaster() 向 ResourceManager 注册,建立心跳连接。
4.2 createAllocator()
registerAM() 之后,调用 createAllocator() 创建资源分配器(与第五篇中 Yarn Client 模式的 runExecutorLauncher() 中的调用方法相同):
1 | // 创建 YarnAllocator 并启动资源申请流程 |
YarnRMClient.createAllocator() 创建 YarnAllocator 实例,然后:
- 设置 AM RPC 端点
YarnAM,接收 Driver 发来的命令(如调整 Executor 数量) - 立即执行一次资源分配
allocator.allocateResources() - 启动汇报线程
launchReporterThread(),通过YarnAllocator.allocateResources()定期向 RM 申请 Container
4.3 AM 申请 Container 启动 Executor(与 Client 模式相同)
Executor 资源的申请、分配和启动流程与 Yarn Client 模式完全一致(详见第五篇文章 [3.6-3.9 节]):
1 | allocateResources() |
与 Client 模式的区别:
- Client 模式:
YarnClientSchedulerBackend.start()调用client.submitApplication()提交 Application,AM 启动后 Driver 已在客户端运行,AM 通过runExecutorLauncher()直接申请 Executor - Cluster 模式:
YarnClusterApplication.start()调用Client.run()→submitApplication(),AM 启动后 Driver 在 AM 内,AM 通过runDriver()先等 Driver 初始化完再申请 Executor
以上流程与 Client 模式完全一致,详见第五篇文章 [3.6-3.9 节]。
Yarn Cluster 模式完整调用链
1 | 客户端: |
Yarn Cluster vs Yarn Client 对比
流程对比
1 | Yarn Client: |
参数对比
| 参数 | Yarn Client | Yarn Cluster |
|---|---|---|
| childMainClass | args.mainClass(用户类) | YarnClusterApplication |
| amClass | ExecutorLauncher | ApplicationMaster |
| AM 命令是否含 –class | 否(args.userClass == null) | 是(args.userClass != null) |
| isClusterMode | false | true |
SparkContext 创建阶段对比
| 创建阶段 | Yarn Client | Yarn Cluster |
|---|---|---|
| TaskScheduler | YarnScheduler | YarnClusterScheduler |
| SchedulerBackend | YarnClientSchedulerBackend | YarnClusterSchedulerBackend |
| postStartHook | 无特殊行为 | 调用 ApplicationMaster.sparkContextInitialized(sc) |
| SchedulerBackend.start() | YarnClientSchedulerBackend: 创建 Client 并向 RM 提交 Application | YarnClusterSchedulerBackend: 绑定 AM ID,初始化 DriverEndpoint(不提交 Application) |
角色定位
| 角色 | Yarn Client | Yarn Cluster |
|---|---|---|
| Driver | 运行在提交客户端机器上 | 运行在 AM 容器内的独立线程 |
| AM | 负责申请 Executor(名为 ExecutorLauncher) | 负责启动 Driver + 申请 Executor |
| Executor | 运行在 NodeManager 容器中 | 运行在 NodeManager 容器中 |
| 客户端 | 需保持连接,响应速度快的 action 立即可见 | 提交后可断开,适合长时间运行的作业 |
总结
Yarn Cluster 模式的完整提交流程:
- 客户端提交:SparkSubmit →
YarnClusterApplication.start()→Client.run()→submitApplication(),向 YARN RM 提交 Application,AM 启动类为ApplicationMaster - AM 启动:
ApplicationMaster.main()→run()→runDriver() - 启动用户线程:
startUserApplication()在独立线程中反射调用用户主类的main方法,用户线程中创建 SparkContext - SparkContext 创建:
YarnClusterManager创建YarnClusterScheduler+YarnClusterSchedulerBackend - 通知 AM:
YarnClusterScheduler.postStartHook()调用ApplicationMaster.sparkContextInitialized(sc),通过 Promise 通知 AM 主线程 - wait/notify:用户线程初始化完 SparkContext 后挂起(wait),等待 AM 完成注册和 Allocator 创建
- AM 注册:
registerAM()向 RM 注册 AM 并建立心跳 - 申请 Executor:
createAllocator()创建 YarnAllocator,通过allocateResources()申请 Container - Executor 启动与注册:
createAllocator中申请的 Container 分配成功后,被分配的 NodeManager 上启动YarnCoarseGrainedExecutorBackend进程,创建 RpcEnv/SparkEnv,向DriverEndpoint发送RegisterExecutor,Driver 确认后回复RegisteredExecutor,创建Executor计算引擎(步骤 9 与步骤 10-12 可能同时运行,但步骤 10-12 之间是顺序的) - 恢复用户线程:
resumeDriver()唤醒用户线程(用户线程被挂起在sparkContextPromise.wait()) - Driver 执行业务代码:用户线程恢复后,继续执行用户主类
main方法,创建 DataFrame/RDD - Job 提交与调度:调用 Action 算子(如
count、collect、save等),触发 DAGScheduler 划分 Stage,TaskScheduler 将 Task 调度到已注册的 Executor 上执行(Action 算子触发 Job 后需等待 Executor 已注册)
本质区别:两种模式的根本差异在于 Driver 的启动顺序与位置。
Yarn Client 模式:Driver 在客户端,先于 AM 启动。
app.start直接反射调用用户main方法(Driver 就在此线程中运行),Driver 初始化 SparkContext 时,在YarnClientSchedulerBackend.start()中才向 RM 提交 Application 启动 AM。因此 AM 启动时 Driver 已在运行,AM 直接通过runExecutorLauncher()申请 Executor,无需等待 Driver 初始化。Driver 与集群端流程始终同时运行。Yarn Cluster 模式:Driver 在 AM 容器内,晚于 AM 启动。
app.start→YarnClusterApplication.start()先向 RM 提交 Application 启动 AM,AM 进入runDriver()后通过startUserApplication()在独立线程中反射调用用户main方法启动 Driver。Driver(用户线程)初始化 SparkContext 后通过 Promise/wait/notify 机制通知 AM 主线程,AM 主线程收到通知后才执行registerAM()和createAllocator(),最后通过resumeDriver()唤醒用户线程继续执行业务代码。
下一次我们将分析 Standalone Client 模式的详细提交流程。