Narrow vs wide transformations
Which operations move data across the network, and why the shuffle is the most expensive step.
On this page
You will learn
- Which transformations are narrow and which are wide
- What physically happens during a shuffle
- Why shuffles dominate the cost of most jobs
- How to remove or shrink shuffles
filter, select). A wide one needs rows from many partitions (groupBy, join, orderBy), so Spark shuffles: writes data to disk, sends it over the network and reads it back.Narrow and wide
| Narrow (no shuffle) | Wide (shuffle) |
|---|---|
select, withColumn, filter | groupBy().agg(), distinct, dropDuplicates |
explode, union, coalesce(n) | join (unless broadcast), orderBy |
sample, limit per partition | repartition, window functions with partitionBy |
The test: can a task produce its output by looking only at its own input partition? If yes, it is narrow. Grouping by customer needs every row of a customer, wherever it is, so it is wide.
What a shuffle does
-
Map side
Each task hashes rows by key into N buckets
-
Shuffle write
Buckets are written to local disk
-
Shuffle read
Each reducer fetches its bucket from every mapper over the network
-
Reduce side
Rows of a key are together; work continues
With M map tasks and N reduce partitions, that is up to M × N blocks moved. Everything shuffled is serialized, usually compressed, written to disk and read back. On a typical job, shuffles account for most of the run time, which is why the Spark UI shows shuffle read and write sizes for every stage.
Seeing shuffles in a plan
In explain(), every shuffle is an Exchange node. Exchange hashpartitioning(customer_id, 200) means "shuffle by customer_id into 200 partitions". BroadcastExchange is different: it ships a small table to every executor instead.
A groupBy in the physical plan
HashAggregate(keys=[country], functions=[sum(amount)])
+- Exchange hashpartitioning(country, 200)
+- HashAggregate(keys=[country], functions=[partial_sum(amount)])
+- FileScan parquet [country,amount]How to shuffle less
- Filter and select early. Fewer rows and columns before the exchange means fewer bytes shuffled. Catalyst does much of this, but not through UDFs or across caching boundaries.
- Broadcast small tables so the large side of a join is not shuffled.
- Aggregate before joining when the join is only used to attach attributes to totals.
- Reuse partitioning. If two operations use the same key, Spark can reuse one shuffle; a different key or partition count forces another.
- Bucket tables that are repeatedly joined on the same key: data written pre-shuffled into buckets can be joined without a shuffle, if both sides have matching bucketing.
Common mistakes
Calling repartition "to be safe"
Joining before filtering
Window functions with many different partitionBy keys
Key takeaways
- Narrow operations work within a partition; wide ones need a shuffle.
- A shuffle writes, transfers and reads every row involved.
- Shuffles appear as Exchange nodes in explain().
- Filter early, broadcast small tables and reuse partitioning to shuffle less.
Check yourself
3 questions1. Which is a narrow transformation?
Show the answer
withColumn. withColumn computes each row from the row itself.
2. What appears in a physical plan for a shuffle?
Show the answer
Exchange. Exchange nodes mark data redistribution.
3. Why does a broadcast join avoid shuffling the large table?
Show the answer
Every executor gets a full copy of the small table, so large-table rows can be joined where they are. With the small side everywhere, no regrouping of the large side is needed.
Go deeper
Jobs, stages and tasksSpark internals
Broadcast joinsSpark internals
Shuffle partitions
Primary sources: RDD guide: shuffle operations · Performance tuning