Skip to content
Good engineers know 3 min read

Narrow vs wide transformations

Which operations move data across the network, and why the shuffle is the most expensive step.

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

Read first

Comfortable with these? Read on.

TL;DR A narrow transformation computes each output partition from one input partition (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, filtergroupBy().agg(), distinct, dropDuplicates
explode, union, coalesce(n)join (unless broadcast), orderBy
sample, limit per partitionrepartition, 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

  1. Map side

    Each task hashes rows by key into N buckets

  2. Shuffle write

    Buckets are written to local disk

  3. Shuffle read

    Each reducer fetches its bucket from every mapper over the network

  4. 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

  1. 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.
  2. Broadcast small tables so the large side of a join is not shuffled.
  3. Aggregate before joining when the join is only used to attach attributes to totals.
  4. Reuse partitioning. If two operations use the same key, Spark can reuse one shuffle; a different key or partition count forces another.
  5. 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.
Note: a shuffle is not a bug. Joins and aggregations of large data need one. The goal is to avoid unnecessary shuffles and to make the necessary ones carry as little data as possible.

Common mistakes

Calling repartition "to be safe"

Each call is a full shuffle.

Joining before filtering

Rows that would be filtered out are shuffled for nothing; filter first or check that Catalyst pushes the filter down.

Window functions with many different partitionBy keys

Each different key is another shuffle.

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 questions

1. 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

Primary sources: RDD guide: shuffle operations · Performance tuning