Jobs, stages and tasks
How one action becomes a DAG of stages split at shuffles, and how to read it in the Spark UI.
On this page
You will learn
- How one action becomes a job, stages and tasks
- Why stage boundaries are always shuffles
- How to read the Spark UI to find a slow stage
- What happens when a task fails
Read first
- Transformations vs actions · 3 min read
- Partitions · 4 min read
Comfortable with these? Read on.
Job, stage, task
-
Action
write,count, ... -
Job
One per action (sometimes more)
-
Stages
Split at shuffles
-
Tasks
One per partition, per stage
Within a stage, Spark pipelines every narrow operation: a task reads its partition, filters, projects and computes in one pass, without materialising anything in between. A stage ends where rows have to be regrouped across the cluster, which is a shuffle.
Walk through a query
Query · the code
The command
orders.filter("status = 'PAID'") \ .join(customers, "customer_id") \ .groupBy("country").agg(F.sum("amount")) \ .write.parquet("out")
Stage 1 · scan orders
customer_id. One task per input partition, say 400 tasks.Stage 2 · scan customers
customer_id. Runs in parallel with stage 1, since it does not depend on it.Stage 3 · join + partial agg
country. 200 tasks before AQE coalescing.Stage 4 · final agg + write
If customers were small enough to broadcast, stages 2 and 3 would merge into stage 1: no shuffle for the join, so no stage boundary.
Reading the Spark UI
| Tab | What to look for |
|---|---|
| Jobs | Which action is slow. Unexpected extra jobs come from schema inference, count() calls or sort sampling. |
| Stages | The slow stage, its shuffle read and write sizes, and spill to memory or disk. |
| Stage detail | The task-time distribution: if the max task takes 10 minutes and the median 5 seconds, you have skew. |
| SQL / DataFrame | The physical plan with row counts per operator, after AQE changes. |
| Executors | GC time, failed tasks, memory use per executor. |
Failures and retries
- A failed task is retried up to
spark.task.maxFailurestimes (4 by default) before the job fails. - If shuffle files from a lost executor are needed, Spark re-runs the parts of the previous stage that produced them: lineage makes this possible without checkpoints.
- With speculation on (
spark.speculation), Spark launches a backup copy of unusually slow tasks and uses whichever finishes first. It helps with slow machines, not with skew.
Exchange nodes in explain()) is the fastest way to estimate how expensive a query is.Common mistakes
Looking at average task time only
Assuming one action is one job
Blaming the last stage
Key takeaways
- An action creates a job; shuffles split it into stages; partitions become tasks.
- Narrow operations are pipelined inside a stage.
- The Spark UI task-time distribution reveals skew.
- Failed tasks are retried, and lost shuffle output is recomputed from lineage.
Check yourself
3 questions1. What always marks a boundary between two stages?
Show the answer
A shuffle. Stages are split wherever rows must be redistributed across the cluster.
2. A stage reads 300 partitions. How many tasks does it run?
Show the answer
300. One task per partition in a stage.
3. In the stage view, max task time is 12 minutes and the median 4 seconds. What is the likely problem?
Show the answer
Data skew. A few tasks processing far more data than the rest is the signature of skew.
Go deeper
Narrow vs wide transformationsSpark internals
Data skew and saltingSpark internals
Catalyst and physical plans
Primary sources: Web UI guide · Cluster mode overview