Apache Spark courseLesson 2 of 4
Apache Spark course · Lesson 2 of 4
Spark Execution Model: Jobs, Stages and Tasks
How Spark turns DataFrame code into jobs, stages and tasks, how the driver and executors divide the work, and how to read the Spark UI to find slow stages.
On this page
Understanding how Spark executes your code is the basis of all performance work. The vocabulary is small: driver, executors, jobs, stages, tasks.
Driver and executors
- The driver runs your program, builds the plan, and schedules work. It also receives anything you
collect(). - Executors are processes on worker nodes that run tasks and hold cached data. Each executor has a number of cores; one core runs one task at a time.
Total parallelism is roughly executors × cores per executor.
From code to execution
- You chain transformations; Spark builds a logical plan.
- An action (
count,write,show, …) triggers optimisation into a physical plan. - Execution is submitted as one or more jobs.
- Each job is split into stages at shuffle boundaries (wide transformations such as
groupByorjoin). - Each stage runs as tasks, one per partition, in parallel across executor cores.
from pyspark.sql import SparkSession, functions as F
spark = SparkSession.builder.master("local[2]").appName("example").getOrCreate()
orders = spark.range(0, 200_000).withColumn("customer", F.col("id") % 1000)
per_customer = orders.groupBy("customer").agg(F.count("*").alias("n"))
per_customer.explain()
The physical plan has one Exchange (the shuffle). Everything before it is the first stage (read, project, partial aggregate); everything after it is the next stage (final aggregate).
print(per_customer.count())
1000
Jobs per action
An action triggers at least one job. With Adaptive Query Execution, Spark often runs each shuffle stage as its own job so it can look at real statistics before planning the next step, so you may see several jobs for one action in the UI. That is normal.
Reading the Spark UI
Open the UI (port 4040 locally, or your platform’s link) and work top-down:
- Jobs tab: which action is slow.
- Stages tab: which stage inside it dominates. Compare median and max task duration. A max far above the median means skew.
- Stage detail: shuffle read/write sizes, spill to disk, GC time, and the slowest tasks.
- SQL tab: the executed plan with row counts per operator, the best place to see where rows explode or where a join strategy changed.
What usually makes a stage slow
| Symptom | Likely cause |
|---|---|
| A few tasks far slower than others | Data skew |
| Large shuffle read/write | Wide transformation on lots of data; filter or aggregate earlier |
| Spill to disk | Partitions too large for executor memory |
| Many tiny tasks | Too many small partitions or small files |
| Huge row counts after a join | Fan-out on a non-unique key |
Common mistakes
- Optimising code without opening the Spark UI.
- Looking at average task time and missing a single straggler task.
- Doing heavy work on the driver (
collectthen Python loops).
Interview relevance
“Explain jobs, stages and tasks” is a staple. See the interview answer.
Key takeaway
Actions create jobs, shuffles split jobs into stages, and partitions become tasks. Use the Spark UI to find the stage and the tasks where time goes.
Progress is saved in this browser only. No account needed.