前言
在第五篇文章 Yarn Client 模式提交流程(五) 中,我们详细分析了 Yarn Client 模式的提交流程。本文继续分析 Yarn Cluster 模式的提交流程,重点介绍与 Client 模式的关键差异——Driver 如何在 AM 内启动。
版本
Spark 3.2.3
Yarn Cluster 模式概述
Yarn Cluster 模式下,Driver 运行在 YARN 集群的 ApplicationMaster 容器中,客户端提交后即可断开,作业在集群中独立运行。
// 提交命令示例spark-submit --master yarn --deploy-mode cluster --classorg.apache.spark.examples.SparkPi$SPARK_HOME/examples/jars/spark-examples_2.12-3.2.3.jar与 Yarn Client 模式的核心区别
app.start 直接反射调用用户 main 方法),Driver 初始化 SparkContext 时,在 YarnClientSchedulerBackend.start() 中提交 Application 启动 AM | app.start → YarnClusterApplication.start() → Client.run() → submitApplication() 提交 Application),AM 的 runDriver() 中再通过 startUserApplication() 反射调用用户 main 方法启动 Driver | |
Promise + wait/notifyYarnClusterScheduler.postStartHook() → sparkContextInitialized() 填充 Promise,通知 AM 主线程 awaitResult 返回;用户线程通过 wait() 挂起,等 AM 注册和申请完资源后通过 resumeDriver() → notify() 唤醒 |
Yarn Cluster 模式完整提交流程
阶段一:客户端提交(SparkSubmit → YarnClusterApplication)
1.1 SparkSubmit 准备
回顾第二篇文章 SparkSubmit 提交流程分析(二),prepareSubmitEnvironment 方法中,Yarn Cluster 模式的配置如下:
// Yarn Cluster 模式if (isYarnCluster) { childMainClass = YARN_CLUSTER_SUBMIT_CLASS// "org.apache.spark.deploy.yarn.YarnClusterApplication"if (args.primaryResource != SparkLauncher.NO_RESOURCE) { childArgs += ("--jar", args.primaryResource) } childArgs += ("--class", args.mainClass)if (args.childArgs != null) { args.childArgs.foreach { arg => childArgs += ("--arg", arg) } }}关键差异:
• Client 模式: childMainClass = args.mainClass,直接在客户端反射调用用户主类的 main 方法• Cluster 模式: childMainClass = "org.apache.spark.deploy.yarn.YarnClusterApplication",用户主类名和 jar 路径作为--class和--jar参数传入
1.2 runMain 加载 YarnClusterApplication
根据第二篇文章的分析,runMain 方法会:
1. 加载 childMainClass(即YarnClusterApplication)2. YarnClusterApplication实现了SparkApplication接口,直接实例化3. 调用 app.start(childArgs.toArray, sparkConf)
// runMain 内部,加载 childMainClass 后启动应用// 如果主类实现了 SparkApplication 接口,直接实例化(如 YarnClusterApplication)// 否则包装成 JavaMainApplication 通过反射调用 main 方法(如用户主类 SparkPi)val app: SparkApplication = if (classOf[SparkApplication].isAssignableFrom(mainClass)) { mainClass.getConstructor().newInstance().asInstanceOf[SparkApplication]} else {newJavaMainApplication(mainClass)}app.start(childArgs.toArray, sparkConf)要点:YarnClusterApplication 实现了 SparkApplication,所以走的是直接实例化分支,而不是 JavaMainApplication 包装。
1.3 YarnClusterApplication.start()
// Yarn Cluster 模式入口类,实现了 SparkApplication 接口// SparkSubmit.prepareSubmitEnvironment 将 childMainClass 设为此类private[spark] classYarnClusterApplicationextendsSparkApplication{overridedefstart(args: Array[String], conf: SparkConf): Unit = {// 清除已在 yarn 缓存中分发的文件配置,避免重复 conf.remove(JARS) conf.remove(FILES) conf.remove(ARCHIVES)// 创建 Client 实例并调用 run() 提交 Application// 第三个参数为 null,因为此时 RPC 环境还没建立newClient(newClientArguments(args), conf, null).run() }}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 启动命令不同。
defrun(): Unit = {this.appId = submitApplication()// fireAndForget 模式:不等待应用完成,打印一次状态就返回,默认 false// 由 spark.yarn.submit.waitAppCompletion=false 控制if (!launcherBackend.isConnected() && fireAndForget) {val report = getApplicationReport(appId)val state = report.getYarnApplicationState logInfo(s"Application report for $appId (state: $state)") logInfo(formatReportDetails(report, getDriverLogsLink(report)))if (state == YarnApplicationState.FAILED || state == YarnApplicationState.KILLED) {thrownewSparkException(s"Application $appId finished with status: $state") }// 默认行为:进入监控循环,轮询 YARN RM 获取应用状态,直到应用结束 } else {valYarnAppReport(appState, finalState, diags) = monitorApplication(appId)if (appState == YarnApplicationState.FAILED || finalState == FinalApplicationStatus.FAILED) { diags.foreach { err => logError(s"Application diagnostics message: $err") }thrownewSparkException(s"Application $appId finished with failed status") }if (appState == YarnApplicationState.KILLED || finalState == FinalApplicationStatus.KILLED) {thrownewSparkException(s"Application $appId is killed") }if (finalState == FinalApplicationStatus.UNDEFINED) {thrownewSparkException(s"The final status of application $appId is undefined") } }}run() 方法的执行逻辑:
1. 先提交:调用 submitApplication()向 YARN RM 提交 Application2. 再等待:根据 spark.yarn.submit.waitAppCompletion配置决定行为:• 默认等待( waitAppCompletion=true):进入monitorApplication循环,定期轮询 RM 获取应用状态,直到应用结束• fireAndForget 模式( waitAppCompletion=false):提交后立即返回,只打印一次状态(Cluster 模式下fireAndForget默认为 true,因为客户端提交后可断开)
submitApplication() 流程与第五篇分析的 Yarn Client 模式相同,但需要注意:两种模式下的 createContainerLaunchContext() 构造的 AM 启动命令不同(见下一节)。以下是完整源码注释:
defsubmitApplication(): ApplicationId = {// ① 校验 Spark 配置中的资源请求是否合法// 检查 spark.yarn.executor.resource.* 等配置项ResourceRequestHelper.validateResources(sparkConf)var appId: ApplicationId = nulltry {// ② 连接 LauncherBackend:用于接收客户端的停止请求// 当用户按下 Ctrl+C 或调用 kill 时,LauncherBackend 能收到通知 launcherBackend.connect()// ③ 初始化 YarnClient:YARN 官方 Hadoop API 的客户端// YarnClient 封装了与 RM 通信的底层细节 yarnClient.init(hadoopConf) yarnClient.start()// ④ 向 RM 请求创建新 Application,RM 返回 ApplicationIdval newApp = yarnClient.createApplication()val newAppResponse = newApp.getNewApplicationResponse() appId = newAppResponse.getApplicationId()// ⑤ 创建 Container 启动上下文// createContainerLaunchContext() 中指定了 AM 启动命令、环境变量、classpath// 这是两种模式的关键差异点——决定了 AM 进程启动什么类、带什么参数val containerContext = createContainerLaunchContext(newAppResponse)// ⑥ 创建 Application 提交上下文// createApplicationSubmissionContext() 中指定了应用名称、队列、资源需求等val appContext = createApplicationSubmissionContext(newApp, containerContext)// ⑦ 正式提交 Application 到 RM// 提交后 RM 会调度资源,在某个 NodeManager 上启动 AM Container logInfo(s"Submitting application $appId to ResourceManager") yarnClient.submitApplication(appContext)// ⑧ 设置 appId,通知 LauncherBackend 提交已完成 launcherBackend.setAppId(appId.toString) reportLauncherState(SparkAppHandle.State.SUBMITTED) appId } catch {case e: Throwable =>// ⑨ 提交失败时清理 staging 目录if (stagingDirPath != null) { cleanupStagingDir() }throw e }}与 Client 模式的关键差异隐藏在 createContainerLaunchContext 中。
1.5 createContainerLaunchContext 模式差异
createContainerLaunchContext 方法中根据 isClusterMode 决定 AM 启动类:
val amClass =if (isClusterMode) {// Cluster 模式:直接启动 ApplicationMasterUtils.classForName("org.apache.spark.deploy.yarn.ApplicationMaster").getName } else {// Client 模式:启动 ExecutorLauncher(内部调用 ApplicationMaster.main)Utils.classForName("org.apache.spark.deploy.yarn.ExecutorLauncher").getName }以下是在 createContainerLaunchContext 中构建 AM 启动命令的完整源码分析。核心分为三步:JVM 参数 → AM 参数 → 拼接最终命令。
① JVM 参数(javaOpts)
val javaOpts = ListBuffer[String]()// AM 堆内存:由 spark.driver.memory(Cluster)或 spark.yarn.am.memory(Client)控制javaOpts += "-Xmx" + amMemory + "m"// 临时目录javaOpts += "-Djava.io.tmpdir=" + tmpDir// Cluster 模式下添加 Driver 特有的 JVM 参数if (isClusterMode) { sparkConf.get(DRIVER_JAVA_OPTIONS).foreach { opts => javaOpts ++= Utils.splitCommandString(opts) .map(Utils.substituteAppId(_, appId.toString)) .map(YarnSparkHadoopUtil.escapeForShell) }// Driver Library Pathval libraryPaths = Seq(sparkConf.get(DRIVER_LIBRARY_PATH), sys.props.get("spark.driver.libraryPath")).flattenif (libraryPaths.nonEmpty) { prefixEnv = Some(createLibraryPathPrefix( libraryPaths.mkString(File.pathSeparator), sparkConf)) }}// 日志目录参数(所有模式通用)javaOpts += ("-Dspark.yarn.app.container.log.dir=" +ApplicationConstants.LOG_DIR_EXPANSION_VAR)② AM 参数(amArgs)
// 启动类:Cluster → ApplicationMaster,Client → ExecutorLauncherval amClass =if (isClusterMode) {Utils.classForName("org.apache.spark.deploy.yarn.ApplicationMaster").getName } else {Utils.classForName("org.apache.spark.deploy.yarn.ExecutorLauncher").getName }// --class:仅 Cluster 模式有值(用户主类名),Client 模式为 Nilval userClass =if (isClusterMode) {Seq("--class", YarnSparkHadoopUtil.escapeForShell(args.userClass)) } else {Nil }// --jar:仅在用户传了 jar 时有值val userJar =if (args.userJar != null) {Seq("--jar", args.userJar) } else {Nil }除了 --jar 参数外,jar 文件还通过 YARN LocalResource 机制分发给 AM。createContainerLaunchContext 前半部分调用了 prepareLocalResources,将 jar 上传到 HDFS staging 目录:
// createContainerLaunchContext 内部,准备 AM 所需的本地资源// 将 jar、配置文件等上传到 HDFS staging 目录,作为 YARN LocalResourceval localResources = prepareLocalResources(stagingDirPath, pySparkArchives)// 将 LocalResource 设置到 ContainerLaunchContext 中amContainer.setLocalResources(localResources.asJava)prepareLocalResources 中对 args.userJar 的处理:
// prepareLocalResources 内部Option(args.userJar).filter(_.trim.nonEmpty).foreach { jar =>// 将 jar 作为 LocalResource 上传(destName = APP_JAR_NAME = "__app__.jar")val (isLocal, localizedPath) = distribute(jar, destName = Some(APP_JAR_NAME))if (isLocal) {// 如果是本地文件,distribute 会上传到 HDFS staging 目录// 设置 spark.yarn.app.jar 供后续使用 sparkConf.set(APP_JAR, localizedPath) }}因此 jar 通过两个途径到达 AM:
1. --jar参数:args.userJar原始路径作为命令行参数传入,ApplicationMasterArguments解析后使用2. YARN LocalResource:由 prepareLocalResources上传到 HDFS staging 目录,AM 启动时 YARN 自动本地化到容器文件系统中
// 拼接完整 AM 参数val amArgs =Seq(amClass) ++ userClass ++ userJar ++ primaryPyFile ++ primaryRFile ++ userArgs ++Seq("--properties-file", buildPath(Environment.PWD.$$(), LOCALIZED_CONF_DIR, SPARK_CONF_FILE)) ++Seq("--dist-cache-conf", buildPath(Environment.PWD.$$(), LOCALIZED_CONF_DIR, DIST_CACHE_CONF_FILE))③ 最终命令拼接
val commands = prefixEnv ++Seq(Environment.JAVA_HOME.$$() + "/bin/java", "-server") ++ javaOpts ++ amArgs ++Seq("1>", ApplicationConstants.LOG_DIR_EXPANSION_VAR + "/stdout","2>", ApplicationConstants.LOG_DIR_EXPANSION_VAR + "/stderr")因此 Cluster 模式 AM 的实际启动命令大致为:
JAVA_HOME/bin/java -server -Xmx<amMemory>m \ -Djava.io.tmpdir=<tmpDir> \ -Dspark.yarn.app.container.log.dir=<logDir> \ org.apache.spark.deploy.yarn.ApplicationMaster \ ← amClass(Client → ExecutorLauncher) --class org.apache.spark.examples.SparkPi \ ← userClass(仅 Cluster) --jar $SPARK_HOME/examples/jars/spark-examples_2.12-3.2.3.jar \ ← userJar(仅 Cluster) --arg <userArgs> \ ← 用户参数 --properties-file <conf_dir>/__spark_conf__.properties \ ← Spark 配置 --dist-cache-conf <conf_dir>/__spark_dist_cache__.properties \ ← 分布式缓存配置 1> <logDir>/stdout 2> <logDir>/stderr这里 --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():
defmain(args: Array[String]): Unit = {val amArgs = newApplicationMasterArguments(args)val sparkConf = newSparkConf()// ... master = newApplicationMaster(amArgs, sparkConf, yarnConf) ugi.doAs(newPrivilegedExceptionAction[Unit]() {overridedefrun(): Unit = System.exit(master.run()) })}ApplicationMasterArguments 解析 AM 命令行参数,支持:
--jar JAR_PATH | |
--class CLASS_NAME | |
--arg ARG | |
--properties-file FILE | |
--dist-cache-conf |
2.2 ApplicationMaster.run()
// ApplicationMaster 主入口方法,根据模式分支执行finaldefrun(): Int = {// Cluster 模式下设置 AM 进程的系统属性// Client 模式 Driver 已在客户端,AM 不需要这些设置val attemptID = if (isClusterMode) {// 设置 Web UI 端口为随机,避免集群中端口冲突if (System.getProperty(UI_PORT.key) == null) {System.setProperty(UI_PORT.key, "0") }System.setProperty("spark.master", "yarn")System.setProperty(SUBMIT_DEPLOY_MODE.key, "cluster")System.setProperty("spark.yarn.app.id", appAttemptId.getApplicationId().toString())Option(appAttemptId.getAttemptId.toString) } else {None }// ...// Cluster 模式:runDriver() 先启动 Driver 线程再申请 Executor// Client 模式:runExecutorLauncher() 直接申请 Executor(Driver 已在客户端)if (isClusterMode) { runDriver() } else { runExecutorLauncher() }}Cluster 模式独有的初始化:
1. 设置 Web UI 端口为随机端口(避免与集群中其他 Spark 进程冲突) 2. 设置 spark.master = yarn3. 设置 spark.submit.deployMode = cluster4. 保存 spark.yarn.app.id
2.3 runDriver() — Cluster 模式核心方法
privatedefrunDriver(): Unit = { addAmIpFilter(None, System.getenv(ApplicationConstants.APPLICATION_WEB_PROXY_BASE_ENV))// ① 启动用户主类线程(独立的 Driver 线程) userClassThread = startUserApplication()// ② 等待 SparkContext 初始化(跨线程同步) logInfo("Waiting for spark context initialization...")val totalWaitTime = sparkConf.get(AM_MAX_WAIT_TIME)try {// awaitResult 阻塞当前(AM 主)线程,等待用户线程中的 SparkContext 初始化完毕。// 用户线程执行到 SparkContext 构造函数末尾 → postStartHook()// → YarnClusterScheduler.postStartHook() → sparkContextInitialized(sc)// → Promise.success(sc),此时 awaitResult 返回,AM 主线程拿到 sc 对象。// 拿到 sc 后主线程才能获取 rpcEnv/DRIVER_PORT 等信息去注册 AM 和申请 Executor。val sc = ThreadUtils.awaitResult(sparkContextPromise.future,Duration(totalWaitTime, TimeUnit.MILLISECONDS))if (sc != null) {val rpcEnv = sc.env.rpcEnvval userConf = sc.getConfval host = userConf.get(DRIVER_HOST_ADDRESS)val port = userConf.get(DRIVER_PORT)// ③ 向 YARN RM 注册 AM registerAM(host, port, userConf, sc.ui.map(_.webUrl), appAttemptId)// ④ 连接 Driver 的 YarnSchedulerBackendval driverRef = rpcEnv.setupEndpointRef(RpcAddress(host, port),YarnSchedulerBackend.ENDPOINT_NAME)// ⑤ 创建资源分配器 createAllocator(driverRef, userConf, rpcEnv, appAttemptId, distCacheConf) } else {thrownewIllegalStateException("User did not initialize spark context!") }// ⑥ 恢复用户线程执行 resumeDriver() userClassThread.join() } catch {case e: SparkExceptionif e.getCause().isInstanceOf[TimeoutException] =>// SparkContext 初始化超时 finish(FinalApplicationStatus.FAILED,ApplicationMaster.EXIT_SC_NOT_INITED,"Timed out waiting for SparkContext.") } finally { resumeDriver() }}执行流程(时序角度):
runDriver() 用户线程 (startUserApplication) | | |---- startUserApplication() --->| 启动用户主类(如 SparkPi) | | | | SparkContext 初始化 | | |<--- sparkContextInitialized ---| 通知 AM(挂起自己) | | |---- registerAM() ------------->| 注册到 RM |---- createAllocator() -------->| 创建资源分配器 |---- resumeDriver() ----------->| 恢复用户线程 | | | | 用户代码继续执行 | | |<--- userClassThread.join() ----| 等待用户线程结束阶段三:用户线程与 SparkContext 初始化
3.1 startUserApplication() 启动用户主类
// 在独立线程中启动用户主类,返回用户线程引用privatedefstartUserApplication(): Thread = { logInfo("Starting the user application in a separate Thread")var userArgs = args.userArgs// ...// 通过反射获取用户主类的 main 方法(如 SparkPi.main)val mainMethod = userClassLoader.loadClass(args.userClass) .getMethod("main", classOf[Array[String]])// 在独立线程中执行用户 main 方法,这个线程就是 Driverval userThread = newThread {overridedefrun(): Unit = {try {// 调用用户主类的 main 方法(如 SparkPi) mainMethod.invoke(null, userArgs.toArray)// 正常结束 → 标记成功 finish(FinalApplicationStatus.SUCCEEDED, ApplicationMaster.EXIT_SUCCESS) } catch {case e: InvocationTargetException => e.getCause match {case _: InterruptedException =>caseSparkUserAppException(exitCode) => finish(FinalApplicationStatus.FAILED, exitCode, msg)case cause: Throwable => finish(FinalApplicationStatus.FAILED,ApplicationMaster.EXIT_EXCEPTION_USER_CLASS, ...) }// 异常时通过 Promise 通知 runDriver() 中 awaitResult 等待的 AM 主线程 sparkContextPromise.tryFailure(e.getCause()) } finally { sparkContextPromise.trySuccess(null) } } } userThread.setContextClassLoader(userClassLoader) userThread.setName("Driver") userThread.start() userThread}关键点:
• 用户主类(如 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 会创建:
// YarnClusterManageroverridedefcreateTaskScheduler(sc: SparkContext, masterURL: String): TaskScheduler = { sc.deployMode match {case"cluster" => newYarnClusterScheduler(sc)case"client" => newYarnScheduler(sc) }}overridedefcreateSchedulerBackend(sc: SparkContext, masterURL: String, scheduler: TaskScheduler): SchedulerBackend = { sc.deployMode match {case"cluster" =>newYarnClusterSchedulerBackend(scheduler.asInstanceOf[TaskSchedulerImpl], sc)case"client" =>newYarnClientSchedulerBackend(scheduler.asInstanceOf[TaskSchedulerImpl], sc) }}YarnCluster 模式创建的对象:
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。
// YarnClusterScheduler 是 Yarn Cluster 模式的 TaskScheduler// 核心作用:在 postStartHook() 中通知 AM 主线程 SparkContext 已初始化private[spark] classYarnClusterScheduler(sc: SparkContext) extendsYarnScheduler(sc) { logInfo("Created YarnClusterScheduler")overridedefpostStartHook(): Unit = {// 信号枪:SparkContext 初始化完毕时(构造函数末尾),此处被回调。// 运行在用户线程中,填充 Promise 后 AM 主线程的 awaitResult 返回,// 主线程即可获取 sc.rpcEnv / port 去执行 registerAM 和 createAllocator。ApplicationMaster.sparkContextInitialized(sc)super.postStartHook() logInfo("YarnClusterScheduler.postStartHook done") }}ApplicationMaster.sparkContextInitialized(sc) 的完整调用链(伴生对象委托给实例方法):
// 伴生对象静态方法(YarnClusterScheduler.postStartHook 调用此方法)private[spark] defsparkContextInitialized(sc: SparkContext): Unit = { master.sparkContextInitialized(sc)}// 私有实例方法privatedefsparkContextInitialized(sc: SparkContext) = { sparkContextPromise.synchronized { sparkContextPromise.success(sc) // 填充 Promise,通知 AM 主线程 sparkContextPromise.wait() // 挂起用户线程,等待 resumeDriver() 唤醒 }}调用链:
1. SparkContext构造函数中顺序执行:先调用_taskScheduler.start(),再调用_taskScheduler.postStartHook()2. _taskScheduler.start()内部调用backend.start()完成 DriverEndpoint 创建3. _taskScheduler.postStartHook()由于多态,实际调用的是YarnClusterScheduler.postStartHook()4. YarnClusterScheduler.postStartHook()内部调用ApplicationMaster.sparkContextInitialized(sc)5. 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()
// YarnClusterSchedulerBackend 是 Yarn Cluster 模式的 SchedulerBackend// 与 Client 模式的 YarnClientSchedulerBackend 不同,不提交 Application// 因为 AM 已在集群中运行,start() 只做绑定和初始化private[spark] classYarnClusterSchedulerBackend( scheduler: TaskSchedulerImpl, sc: SparkContext)extendsYarnSchedulerBackend(scheduler, sc) {overridedefstart(): Unit = {val attemptId = ApplicationMaster.getAttemptId bindToYarn(attemptId.getApplicationId(), Some(attemptId))super.start() // 创建 DriverEndpoint,等待 Executor 注册 totalExpectedExecutors = SchedulerBackendUtils.getInitialTargetExecutorNumber(sc.conf) startBindings() }}与 Yarn Client 模式的 YarnClientSchedulerBackend.start()(向 RM 提交 Application)不同,Cluster 模式下 AM 已经运行在集群中,所以 start() 只是做绑定和初始化工作:
1. 获取当前 AM 的 Attempt ID 2. bindToYarn()绑定 Application ID 和 Attempt ID3. super.start()调用CoarseGrainedSchedulerBackend.start()创建DriverEndpoint(RPC 端点,等待 Executor 注册)4. 设置预期 Executor 总数 5. startBindings()完成绑定
3.5 sparkContextInitialized wait/notify 机制
ApplicationMaster.sparkContextInitialized(sc) 的实现:
// 私有实例方法(通过伴生对象委托调用)privatedefsparkContextInitialized(sc: SparkContext) = { sparkContextPromise.synchronized { sparkContextPromise.success(sc) // 填充 Promise,通知 runDriver() 中 awaitResult 等待的 AM 主线程 sparkContextPromise.wait() // 挂起当前线程(用户线程) }}wait/notify 配对流程:
用户线程 (Driver) AM 主线程 (runDriver) | | | SparkContext 初始化完毕 | | postStartHook() | | sparkContextInitialized(sc) | | → success(sc) --------通知-------------> | awaitResult 返回 | → wait() 挂起 | | | registerAM() | | createAllocator() | | resumeDriver() |<-------- 唤醒 (notify) --------------------| | | | 继续执行业务代码 |这种 wait/notify 机制确保:
1. 用户线程在 SparkContext 初始化后暂停,等待 AM 完成注册和创建 Allocator 2. 避免用户线程过早执行业务代码,而 Executor 尚未分配 3. AM 初始化完后通过 resumeDriver()唤醒用户线程
阶段四:AM 注册与资源申请
4.1 registerAM()
runDriver() 中获取 SparkContext 后,调用 registerAM() 向 YARN RM 注册:
// ApplicationMaster.registerAM()privatedefregisterAM( host: String, port: Int, _sparkConf: SparkConf, uiAddress: Option[String], appAttemptId: ApplicationAttemptId): Unit = {// 创建 YarnRMClient(封装 AMRMClient 的注册/注销逻辑) client = newYarnRMClient(appAttemptId)// 向 RM 注册 AM,建立心跳 client.register(host, port, uiAddress)// ...}YarnRMClient.register() 内部调用 Hadoop YARN 的 AMRMClient.registerApplicationMaster() 向 ResourceManager 注册,建立心跳连接。
4.2 createAllocator()
registerAM() 之后,调用 createAllocator() 创建资源分配器(与第五篇中 Yarn Client 模式的 runExecutorLauncher() 中的调用方法相同):
// 创建 YarnAllocator 并启动资源申请流程privatedefcreateAllocator(driverRef: RpcEndpointRef, _sparkConf: SparkConf, rpcEnv: RpcEnv, appAttemptId: ApplicationAttemptId, distCacheConf: String): Unit = {// 构建 Driver 的 RPC 地址,供 Executor 反向注册使用val driverUrl = RpcEndpointAddress(driverRef.address.host, driverRef.address.port,CoarseGrainedSchedulerBackend.ENDPOINT_NAME).toString// 创建 YarnAllocator,负责与 RM 通信申请 Container allocator = client.createAllocator(yarnConf, _sparkConf, appAttemptId, driverUrl, driverRef, securityMgr, localResources)// 注册 AM RPC 端点,接收 Driver 发来的命令 rpcEnv.setupEndpoint("YarnAM", newAMEndpoint(rpcEnv, driverRef))// 立即执行一次资源分配 allocator.allocateResources()// 启动后台汇报线程,定期向 RM 申请 Container reporterThread = launchReporterThread()}YarnRMClient.createAllocator() 创建 YarnAllocator 实例,然后:
1. 设置 AM RPC 端点 YarnAM,接收 Driver 发来的命令(如调整 Executor 数量)2. 立即执行一次资源分配 allocator.allocateResources()3. 启动汇报线程 launchReporterThread(),通过YarnAllocator.allocateResources()定期向 RM 申请 Container
4.3 AM 申请 Container 启动 Executor(与 Client 模式相同)
Executor 资源的申请、分配和启动流程与 Yarn Client 模式完全一致(详见第五篇文章 [3.6-3.9 节]):
allocateResources() → handleAllocatedContainers() → runAllocatedContainers() → ExecutorRunnable.run() → startContainer() → NMClient → NodeManager → YarnCoarseGrainedExecutorBackend.main() → CoarseGrainedExecutorBackend.run() → 创建 RpcEnv / SparkEnv → setupEndpoint("Executor", backend) → onStart() → 发送 RegisterExecutor → DriverEndpoint 接收并注册 → 回复 RegisteredExecutor → new Executor(...)与 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 模式完整调用链
客户端:SparkSubmit.main → submit → runMain → 加载 YarnClusterApplication(实现了 SparkApplication) → app.start → YarnClusterApplication.start() → new Client(args, conf, null).run() → client.submitApplication() → createContainerLaunchContext() → amClass = ApplicationMaster(含 --class SparkPi) → yarnClient.submitApplication(appContext) → YARN RM 接收 → monitorApplication(appId) → 等待结束AM 容器 (集群):ApplicationMaster.main → ApplicationMaster.run() → runDriver() [1] startUserApplication() → 反射调用 SparkPi.main("Driver" 线程) → SparkContext 初始化 → createTaskScheduler("yarn", "cluster") → YarnClusterScheduler → YarnClusterSchedulerBackend → _taskScheduler.start() → YarnClusterSchedulerBackend.start() → bindToYarn() → super.start() → 创建 DriverEndpoint → postStartHook() → ApplicationMaster.sparkContextInitialized(sc) → sparkContextPromise.success(sc) → 通知 AM 主线程 → sparkContextPromise.wait() → 挂起用户线程 [2] awaitResult 返回(收到 SparkContext) [3] registerAM(host, port, ...) → YarnRMClient.register() → amClient.registerApplicationMaster() → 向 RM 注册 [4] createAllocator() → YarnAllocator.allocateResources() → handleAllocatedContainers() → runAllocatedContainers() → ExecutorRunnable.run() → startContainer() → NMClient → NodeManager 启动 YarnCoarseGrainedExecutorBackend [5] resumeDriver() → sparkContextPromise.notify() → 唤醒用户线程 [6] userClassThread.join() → 等待用户代码结束用户线程恢复执行:SparkPi 继续执行 → DAGScheduler.runJob → Job 提交与调度Yarn Cluster vs Yarn Client 对比
流程对比
Yarn Client:客户端: SparkSubmit → 反射调用用户主类 (Driver 在客户端) → SparkContext → _taskScheduler.start() → YarnClientSchedulerBackend.start() → Client.submitApplication()AM: ExecutorLauncher → ApplicationMaster.main → runExecutorLauncher() → registerAM() + createAllocator() → 申请 ExecutorYarn Cluster:客户端: SparkSubmit → YarnClusterApplication.start() → Client.run() → submitApplication() ← 仅提交AM: ApplicationMaster → runDriver() → startUserApplication() ← 启动 Driver 线程 → 等待 SparkContext 初始化 → registerAM() + createAllocator() ← 再申请 Executor → 恢复 Driver 线程参数对比
SparkContext 创建阶段对比
角色定位
总结
Yarn Cluster 模式的完整提交流程:
1. 客户端提交:SparkSubmit → YarnClusterApplication.start()→Client.run()→submitApplication(),向 YARN RM 提交 Application,AM 启动类为ApplicationMaster2. AM 启动: ApplicationMaster.main()→run()→runDriver()3. 启动用户线程: startUserApplication()在独立线程中反射调用用户主类的main方法,用户线程中创建 SparkContext4. SparkContext 创建: YarnClusterManager创建YarnClusterScheduler+YarnClusterSchedulerBackend5. 通知 AM: YarnClusterScheduler.postStartHook()调用ApplicationMaster.sparkContextInitialized(sc),通过 Promise 通知 AM 主线程6. wait/notify:用户线程初始化完 SparkContext 后挂起(wait),等待 AM 完成注册和 Allocator 创建 7. AM 注册: registerAM()向 RM 注册 AM 并建立心跳8. 申请 Executor: createAllocator()创建 YarnAllocator,通过allocateResources()申请 Container9. Executor 启动与注册: createAllocator中申请的 Container 分配成功后,被分配的 NodeManager 上启动YarnCoarseGrainedExecutorBackend进程,创建 RpcEnv/SparkEnv,向DriverEndpoint发送RegisterExecutor,Driver 确认后回复RegisteredExecutor,创建Executor计算引擎(步骤 9 与步骤 10-12 可能同时运行,但步骤 10-12 之间是顺序的)10. 恢复用户线程: resumeDriver()唤醒用户线程(用户线程被挂起在sparkContextPromise.wait())11. Driver 执行业务代码:用户线程恢复后,继续执行用户主类 main方法,创建 DataFrame/RDD12. 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 模式的详细提交流程。
🧐 分享、点赞、在看,给个3连击呗!👇
夜雨聆风