What you will be able to do
- Describe a Spark application as one driver program and its own set of executors
- Explain that a job is created only when an action runs, never by a transformation
- Explain how a job is split into stages at shuffle boundaries, and why each stage waits for the one before it
- Relate tasks to partitions, executors and CPU cores, and read a stage's task count as a sign of parallelism
Key concept
Execution hierarchy (application → job → stage → task) — A Spark application runs one job for each action. The driver splits each job into stages wherever data has to be shuffled, and splits each stage into tasks, one per partition. Executors run those tasks in parallel.
1.The top level: one application, one driver, its own executors
The execution hierarchy starts with the application. When you run a Spark program, through spark-submit, a notebook or an interactive shell, you start one application. It is made of two kinds of process. The driver program runs your main code and holds the SparkContext (in modern code, the SparkContext that sits under your SparkSession). The executors are processes on worker nodes that do the real computation. The driver plans the work, and the executors carry it out. Everything lower in the hierarchy (jobs, stages, tasks) belongs to one application and is scheduled by that application's driver.
Starting up follows a fixed order. The SparkContext connects to a cluster manager (Standalone, YARN or Kubernetes). The cluster manager allocates resources. Spark then acquires executors on worker nodes, ships your application code (JARs or Python files) to them, and finally sends them tasks. The executors stay up for the whole life of the application and run tasks on multiple threads, so one executor can run several tasks at the same time, one per CPU core it has been given.
Checkpoint 1 of 7· Put it in order
Put these steps of running an application on a cluster in order.
- 1.The SparkContext in the driver connects to a cluster manager
- 2.The application code is sent to the executors
- 3.Spark acquires executors on nodes in the cluster
- 4.The SparkContext sends tasks to the executors to run
Spark needs a cluster manager before it can get executors, and needs executors before it can ship code to them. Tasks are sent last, once the executors have the code they need to run.
“Once connected, Spark acquires executors on nodes in the cluster, which are processes that run computations and store data for your application.”Source: spark.apache.org
Executors are not shared. Each application gets its own executors, and each driver schedules only its own tasks. This keeps applications apart: a task from application A never runs in the same JVM as a task from application B. The cost is that two applications cannot share in-memory data directly. They have to go through external storage. Databricks' Spark UI guidance also says to check that tasks are spread across several executors, because that spread is what gives you parallelism.
Checkpoint 2 of 7· Check yourself
Two separate Spark applications are submitted to the same YARN cluster. Which statement is true?
Executors are launched per application, and each driver schedules its own tasks. Because of this isolation, data cannot be shared between applications without external storage.
“Each application has its own executors.”Source: spark.apache.org
2.Jobs: one per action
Below the application is the job. A job is a parallel computation that Spark starts in response to an action, such as collect, count, show or a write. Transformations such as filter, select, join and groupBy do not start jobs. They only add steps to the plan. The Databricks FAQ puts the split this way: transformations add processing logic to the plan, and actions make that logic run and produce a result. So the number of jobs an application runs depends on how many actions it calls, not on how many lines of transformation code it has. One application can run many jobs, one after another or at the same time from several threads, and they all share the application's executors.
This matters when you read the Spark UI. The Jobs tab lists one entry per action, and its event timeline shows when each job was submitted. A job shows as running from the moment it is submitted, even if none of its tasks have started yet. So several jobs submitted together from different threads can all appear to start at the same moment. That is expected, and on its own it does not mean they are getting executor time.
Checkpoint 3 of 7· Exam question
A data engineer runs the following PySpark code on a DataFrame `orders_df` that has not been cached or persisted: ```python filtered_df = orders_df.filter(orders_df.status == "SHIPPED") mapped_df = filtered_df.withColumn("total", filtered_df.qty * filtered_df.price) count_result = mapped_df.count() collected_rows = mapped_df.collect() ``` How many Spark jobs does this code trigger?
Correct answer: C — 2 jobs
- A. This is incorrect because `count()` and `collect()` are both actions, and every action that requests a result from the driver triggers at least one job.
- B. This is incorrect because `filter` and `withColumn` are lazy transformations that build a plan but do not run it; only `count()` and `collect()` force execution, and there are two of them.
- C. This is correct because `count()` and `collect()` are each separate actions, and since `mapped_df` was never cached, each action independently triggers its own job that re-reads and re-applies the filter and column transformations from `orders_df`.
- D. This is incorrect because there are only two actions in the snippet; `filter` and `withColumn` are transformations and do not add to the job count on their own.
- E. This is incorrect because it overcounts the actions present; only `count()` and `collect()` request a result, so only two jobs are triggered.
3.Stages: where a job is cut at shuffles
The driver splits each job into stages. A stage is a set of tasks that can run as one pipeline, with no data moving between executors. Operations that work on each partition on its own, like a filter followed by a projection, are joined together into one stage. The pipeline breaks when data has to be redistributed across the cluster. That happens at a shuffle, for example for a groupBy aggregation or most joins. The job details page in the Spark UI draws this as a DAG, a graph that shows the order of operations and the dependencies between them. Do not confuse these stages with the "query stages" of adaptive query execution (AQE). That is a separate AQE term for the pieces of a query bounded by a shuffle or broadcast exchange, where AQE can re-optimize using runtime statistics. Shuffles have their own lesson; here the only point is that a shuffle ends one stage and begins the next.
Stages depend on each other. A stage on the reading side of a shuffle cannot start until every task in the stage that writes the shuffle data has finished, because any one of those tasks might hold rows the next stage needs. That is why a job with one shuffle usually shows two stages in the Spark UI, and a job with two shuffles shows three. On the stage list, the Shuffle Write column shows how much data a stage wrote for the stage after it, and Shuffle Read shows how much a stage read from the stage before it. You may also see greyed-out skipped stages. Spark skips a stage when its output is already available, for example because it was cached or checkpointed.
Checkpoint 4 of 7· Check yourself
In the Spark UI, a job's DAG shows Stage 4 greyed out as "skipped". What is the most likely reason?
Spark skips a stage when the data it would produce is already available, such as from a cache or checkpoint, so it does not run that work again.
“Spark is smart enough to skip some stages if they don't need to be recomputed.”Source: docs.databricks.com
Checkpoint 5 of 7· Exam question
A developer executes the following code on a DataFrame with 8 partitions: ```python result_df = sales_df.filter(sales_df.region == "EU") \ .groupBy("product_id") \ .agg(sum("amount").alias("total_amount")) result_df.show() ``` How many stages will the resulting Spark job contain?
Correct answer: B — 2 stages
- A. This is incorrect because `groupBy().agg()` requires a shuffle to co-locate rows with the same key, and a shuffle always introduces a boundary between stages, so the job cannot run as a single stage.
- B. This is correct because `filter` is a narrow transformation that runs in the first stage alongside the map-side portion of the aggregation, and the shuffle required by `groupBy().agg()` starts a second stage that finalizes the aggregation, giving exactly two stages.
- C. This is incorrect because there is only one shuffle-causing operation in this plan; `filter` does not introduce its own stage boundary since it is a narrow transformation applied within the first stage.
- D. This is incorrect because it assumes multiple shuffle boundaries, but only `groupBy().agg()` requires data movement between partitions here.
- E. This is incorrect because it overstates the number of shuffle boundaries; a single wide operation like this aggregation produces one boundary, not several.
4.Tasks: one partition, one executor, one core
The task is the bottom of the hierarchy. A task is a unit of work sent to one executor. Inside a stage, Spark creates one task for each partition of the data that stage handles. Every task in a stage runs the same code on a different slice of the data. An executor runs its tasks on threads, about one task per CPU core it has been given. So the most tasks the cluster can run at once is roughly the total number of executor cores. A stage with more tasks than that runs in waves.
| Level | What it is | What creates it |
|---|---|---|
| Application | A driver program plus its own executors on the cluster | Submitting or starting a Spark program |
| Job | A parallel computation made of multiple tasks | A Spark action (e.g. save, collect) |
| Stage | A set of tasks in a job that depends on other stages | The driver splitting a job where data must be shuffled |
| Task | A unit of work sent to one executor | One for each partition in a stage |
Because a task is the unit of parallelism, a stage's task count tells you something when you are diagnosing a slow job. Databricks' guide to slow jobs says to start with the longest stage and look at how many tasks it has. A long stage with only one task is a warning sign. While that single task runs, only one CPU is busy and the rest of the cluster may be idle. Typical causes include a window function with no PARTITION BY, an unsplittable file such as gzip, and an explicit repartition(1) or coalesce(1). Adding executors does not speed up such a stage, because there is only one partition to work on.
About 16, one per executor core. The other tasks wait and are scheduled as cores become free, so the stage runs in several waves. A stage with only 1 task would use one core and leave the other 15 idle.
Checkpoint 6 of 7· Match them up
Match each level of the hierarchy to its definition.
Tap a term, then the definition that fits it.
Each level contains the one below it: an application runs jobs, a job is split into stages, and a stage is made of tasks that run on executors.
“A unit of work that will be sent to one executor”Source: spark.apache.org
Checkpoint 7 of 7· Exam question
A developer runs the following code on an orders DataFrame: ```python repartitioned_df = orders_df.repartition("customer_id") filtered_df = repartitioned_df.filter(repartitioned_df.status == "OPEN") summary_df = filtered_df.groupBy("customer_id").agg(count("*").alias("open_orders")) summary_df.collect() ``` How many stages will the resulting Spark job contain?
Correct answer: C — 3 stages
- A. This is incorrect because the plan contains two operations, `repartition` and `groupBy().agg()`, that each require a shuffle, and each shuffle introduces its own stage boundary.
- B. This is incorrect because it accounts for only one of the two shuffles in this plan; both `repartition` and the aggregation move data between partitions.
- C. This is correct because `repartition` forces a shuffle that ends the first stage, `filter` then runs as a narrow transformation inside the resulting stage, and `groupBy().agg()` forces a second shuffle that starts a third stage, giving three stages total.
- D. This is incorrect because it overcounts the boundaries; `filter` is a narrow transformation and does not add a stage on its own between the two shuffles.
- E. This is incorrect because it assumes more shuffle-causing operations than are actually present in this plan.
Not in the way a one-task stage is. With one task per partition, the 8 tasks can all run at once on 8 of the 16 cores, so there are no waves, though 8 cores sit idle. A stage with just 1 task is the real warning sign, because it leaves nearly the whole cluster idle.
Exam traps
Each one states something that sounds right. Open it to see what is actually true.
1.Each transformation (filter, select, groupBy) starts its own Spark job.Why is that wrong?
Transformations only add to the plan. A job starts only when an action such as collect, count, show or a write needs a result.
Covered in Jobs: one per action
2.If a job shows as running in the event timeline, its tasks are already running on executors.Why is that wrong?
A job counts as running from the moment it is submitted, so jobs submitted at the same time can all show as running before any of their tasks have started.
Covered in Jobs: one per action
3.Adding executors will speed up a slow stage, however many tasks it has.Why is that wrong?
Each task handles one partition and uses one core. A stage with one task keeps one CPU busy while the rest of the cluster sits idle, so more executors do not help.
Covered in Tasks: one partition, one executor, one core
4.Executors are a shared pool that all applications on a cluster use, so applications can share in-memory data.Why is that wrong?
Each application gets its own executors, which stay up for the whole application. Sharing data between applications requires external storage.
Covered in The top level: one application, one driver, its own executors
Sources
Every claim above is drawn from one of these pages, quoted as it was written on the date shown.
- 1.
“Ensure that the tasks are executed on multiple executors (nodes) in your compute to have enough parallelism while processing.”
↩︎ The top level: one application, one driver, its own executors“The job details page shows a DAG visualization.”
↩︎ Stages: where a job is cut at shuffles“Spark is smart enough to skip some stages if they don't need to be recomputed.”
↩︎ Checkpoint - 2.https://spark.apache.org/docs/latest/cluster-overview.htmlSecondary source
“Spark applications run as independent sets of processes on a cluster, coordinated by the SparkContext object in your main program (called the driver program).”
↩︎ The top level: one application, one driver, its own executors“A parallel computation consisting of multiple tasks that gets spawned in response to a Spark action (e.g. save, collect)”
↩︎ Jobs: one per action“Each job gets divided into smaller sets of tasks called stages that depend on each other”
↩︎ Stages: where a job is cut at shuffles“Each job gets divided into smaller sets of tasks called stages that depend on each other”
↩︎ Key concept“Each application gets its own executor processes, which stay up for the duration of the whole application and run tasks in multiple threads.”
↩︎ Exam trap 4“Once connected, Spark acquires executors on nodes in the cluster, which are processes that run computations and store data for your application.”
↩︎ Checkpoint“Each application has its own executors.”
↩︎ Checkpoint“A unit of work that will be sent to one executor”
↩︎ Checkpoint - 3.https://docs.databricks.com/aws/en/spark/faqOfficial docs
“Actions: trigger processing logic to evaluate and output a result.”
↩︎ Jobs: one per action“none of the logic defined by a collection of operations are evaluated until an action is triggered”
↩︎ Exam trap 1“none of the logic defined by a collection of operations are evaluated until an action is triggered”
↩︎ Prediction - 4.
“A job appears as running from the moment it's submitted, not from when its first task starts executing.”
↩︎ Jobs: one per action“A job appears as running from the moment it's submitted, not from when its first task starts executing.”
↩︎ Exam trap 2 - 5.https://spark.apache.org/docs/latest/sql-performance-tuning.htmlSecondary source
“Spark will list the files by using Spark distributed job.”
↩︎ Jobs: one per action - 6.https://docs.databricks.com/aws/en/optimizations/aqeOfficial docs
“at the end of a shuffle and broadcast exchange (referred to as a query stage in AQE)”
↩︎ Stages: where a job is cut at shuffles - 7.
“Shuffle Write: How much shuffle data this stage wrote.”
↩︎ Stages: where a job is cut at shuffles“The number of tasks in the long stage can point you in the direction of your issue.”
↩︎ Tasks: one partition, one executor, one core - 8.https://www.databricks.com/blog/2020/05/29/adaptive-query-execution-speeding-up-spark-sql-at-runtime.htmlSecondary source
“the following stage can only proceed if all the parallel processes running the materialization have completed”
↩︎ Stages: where a job is cut at shuffles - 9.
“While this one task is running only one CPU is utilized and the rest of the cluster may be idle.”
↩︎ Tasks: one partition, one executor, one core“While this one task is running only one CPU is utilized and the rest of the cluster may be idle.”
↩︎ Exam trap 3 - 10.https://spark.apache.org/docs/latest/rdd-programming-guide.htmlSecondary source
“Spark will run one task for each partition of the cluster.”
↩︎ Tasks: one partition, one executor, one core