All concepts

Spark's Execution Model

Nothing runs until you ask for a result — then the plan is cut into stages at every point data has to cross the network.

Distributed Processing · Intermediate · ~6 min

In plain English

Writing a shopping list is free; the walk to the shops is not. Spark plans the whole list, then works out how many trips it needs.

Why it's worth your time

Every Spark performance problem is either too many trips (stages) or one trolley doing all the carrying (skew).

If you remember three things

  • Nothing executes until an action; the plan is optimised first
  • Stage boundaries are shuffles — count the Exchanges to see the real cost
  • Narrow transformations fuse into one pass; wide ones cross the network

Overview

Spark builds a lazy plan. Transformations like filter, select and join add nodes to a logical plan and execute nothing; only an action — count, write, collect — triggers work. The optimiser then rewrites that plan, and the scheduler cuts it into stages at every shuffle boundary, because a shuffle is the one operation where data must move between machines. Within a stage, work is a set of independent tasks — one per partition — that run wherever there's a free core. Understanding this is the difference between tuning Spark and guessing: almost every performance problem is either 'too many stages' or 'one task doing all the work'.

In an interview

Spark transformations are lazy and build a plan; an action executes it. The optimiser rewrites the plan, then the scheduler splits it into stages at shuffle boundaries, since a shuffle is where data crosses the network. Each stage is a set of tasks, one per partition, run in parallel across executors. Narrow transformations pipeline inside a stage; wide ones force a new stage.

Production defaults

First move
read the physical plan — Exchanges and PushedFilters — before tuning anything
Adaptive execution
leave it on; it sizes post-shuffle partitions better than a static number
Caching
only for data reused across actions, and unpersist afterwards

What breaks

  • Driver OOM — collect() on a distributed dataset. Write to storage and read what you need.
  • Executors OOM after adding cache() — Cached data took memory the shuffle needed. Cache less, or use MEMORY_AND_DISK.

Watch it explained

1.3 Apache Spark Architecture | Spark Execution Model | Spark tutorial — Data Savvy, 7:47

Related