Skip to content
Everyone knows 3 min read

Jobs, stages and tasks

How one action becomes a DAG of stages split at shuffles, and how to read it in the Spark UI.

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

Comfortable with these? Read on.

TL;DR Each action starts a job. Spark cuts the job into stages wherever data must be shuffled, and runs each stage as one task per partition. Stages run in order; tasks inside a stage run in parallel.

Job, stage, task

  1. Action

    write, count, ...

  2. Job

    One per action (sometimes more)

  3. Stages

    Split at shuffles

  4. 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

Revenue per country for paid orders, joined to a large customers table:

The command

orders.filter("status = 'PAID'") \
    .join(customers, "customer_id") \
    .groupBy("country").agg(F.sum("amount")) \
    .write.parquet("out")

Stage 1 · scan orders

Read orders, filter paid rows, and write shuffle files partitioned by customer_id. One task per input partition, say 400 tasks.

Stage 2 · scan customers

Read customers and write shuffle files by customer_id. Runs in parallel with stage 1, since it does not depend on it.

Stage 3 · join + partial agg

Read both shuffles, sort-merge join them by customer, and pre-aggregate by country. Then shuffle again, by country. 200 tasks before AQE coalescing.

Stage 4 · final agg + write

Read the country shuffle, finish the sums, and write the output files. Only a handful of tasks are needed: there are few countries, and AQE coalesces the partitions.

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

TabWhat to look for
JobsWhich action is slow. Unexpected extra jobs come from schema inference, count() calls or sort sampling.
StagesThe slow stage, its shuffle read and write sizes, and spill to memory or disk.
Stage detailThe task-time distribution: if the max task takes 10 minutes and the median 5 seconds, you have skew.
SQL / DataFrameThe physical plan with row counts per operator, after AQE changes.
ExecutorsGC time, failed tasks, memory use per executor.

Failures and retries

  • A failed task is retried up to spark.task.maxFailures times (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.
Why this matters: stage boundaries are where data is written to local disk and read back over the network. Counting shuffles in a plan (the Exchange nodes in explain()) is the fastest way to estimate how expensive a query is.

Common mistakes

Looking at average task time only

Skew hides in the maximum. Check the distribution.

Assuming one action is one job

Schema inference, sort sampling and broadcasts can add jobs.

Blaming the last stage

The slow part is often an earlier stage; the UI shows each stage's time.

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 questions

1. 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

Primary sources: Web UI guide · Cluster mode overview