3 · Under the hood · lesson 12 of 20

Jobs, Stages & Tasks

How one action becomes a DAG that becomes work on executors.

When you fire an action, Spark builds a DAG of transformations, cuts it at every shuffle boundary into stages, and then launches one task per partition in each stage.

  • ▸Job — one per action (show, count, write, collect).
  • ▸Stage — a run of transformations with no shuffle between them.
  • ▸Task — one stage running on one partition on one executor.
NOTE
The Spark UI (http://localhost:4040 by default) shows the full DAG for every job. It's the single most useful debugging tool you have.
Loading 3D scene…
Key takeaways
  • ✓1 action → 1 job → N stages → many tasks.
  • ✓Stage boundaries are shuffle boundaries.
  • ✓The Spark UI is the ground truth for what actually ran.