Menu

Apache Spark course · Lesson 3 of 5

Spark Execution Model: Lazy Evaluation, Jobs, Stages and Tasks

How Spark turns lazy DataFrame code into a DAG of jobs, stages and tasks, why shuffles split stages, and what to look for in each tab of the Spark UI.

  • Intermediate
  • 15 min read
  • Updated Oct 2026
On this page
  1. Sample data
  2. Lazy evaluation and the DAG
  3. What it is
  4. Why Spark is lazy
  5. Pitfalls
  6. In interviews
  7. Narrow and wide transformations
  8. What it is
  9. Not every join or aggregation shuffles
  10. Pitfalls
  11. In interviews
  12. Job, stage and task hierarchy
  13. What it is
  14. Counting them, without AQE
  15. Counting them, with AQE (the default)
  16. Pitfalls
  17. In interviews
  18. The Spark UI: finding slow stages and tasks
  19. What it is
  20. What each tab tells you
  21. A reading routine for a slow job
  22. Pitfalls
  23. In interviews
  24. Practice questions
  25. Key takeaways

Spark does not run your code line by line. It records what you asked for, waits for an action, then plans the whole computation at once and cuts it into jobs, stages and tasks. Once you can predict that breakdown for a query, the Spark UI stops being a wall of numbers and becomes the fastest way to find out why a job is slow.

Sample data

Two hundred thousand synthetic orders across 1,000 customers. spark.sql.shuffle.partitions is set to 4 (instead of the default 200) so every count below is small enough to read. The helper api() reads the Spark UI’s REST API, which is how the job and stage numbers in this lesson were captured.

import json, urllib.request
from pyspark.sql import SparkSession, functions as F

spark = (SparkSession.builder.master("local[2]").appName("execution-model")
         .config("spark.sql.shuffle.partitions", "4")
         .getOrCreate())
sc = spark.sparkContext

_opener = urllib.request.build_opener(urllib.request.ProxyHandler({}))
def api(path):
    port = sc.uiWebUrl.rsplit(":", 1)[1]
    url = f"http://localhost:{port}/api/v1/applications/{sc.applicationId}/{path}"
    return json.load(_opener.open(url))

orders = (spark.range(0, 200_000, numPartitions=4)
          .withColumn("customer", F.col("id") % 1000)
          .withColumn("amount", (F.col("id") % 97).cast("double")))

Lazy evaluation and the DAG

What it is

Lazy evaluation means transformations (filter, select, groupBy, join) only describe a result. Spark records them and runs nothing until an action (count, collect, show, write) needs a value. The recorded steps form a DAG, a directed acyclic graph: each node is a dataset, each edge a transformation, and there are no cycles because every step produces a new immutable dataset.

jobs_before = len(api("jobs"))

per_customer = (orders.filter("amount > 10")
                .groupBy("customer")
                .agg(F.sum("amount").alias("total")))
ranked = per_customer.orderBy(F.desc("total"))

print("jobs started by transformations:", len(api("jobs")) - jobs_before)
jobs started by transformations: 0

Three transformations, zero jobs. Nothing has read a single row yet.

Why Spark is lazy

  • Whole-plan optimisation. Seeing the full chain lets Catalyst push filters down, drop unused columns and choose join strategies before any data moves.
  • Pipelining. Consecutive narrow steps run together inside one task, row by row, with no intermediate dataset written anywhere.
  • Skipping unused work. A branch whose result is never consumed by an action is never computed.
  • Fault tolerance. The DAG doubles as the lineage used to recompute lost partitions.

explain() shows the plan without running it:

per_customer.explain()
== Physical Plan ==
AdaptiveSparkPlan isFinalPlan=false
+- HashAggregate(keys=[customer#1L], functions=[sum(amount#2)])
   +- Exchange hashpartitioning(customer#1L, 4), ENSURE_REQUIREMENTS, [plan_id=19]
      +- HashAggregate(keys=[customer#1L], functions=[partial_sum(amount#2)])
         +- Project [(id#0L % 1000) AS customer#1L, cast((id#0L % 97) as double) AS amount#2]
            +- Filter (cast((id#0L % 97) as double) > 10.0)
               +- Range (0, 200000, step=1, splits=4)

Read it from the bottom up: generate rows, filter, compute columns, partially sum per customer, shuffle (Exchange), then finish the sum. isFinalPlan=false means Adaptive Query Execution may still change it at run time.

Pitfalls

  • Errors surface at the action. A bad cast written on line 10 fails on line 40 where you call write. Read the stack trace for the operator, not just the line number.
  • Some “transformations” run jobs. Reading CSV or JSON with schema inference, sortBy on RDDs (to sample range boundaries) and df.rdd.getNumPartitions() on some plans start jobs before your action.
  • Repeated actions repeat work. Each action replays the DAG from the source unless you cache an intermediate result.

In interviews

“What is lazy evaluation and why does Spark use it?” Give the definition, then two concrete benefits (optimisation across the whole plan, pipelining of narrow steps), and one cost (errors appear late; repeated actions recompute).

Narrow and wide transformations

What it is

The DAG has two kinds of edges.

  • A narrow transformation computes each output partition from one input partition (or a fixed small set). No data crosses executors. Examples: filter, select, withColumn, map, flatMap, mapPartitions, union, coalesce (when reducing partitions).
  • A wide transformation needs data from many input partitions for each output partition, so rows must be redistributed by key across the cluster: a shuffle. Examples: groupBy/agg, distinct, join (unless one side is broadcast), orderBy, repartition, window functions with partitionBy, RDD reduceByKey and groupByKey.

A chain of narrow steps has no Exchange and runs as one pipelined stage. The *(1) prefix marks operators fused into the same generated code:

narrow = orders.filter("amount > 10").select("customer", (F.col("amount") * 1.2).alias("gross"))
narrow.explain()
== Physical Plan ==
*(1) Project [(id#0L % 1000) AS customer#1L, (cast((id#0L % 97) as double) * 1.2) AS gross#10]
+- *(1) Filter (cast((id#0L % 97) as double) > 10.0)
   +- *(1) Range (0, 200000, step=1, splits=4)

Not every join or aggregation shuffles

  • A broadcast hash join copies the small table to every executor, so the large side does not move; that join is narrow for the large side.
  • If the data is already partitioned on the grouping or join key (for example after an earlier repartition on that key, or from bucketed tables), Spark reuses the existing distribution and skips the Exchange.

Always confirm with explain() rather than assuming.

Pitfalls

  • Treating repartition as free. It is a full shuffle, even when it reduces the partition count; coalesce merges partitions without one.
  • Chaining two wide operations on different keys (group by customer, then by country) costs two shuffles. Sometimes one well-chosen key does both jobs.
  • Calling orderBy before writing data that nobody reads in order: a global sort is a wide transformation with an extra sampling job.

In interviews

Expect “What causes a shuffle?” (see the shuffle question). List the wide operations, explain that a shuffle is needed whenever rows with the same key must meet, and mention the exceptions: broadcast joins and data that is already partitioned correctly.

Job, stage and task hierarchy

What it is

Level Created by Boundary Runs as
Job An action (one action can create several jobs) One action, or one AQE query stage One or more stages
Stage The DAG scheduler Shuffle boundaries: each wide dependency starts a new stage A set of identical tasks
Task The task scheduler One per partition of the stage One thread on one executor core

A stage that writes shuffle output for a later stage is a shuffle map stage; the last stage of a job, which produces the action’s result, is the result stage. Inside a stage, all narrow steps are pipelined: a task reads its partition, filters, projects and partially aggregates in one pass.

Counting them, without AQE

With Adaptive Query Execution turned off, count() on the aggregated DataFrame gives one job:

spark.conf.set("spark.sql.adaptive.enabled", "false")
jobs_before = len(api("jobs"))
sc.setJobDescription("count without AQE")
print(per_customer.count())
sc.setJobDescription(None)

for j in sorted(api("jobs"), key=lambda j: j["jobId"])[jobs_before:]:
    print(f'job {j["jobId"]}: stages {j["stageIds"]}, tasks {j["numTasks"]}')
spark.conf.set("spark.sql.adaptive.enabled", "true")
1000
job 0: stages [0, 1, 2], tasks 9

Three stages, nine tasks:

  1. Stage 0 (4 tasks, one per input partition): generate, filter, project, partial sum, then write shuffle files hashed by customer.
  2. Stage 1 (4 tasks, one per shuffle partition): read the shuffle, finish the sum per customer, partially count rows, write a second small shuffle. count() itself adds an aggregation.
  3. Stage 2 (1 task): read the partial counts and add them up.

Counting them, with AQE (the default)

Adaptive Query Execution runs each shuffle stage as a separate job so it can look at real statistics before planning the next step:

jobs_before = len(api("jobs"))
sc.setJobDescription("count with AQE")
print(per_customer.count())
sc.setJobDescription(None)

for j in sorted(api("jobs"), key=lambda j: j["jobId"])[jobs_before:]:
    print(f'job {j["jobId"]}: stages {j["stageIds"]}, tasks {j["numTasks"]}, skipped stages {j["numSkippedStages"]}')
1000
job 1: stages [3], tasks 4, skipped stages 0
job 2: stages [4, 5], tasks 5, skipped stages 1
job 3: stages [6, 7, 8], tasks 6, skipped stages 2

One action, three jobs. Later jobs list earlier stages again, but those are skipped: their shuffle output already exists, so they are not rerun. That is why the UI shows “skipped” stages in grey. The task totals count the skipped stages’ tasks too: job 2 is the skipped 4-task map stage plus one new task, not four, because AQE saw that the four shuffle partitions were tiny and coalesced them into one. The AQE lesson covers this in detail.

You can see the final plan after execution: ShuffleQueryStage marks each stage boundary AQE materialised, and AQEShuffleRead coalesced shows the merge.

w = orders.groupBy("customer").count()
w.collect()
w.explain()
== Physical Plan ==
AdaptiveSparkPlan isFinalPlan=true
+- == Final Plan ==
   ResultQueryStage 1
   +- *(2) HashAggregate(keys=[customer#1L], functions=[count(1)])
      +- AQEShuffleRead coalesced
         +- ShuffleQueryStage 0
            +- Exchange hashpartitioning(customer#1L, 4), ENSURE_REQUIREMENTS, [plan_id=209]
               +- *(1) HashAggregate(keys=[customer#1L], functions=[partial_count(1)])
                  +- *(1) Project [(id#0L % 1000) AS customer#1L]
                     +- *(1) Range (0, 200000, step=1, splits=4)
+- == Initial Plan ==
   (trimmed: the same plan without query stages)

Pitfalls

  • Expecting one job per action. AQE, schema inference, broadcast joins (the broadcast side runs its own job) and orderBy sampling all add jobs. That is normal.
  • Confusing parallelism with task count. A stage with 2,000 tasks on 100 cores runs in about 20 waves; a stage with 50 tasks on 100 cores leaves half the cores idle.
  • Ignoring the last stage. A final stage with one task (a global sort to a single file, coalesce(1)) runs on one core no matter how big the cluster is.

In interviews

“Explain jobs, stages and tasks” is a staple (see the interview answer). Give the hierarchy, say that stages are cut at shuffle boundaries and tasks map one-to-one to partitions, then add the modern nuance: with AQE one action often produces several jobs and skipped stages. Walking through the stage count of a concrete query, as above, is the strongest answer.

The Spark UI: finding slow stages and tasks

What it is

Every running application serves a web UI from the driver, on port 4040 by default (4041, 4042 and so on if the port is taken). Managed platforms link to it from the job page. After the application ends, the same pages are available from the History Server if event logging was enabled. Work through it top-down: query, then job, then stage, then task.

What each tab tells you

Tab What it shows What to look for
Jobs Every job with its status, duration and stage progress; an event timeline of jobs and executors being added or removed The slowest job; failed jobs; gaps in the timeline where nothing ran (driver-side work such as listing files or collect processing); executors being lost
Stages Every stage with duration, task progress, input, output, shuffle read and shuffle write Stages with large shuffle read or write; stages that ran many times (retries after fetch failures)
Stage detail Summary metrics (min, 25th percentile, median, 75th percentile, max) per task for duration, GC time, input, shuffle read and spill; the task table; an event timeline per task Max far above median means skew. Non-zero spill (memory/disk) means partitions too large for execution memory. High GC time relative to duration means memory pressure. Long scheduler delay or task deserialization time means too many tiny tasks or huge closures
SQL / DataFrame Each query with its job ids and a graph of physical operators annotated with metrics (rows output, data size, time, spill); the initial and final AQE plans Where rows explode (a join whose output is far larger than its inputs); which join strategy actually ran; whether filters pruned files; AQE coalescing and skew splits
Storage Cached RDDs and DataFrames: storage level, fraction cached, size in memory and on disk Datasets cached but never reused; partially cached data being recomputed; cache that should have been released
Executors Every executor with cores, storage memory used, task counts, shuffle totals, GC time, failed tasks; links to logs and thread dumps One executor with far more failed tasks (bad node); high GC time across executors; uneven task counts; driver memory use
Environment Effective configuration, JVM and system properties, classpath Confirm the settings you think you set (shuffle partitions, memory, AQE) actually reached the application

A reading routine for a slow job

  1. SQL tab: open the slow query; find the operator with the largest time or row count, and note its stage.
  2. Stages tab: open that stage. Compare median and max task duration and shuffle read size.
  3. Decide which pattern you are looking at:
Symptom Likely cause Lesson
A few tasks far slower than the rest Data skew on a hot key Partitions, shuffles and skew
Large shuffle read and write Wide transformation over too much data; filter or aggregate earlier, broadcast the small side Partitions, shuffles and skew
Spill to memory and disk Partitions too large for execution memory Memory and executor tuning
Thousands of tasks lasting milliseconds Too many small partitions or small files Partitions, shuffles and skew
Huge row counts after a join Fan-out on a non-unique key Joins and join strategy
High GC time Memory pressure, large cached objects, wide rows Memory and executor tuning

Pitfalls

  • Reading averages. One straggler hides in an average; always compare the max with the median.
  • Trusting the pre-execution plan. explain() before running shows the initial AQE plan; the SQL tab shows what actually ran.
  • Losing the UI. The live UI disappears when the application ends. Enable event logging so the History Server can show it later.

In interviews

“How would you debug a slow Spark job?” A strong answer is a routine, not a guess: open the SQL tab to find the expensive operator, go to its stage, compare task metric percentiles, then match the symptom (skew, spill, shuffle volume, tiny tasks, GC) to a fix, and confirm the improvement in the UI afterwards.

Practice questions

How many stages does df.filter(…).groupBy(“k”).count().collect() produce, and why?

Without AQE, two stages: one that reads, filters and partially aggregates, ending at the shuffle by k; and one that reads the shuffle and finishes the aggregation, producing the result. The filter does not add a stage because it is narrow and pipelined into the first one. With AQE on, you may see two jobs (one per query stage) and a skipped stage, but the same two pieces of work.

Why can one action create several jobs?

Adaptive Query Execution submits each shuffle stage as its own job so it can re-plan using real statistics. Other causes: a broadcast join collects the small side in a separate job, a global sort runs a sampling job to pick range boundaries, and file sources may run a job to list or infer schema. Stages already computed by an earlier job are shown as skipped.

Classify these as narrow or wide: filter, join with a broadcast table, distinct, coalesce(10), repartition(10), withColumn, window with partitionBy.

Narrow: filter, withColumn, coalesce(10) (it merges partitions without a shuffle), and the broadcast join (the large side does not move). Wide: distinct, repartition(10) (a full shuffle even when reducing) and a window with partitionBy (rows are shuffled by the partition key).

A stage has 200 tasks: the median takes 4 seconds and the max takes 9 minutes. What do you check next?

This is the signature of skew. In the stage detail, compare max and median shuffle read size and records: if the slow task read far more data, one key dominates. Find the key with a groupBy(key).count().orderBy(desc("count")) on the input, then choose a fix: let AQE split the skewed join partition, broadcast the smaller side, filter or separately handle NULL or default keys, or salt the key.

Your notebook code fails with a cast error at df.write(…), but the cast is 30 lines earlier. Why?

Lazy evaluation. The cast was only recorded in the plan; it executed when the write action ran the job. The stack trace shows the action, so look at the operator in the error message (or run explain()) to locate the transformation.

Which Spark UI tab confirms that a configuration change took effect, and which shows whether a cached DataFrame is actually in memory?

The Environment tab lists the effective Spark properties for the application. The Storage tab lists cached datasets with their storage level, the fraction of partitions cached, and the size in memory and on disk.

Key takeaways

  • Transformations are lazy and build a DAG; actions trigger optimisation and execution.
  • Narrow transformations pipeline inside one stage; wide transformations need a shuffle, which starts a new stage.
  • An action creates jobs, a shuffle splits a job into stages, and each partition becomes one task on one core.
  • With AQE (the default), one action often runs as several jobs with skipped stages; that is expected.
  • In the Spark UI, go SQL tab, then stage, then task metrics, and compare max with median to spot skew, spill and GC problems.
  • Label jobs with setJobDescription and enable event logs so you can still read the UI after the application ends.

By Data Career Hub Editorial · Last reviewed Oct 2026 · Examples run on PySpark 4.2.0 in local mode (local[2]) with spark.sql.shuffle.partitions set to 4 so the counts stay readable. Job and stage counts were read from the Spark UI REST API during the run. The Spark UI tour is described in words because the UI is not captured here.

Progress is saved in this browser only. No account needed.

Search
Filter by type