前言
在第二篇文章 Spark 源码 | SparkSubmit 提交流程分析(二) 中,我们分析了 SparkSubmit 的提交流程。在 Yarn Client 模式下:
• runMain方法会调用app.start(Yarn Client 模式下为 JavaMainApplication),在 app.start 中通过反射调用用户主类的 main 方法(具体调用过程在不同模式中有所不同,后续分析各模式时再详细说明)• 用户主类(如 SparkPi)的 main 方法通常会创建 SparkSession 或 SparkContext • 在创建 SparkContext 时,会创建 TaskScheduler 和 SchedulerBackend
本文说明:虽然本文以 Yarn Client 模式为例(YarnScheduler、YarnClientSchedulerBackend),但 SparkContext 初始化和 TaskScheduler 创建流程是各模式的公共前提。后续分析 Yarn Cluster、Standalone 等模式时可直接参考本篇内容。
为什么要先写这篇:TaskScheduler 和 SchedulerBackend 的创建是各模式的公共前提,内容较多。如果不先单独分析,在分析各模式提交流程时就需要同时包含这些内容,导致每篇文章都很长。因此本文先专门分析它们的创建过程,后续分析各模式时就可以直接聚焦于该模式特有的流程(如 Executor 启动)。
版本
Spark 3.2.3
SparkContext 初始化入口
用户代码通常这样创建 SparkContext:
val spark = SparkSession.builder() .appName("WordCount") .master("yarn") .getOrCreate()val sc = spark.sparkContextSparkSession.getOrCreate
defgetOrCreate(): SparkSession = synchronized {// ...val sparkContext = userSuppliedContext.getOrElse {SparkContext.getOrCreate(sparkConf) }// ...}可以看到 SparkSession.getOrCreate() 内部调用的是 SparkContext.getOrCreate(sparkConf)。
SparkContext.getOrCreate
defgetOrCreate(config: SparkConf): SparkContext = {SPARK_CONTEXT_CONSTRUCTOR_LOCK.synchronized {if (activeContext.get() == null) { setActiveContext(newSparkContext(config)) } else {if (config.getAll.nonEmpty) { logWarning("Using an existing SparkContext; some configuration may not take effect.") } } activeContext.get() }}最终会创建 new SparkContext(config)。
new SparkContext(config)
在 SparkContext 构造函数中,会进行大量的初始化工作,关键步骤如下:
1. 验证配置(master、app.name)
if (!_conf.contains("spark.master")) {thrownewSparkException("A master URL must be set in your configuration")}if (!_conf.contains("spark.app.name")) {thrownewSparkException("An application name must be set in your configuration")}2. 初始化环境、内存、序列化器等
// 创建 Spark 执行环境_env = createSparkEnv(_conf, isLocal, listenerBus)SparkEnv.set(_env)3. 注册 HeartbeatReceiver
// 必须在 createTaskScheduler 之前注册_heartbeatReceiver = env.rpcEnv.setupEndpoint(HeartbeatReceiver.ENDPOINT_NAME, newHeartbeatReceiver(this))4. 创建 TaskScheduler、SchedulerBackend 和 DAGScheduler
val (sched, ts) = SparkContext.createTaskScheduler(this, master, deployMode)_schedulerBackend = sched_taskScheduler = ts_dagScheduler = newDAGScheduler(this)5. 启动 TaskScheduler
_taskScheduler.start()整个初始化过程在 try-catch 块中,异常时调用 stop() 方法确保资源释放。
createTaskScheduler 方法
createTaskScheduler 是 SparkContext 内部方法,根据 master URL 和 deployMode 创建对应的 TaskScheduler 和 SchedulerBackend。
privatedefcreateTaskScheduler( sc: SparkContext, master: String, deployMode: String): (SchedulerBackend, TaskScheduler) = { master match {case"local" =>// Local 模式val scheduler = newTaskSchedulerImpl(sc, MAX_LOCAL_TASK_FAILURES, isLocal = true)val backend = newLocalSchedulerBackend(sc.getConf, scheduler, 1) scheduler.initialize(backend) (backend, scheduler)caseSPARK_REGEX(sparkUrl) =>// Standalone 模式val scheduler = newTaskSchedulerImpl(sc)val masterUrls = sparkUrl.split(",").map("spark://" + _)val backend = newStandaloneSchedulerBackend(scheduler, sc, masterUrls) scheduler.initialize(backend) (backend, scheduler)case masterUrl =>// 其他模式(YARN、Kubernetes、Mesos):通过 ExternalClusterManager 加载val cm = getClusterManager(masterUrl) match {caseSome(clusterMgr) => clusterMgrcaseNone => thrownewSparkException("Could not parse Master URL: '" + master + "'") }val scheduler = cm.createTaskScheduler(sc, masterUrl)val backend = cm.createSchedulerBackend(sc, masterUrl, scheduler) cm.initialize(scheduler, backend) (backend, scheduler) }}getClusterManager 加载外部集群管理器
对于 YARN、Kubernetes 等模式,通过 Java ServiceLoader 机制加载 ExternalClusterManager 实现类:
privatedefgetClusterManager(url: String): Option[ExternalClusterManager] = {val loader = Utils.getContextOrSparkClassLoader// ServiceLoader.load() 加载所有实现了 ExternalClusterManager 接口的类// .filter(_.canCreate(url)) 过滤,只保留 canCreate 返回 true 的 val serviceLoaders =ServiceLoader.load(classOf[ExternalClusterManager], loader).asScala.filter(_.canCreate(url))if (serviceLoaders.size > 1) {thrownewSparkException(s"Multiple external cluster managers registered for the url $url: $serviceLoaders") } serviceLoaders.headOption}ServiceLoader 机制:Spark 使用 Java SPI(Service Provider Interface)机制加载外部集群管理器。在对应模块的 META-INF/services 目录下创建文件。
例如 YARN 模块的配置文件(resource-managers/yarn/src/main/resources/META-INF/services/org.apache.spark.scheduler.ExternalClusterManager):
org.apache.spark.scheduler.cluster.YarnClusterManager当 master = "yarn" 时,getClusterManager("yarn") 会加载 YarnClusterManager。这是因为:
1. ServiceLoader.load()加载所有实现了ExternalClusterManager接口的类2. .filter(_.canCreate(url))过滤,只保留canCreate返回true的3. 当 url = "yarn"时:• YarnClusterManager.canCreate("yarn")→true(声明能管理 yarn)• 其他集群管理器(Mesos、K8s)的 canCreate("yarn")→false4. 过滤后只剩下 YarnClusterManager,返回它
Yarn 模式的集群管理器
YarnClusterManager
当 master = "yarn" 时,加载的是 YarnClusterManager,它根据 deployMode 创建不同的 TaskScheduler 和 SchedulerBackend:
classYarnClusterManagerextendsExternalClusterManager{overridedefcanCreate(masterURL: String): Boolean = { masterURL == "yarn" }overridedefcreateTaskScheduler(sc: SparkContext, masterURL: String): TaskScheduler = { sc.deployMode match {case"cluster" => newYarnClusterScheduler(sc)case"client" => newYarnScheduler(sc)case _ => thrownewSparkException(s"Unknown deploy mode '${sc.deployMode}' for Yarn") } }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) } }overridedefinitialize(scheduler: TaskScheduler, backend: SchedulerBackend): Unit = { scheduler.asInstanceOf[TaskSchedulerImpl].initialize(backend) }}Yarn Client 模式创建的类
YarnClientSchedulerBackend 继承关系:
YarnClientSchedulerBackend → YarnSchedulerBackend → CoarseGrainedSchedulerBackend → SchedulerBackend(接口)SparkContext 初始化 TaskScheduler
在 SparkContext 构造函数中,创建 TaskScheduler 后会调用 start() 方法:
// 创建并启动 schedulerval (sched, ts) = SparkContext.createTaskScheduler(this, master, deployMode)_schedulerBackend = sched_taskScheduler = ts_dagScheduler = newDAGScheduler(this)_heartbeatReceiver.ask[Boolean](TaskSchedulerIsSet)// 启动 TaskScheduler_taskScheduler.start()TaskScheduler.start() 调用链
对于 Yarn Client 模式,整个调用链如下:
1. YarnClusterManager 创建的对象:
// createTaskScheduler 创建val scheduler = newYarnScheduler(sc) // YarnScheduler 继承自 TaskSchedulerImpl// createSchedulerBackend 创建val backend = newYarnClientSchedulerBackend(...)// initialize 中设置关联cm.initialize(scheduler, backend) → TaskSchedulerImpl.initialize(backend) // 将 backend 赋值给 TaskSchedulerImpl 的 backend 成员2. SparkContext 启动:
_taskScheduler.start()3. YarnScheduler.start():
YarnScheduler 没有重写 start() 方法,继承自 TaskSchedulerImpl:
// TaskSchedulerImpl.start()overridedefstart(): Unit = { backend.start() // 这里的 backend 就是在 initialize() 中设置的 YarnClientSchedulerBackend// ...}4. YarnClientSchedulerBackend.start():
最终会调用 YarnClientSchedulerBackend.start(),主要完成:
• 创建 Client 对象 • 通过 client.submitApplication()向 YARN RM 提交 Application• 通过 bindToYarn()绑定 appId• 等待 Application 运行
关于 YarnClientSchedulerBackend.start() 的详细源码分析,将在下一篇文章《Yarn Client 模式提交流程》中详细介绍。
关键源码(详细分析见下一篇文章):
overridedefstart(): Unit = {super.start()val args = newClientArguments(argsArrayBuf.toArray) totalExpectedExecutors = SchedulerBackendUtils.getInitialTargetExecutorNumber(conf)// 创建 Client,向 YARN RM 提交 Application client = newClient(args, conf, sc.env.rpcEnv)// submitApplication() 提交 Application,返回 appId// bindToYarn() 绑定 appId bindToYarn(client.submitApplication(), None)// 等待 Application 运行 waitForApplication()// 启动监控线程 monitorThread = asyncMonitorApplication() monitorThread.start() startBindings()}完整流程图
用户代码 ↓SparkSession.builder.getOrCreate() ↓SparkContext.getOrCreate() ↓new SparkContext(config) ↓SparkContext.createTaskScheduler(this, master, deployMode) ↓master = "yarn", deployMode = "client" ↓加载 YarnClusterManager(ExternalClusterManager) ↓YarnClusterManager.createTaskScheduler → YarnScheduler ↓YarnClusterManager.createSchedulerBackend → YarnClientSchedulerBackend ↓YarnClusterManager.initialize → TaskSchedulerImpl.initialize(backend) ↓_taskScheduler.start() ↓schedulerBackend.start() ↓YarnClientSchedulerBackend 与 YARN RM 通信总结
本文分析了 SparkContext 初始化和 TaskScheduler 创建流程(以 Yarn Client 模式为例):
1. SparkSession.getOrCreate() → 创建 SparkContext 2. SparkContext 构造函数 → 调用 createTaskScheduler 3. createTaskScheduler → 根据 master 和 deployMode 创建对应实现 4. ClusterManager → 通过 SPI 加载对应的 ExternalClusterManager(如 YarnClusterManager),根据 deployMode 创建: • client 模式:YarnScheduler + YarnClientSchedulerBackend • cluster 模式:YarnClusterScheduler + YarnClusterSchedulerBackend 5. TaskScheduler.start() → 启动 SchedulerBackend,与集群管理器通信
SparkContext 初始化是所有部署模式(Yarn Client/Cluster、Standalone Client/Cluster 等)的共同前置流程,理解此流程后分析各模式的具体提交流程将更清晰。
下一次
下一次我们将分析 Yarn Client 模式的详细提交流程,包括 YarnClientSchedulerBackend.start() 的详细分析以及 Executor 启动流程。
🧐 分享、点赞、在看,给个3连击呗!👇
夜雨聆风