乐于分享
好东西不私藏

Spark 源码 | SparkContext 初始化与 TaskScheduler 创建流程(四)

Spark 源码 | SparkContext 初始化与 TaskScheduler 创建流程(四)

前言

在第二篇文章 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.sparkContext

SparkSession.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_NAMEnewHeartbeatReceiver(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): (SchedulerBackendTaskScheduler) = {  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. 1. ServiceLoader.load() 加载所有实现了 ExternalClusterManager 接口的类
  2. 2. .filter(_.canCreate(url)) 过滤,只保留 canCreate 返回 true 的
  3. 3. 当 url = "yarn" 时:
    • • YarnClusterManager.canCreate("yarn") → true(声明能管理 yarn)
    • • 其他集群管理器(Mesos、K8s)的 canCreate("yarn") → false
  4. 4. 过滤后只剩下 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 模式创建的类

组件
类名
TaskScheduler
YarnScheduler(继承自 TaskSchedulerImpl)
SchedulerBackend
YarnClientSchedulerBackend

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. 1. SparkSession.getOrCreate() → 创建 SparkContext
  2. 2. SparkContext 构造函数 → 调用 createTaskScheduler
  3. 3. createTaskScheduler → 根据 master 和 deployMode 创建对应实现
  4. 4. ClusterManager → 通过 SPI 加载对应的 ExternalClusterManager(如 YarnClusterManager),根据 deployMode 创建:
    • • client 模式:YarnScheduler + YarnClientSchedulerBackend
    • • cluster 模式:YarnClusterScheduler + YarnClusterSchedulerBackend
  5. 5. TaskScheduler.start() → 启动 SchedulerBackend,与集群管理器通信

SparkContext 初始化是所有部署模式(Yarn Client/Cluster、Standalone Client/Cluster 等)的共同前置流程,理解此流程后分析各模式的具体提交流程将更清晰。

下一次

下一次我们将分析 Yarn Client 模式的详细提交流程,包括 YarnClientSchedulerBackend.start() 的详细分析以及 Executor 启动流程。

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