The DAG of Transformations

Building the DAG of Transformations

When an action runs, the physical plan becomes a directed acyclic graph (DAG) of RDDs, and the kind of dependency between neighbors decides the cost. A narrow step (where, select, the probe side of a broadcast join) needs one input partition per output partition and moves no data. A wide step (groupBy, a sort-merge join, orderBy) needs every input partition and must redistribute rows by key.

Narrow dependencies stay inside a partition; wide dependencies need a shuffle
Narrow dependencies stay inside a partition; wide dependencies need a shuffle

The DAGScheduler cuts the graph at every wide dependency and pipelines each chain of narrow steps, so a row flows from the Parquet 129 reader through the filter and join into a partial aggregate without being stored.