My Spark source core SparkContext walkthrough full record
My Spark source core SparkContext walkthrough full record
Dirver Program (SparkConf) package org.apache.spark.SparkConf
Master package org.apache.spark.deploy.master
SparkContext package org.apache.spark.SparkContext
Stage package org.apache.spark.scheduler.Stage
Task package org.apache.spark.scheduler.Task
DAGScheduler package org.apache.spark.scheduler
TaskScheduler package org.apache.spark.scheduler.TaskScheduler
TaskSchedulerImpl package org.apache.spark.scheduler
Worker package org.apache.spark.deploy.worker
Executor package org.apache.spark.executor
BlockManager package org.apache.spark.storage
TaskSet package org.apache.spark.scheduler
/ / Create and start the scheduler
Val (sched, ts) = SparkContext.createTaskScheduler (this, master)
_ schedulerBackend = sched
_ taskScheduler = ts
_ dagScheduler = new DAGScheduler (this)
_ heartbeatReceiver.send (TaskSchedulerIsSet)
/ * *
* Create a task scheduler based on a given master URL.
* Return a 2-tuple of the scheduler backend and the task scheduler.
, /
Private def createTaskScheduler (
Sc: SparkContext
Master: String): (SchedulerBackend, TaskScheduler) = {
Master match {
Case "local" = >
Instantiate a
Val scheduler = new TaskSchedulerImpl (sc)
Build the masterUrls:
Val masterUrls = localCluster.start ()
It is said to be a very critical backend:
Val backend = new SparkDeploySchedulerBackend (scheduler, sc, masterUrls)
Scheduler.initialize (backend)
Backend.shutdownCallback = (backend: SparkDeploySchedulerBackend) = > {
LocalCluster.stop ()
}
(backend, scheduler)