> For the complete documentation index, see [llms.txt](https://nag-9-s.gitbook.io/spark-notes/llms.txt). Markdown versions of documentation pages are available by appending `.md` to page URLs; this page is available as [Markdown](https://nag-9-s.gitbook.io/spark-notes/spark-execution-flow.md).

# Spark Execution Flow

HighPerformance Spark 1st Edition

**what happens when we start a** `SparkContext ?`

.

First, the driver program pings the cluster manager.

The cluster manager launches a number of Spark executors (JVMs shown as black boxes in the diagram) on the worker nodes of the cluster (shown as blue circles).

One node can have multiple Spark executors, but an executor cannot span multiple nodes.

An RDD will be evaluated across the executors in partitions (shown as red rectangles).

Each executor can have multiple partitions, but a partition cannot be spread across multiple executors.

![](https://3957382301-files.gitbook.io/~/files/v0/b/gitbook-legacy-files/o/assets%2F-LsbZSfH3V7-1Ycuzu80%2F-LsbZZL4J3VRHnXA0HfC%2F-LsbZ_RRzVJc0owyLn-C%2FstartApp.png?generation=1572621933712260\&alt=media)In the Spark lazy evaluation paradigm, a Spark application doesn’t “do anything” until the driver program calls an action.With each action, the Spark scheduler builds an execution graph and launches a*Spark job*.Each job consists of*stages*, which are steps in the transformation of the data needed to materialize the final RDD. Each stage consists of a collection of\_tasks\_that represent each parallel computation and are performed on the executors.

[Figure 2-5](file:///G:/ScrapBook%20Data/Hadoop/data/20170703002220/index.html#spark_app_treefig) shows a tree of the different components of a Spark application and how these correspond to the API calls. An application corresponds to starting a `SparkContext`/`SparkSession`. Each *application* may contain many jobs that correspond to one RDD action. Each *job* may contain several stages that correspond to each wide transformation. Each *stage* is composed of one or many tasks that correspond to a parallelizable unit of computation done in each stage. There is one *task* for each partition in the resulting RDD of that stage.

![](https://3957382301-files.gitbook.io/~/files/v0/b/gitbook-legacy-files/o/assets%2F-LsbZSfH3V7-1Ycuzu80%2F-LsbZZL4J3VRHnXA0HfC%2F-LsbZ_RTWM1XSaeauF8y%2FsparkAppFlow.png?generation=1572621933669095\&alt=media)

![](https://3957382301-files.gitbook.io/~/files/v0/b/gitbook-legacy-files/o/assets%2F-LsbZSfH3V7-1Ycuzu80%2F-LsbZZL4J3VRHnXA0HfC%2F-LsbZ_RV5pxuYR2Wtn3y%2FexeFlow1.png?generation=1572621938097500\&alt=media)

<https://mapr.com/blog/getting-started-spark-web-ui/>

![](https://3957382301-files.gitbook.io/~/files/v0/b/gitbook-legacy-files/o/assets%2F-LsbZSfH3V7-1Ycuzu80%2F-LsbZZL4J3VRHnXA0HfC%2F-LsbZ_RXXzXiGgHCGuGS%2Fflow1.png?generation=1572621935063307\&alt=media)
