Narrow transformation
Each input partition contributes to at most one output partition. No shuffle of data across the network. Examples: map, filter, flatMap, withColumn, select, union (often).
Wide transformation
Input partitions contribute to many output partitions. Requires a shuffle (exchange). Examples: groupByKey/groupBy.agg, reduceByKey, join (non-broadcast), distinct, repartition, orderBy/sort.
NARROW (pipelined in one stage) WIDE (shuffle boundary) +-----+ +-----+ +-----+ shuffle +-----+ | P0 |--->| P0' | | P0 | ---------------> | Q0 | | P1 |--->| P1' | | P1 | ---------------> | Q1 | | P2 |--->| P2' | | P2 | ---------------> | Q2 | +-----+ +-----+ +-----+ +-----+
Code
# Narrow: stays in same stage when possible
df2 = df.filter("amount > 0").withColumn("tax", df.amount * 0.1)
# Wide: introduces a shuffle
by_user = df2.groupBy("user_id").sum("amount")Why that matters
- Narrow ops are cheaper and can be pipelined.
- Wide ops create new stages, more disk/network I/O, and are common slow-job culprits.
- Prefer reducing data before wide ops (filter early).