Spark DAG Scheduler: From Logical Plan to Physical Execution

When you submit a Spark job, the DAG scheduler translates the logical plan (chain of RDD transformations) into a physical execution plan:

  1. Build the DAG from the sequence of transformations
  2. Identify stage boundaries at shuffle points (wide dependencies like groupByKey, reduceByKey)
  3. Topologically sort stages to determine execution order
  4. Submit tasks for each stage’s partitions to the task scheduler

Stages within the same pipeline (connected by narrow dependencies) can fuse their transformations, avoiding intermediate materialization.

A DAG of ten RDDs and an action, with two joins: the scheduler walks back from the action, cuts the graph into stages at its four wide (shuffle) dependencies, sorts the stages topologically and runs them in waves, and clicking an arrow switches it between narrow and wide.

Step through the plan, or click an arrow to turn a narrow dependency into a shuffle (or a join side into a co-partitioned one) and watch the stages re-split and the waves change. Clicking an RDD names its stage and the dependency that put the boundary where it is.