Shuffles: Beginner
Narrow Transformations
A purely narrow chain
(products
.filter(F.col("price") > 50)
.select("product_name", "category", "price")
.orderBy("product_name"))Wide Transformations
- filter, select, withColumn, map
- Each executor works alone
- No network, no coordination
- Scales almost for free
- groupBy, join, distinct, orderBy
- Rows must be regrouped by key
- Data crosses the network
- Where the cost concentrates
Watch a wide op regroup the data
(orders
.groupBy("region")
.agg(F.count(F.lit(1)).alias("order_count"))
.orderBy(F.col("order_count").desc()))What a Shuffle Actually Is
Why Wide Costs and Narrow Does Not
| Narrow operation | Wide operation (shuffle) | |
|---|---|---|
| Data movement | None; stays on the executor | Across the network, all-to-all |
| Disk | None | Write buckets out, read them back |
| Coordination | None; fully independent | Every executor waits for the exchange |
| Cost | CPU only, near free | Disk + network + serialisation, dominant |
Spotting Shuffles in Your Code
Tag the wide operation
> From products, for in-stock items only (in_stock = 1), return the number of products in each category, highest count first. The filter is narrow; the grouping is the one wide operation that shuffles.
(products .filter(F.col("in_stock") == 1) .("category") .agg(F.(F.lit(1)).alias("n")) .orderBy(F.col("n").desc(), F.col("category").asc()))
- Tag each operation narrow or wide as you read a chain; the wide ones are the cost.
- Pile on narrow work freely; it is per-partition and nearly free.
- Treat every groupBy, join, distinct, and global sort as a deliberate shuffle.
- Filter early so the wide operation that follows shuffles less data.
- Don't assume all operations cost the same; narrow is free, wide is dominant.
- Don't add a wide operation casually; each one writes to disk, hits the network, and serialises.
- Don't forget a global orderBy is a shuffle; it needs a total order across partitions.
- Don't ignore that a shuffle is a barrier; the next stage waits for the slowest task.
> You are handed a Spark job that reads a large events table, filters it to one country, joins it to a small lookup of content titles, and counts views per title. It is slow, and you have not run it yet.
One category is free. The other can run your whole bill.
- Category
- Spark
- Difficulty
- beginner
- Duration
- 13 minutes
- Challenges
- 3 hands-on challenges
Topics covered: Narrow Transformations, Wide Transformations, What a Shuffle Actually Is, Why Wide Costs and Narrow Does Not, Spotting Shuffles in Your Code
Lesson Sections
- Narrow Transformations (concepts: paSparkExecutionModel)
A narrow transformation is one where each output partition is built from exactly one input partition. Think of filter: to decide which rows of a partition to keep, an executor only needs the rows already sitting in front of it. It never has to look at any other partition, on any other machine. The same is true of select, which picks columns, and withColumn, which computes a new one, and map, which transforms each row. Every one of these works on a partition in place, using only what is already t
- Wide Transformations (concepts: paSparkExecutionModel)
A wide transformation is one where an output partition needs data from many input partitions. The classic example is groupBy. To sum profit per region, every row for a given region has to end up in the same place so it can be added together, but those rows are scattered across every partition on every machine. Spark has to gather all the rows for each region together first, and that gathering pulls data from everywhere. Join is the same story. To match orders to products on a product id, every o
- What a Shuffle Actually Is (concepts: paShuffleOptimization)
The movement that a wide operation forces has a name, and you will hear it constantly: the shuffle. A shuffle is the physical redistribution of data across the cluster so that rows which need to be together end up together. It is the most expensive thing Spark does, and it follows directly from a wide transformation. Every wide operation triggers a shuffle; every shuffle is triggered by a wide operation. The two are the same event seen from two angles, the logical operation and its physical cost
- Why Wide Costs and Narrow Does Not (concepts: paShuffleOptimization)
Being precise about why the two categories differ so dramatically in cost is what lets you predict performance instead of memorising rules. A narrow operation reads a partition that is already in memory on an executor, applies a function, and produces output in memory on the same executor. The data never leaves the machine. The cost is just the CPU work of the function, which is usually trivial compared to anything involving disk or network. A wide operation, by contrast, pays 3 expensive costs
- Spotting Shuffles in Your Code (concepts: paShuffleOptimization)
The practical skill this lesson builds is reading your own code and seeing, before you run it, where the shuffles are. It rests on a small vocabulary of operations that signal a wide transformation. When you see any of these in a chain, a shuffle is coming, and that is where the cost will be. Everything else is narrow and cheap. Run your eye down a chain and tag each line. filter, select, withColumn: narrow, free. Then a groupBy: there is your shuffle, there is the cost. The exercise sounds triv