Lesson 13 / 25
Narrow and Wide Transformations: the Shuffle
Understand why some operations move data across the network and cost the most.
The expensive step
A narrow transformation (filter, select, withColumn, map) can compute each output partition from a single input partition, so it runs without moving data. A wide transformation (groupBy, join, distinct, repartition, orderBy, window with partitionBy) needs rows with related keys to meet, so Spark performs a shuffle: it writes intermediate data to local disk, sends it across the network and reads it back. A shuffle ends a stage and starts the next one. Most Spark tuning is about doing fewer, smaller shuffles. In the plan, a shuffle appears as an Exchange operator.
Seeing the shuffle in the plan, run
I ran this on Apache Spark 4.0.0 (PySpark, local mode, in the official Docker image). The Exchange hashpartitioning(city, 4) line is the shuffle. Spark first does a partial aggregation on each partition (partial_sum), shuffles by city into 4 partitions, then does the final aggregation. # numbers are internal IDs and change between runs.
orders.groupBy("city").agg(F.sum("amount")).explain()
Output:
== Physical Plan ==
AdaptiveSparkPlan isFinalPlan=false
+- HashAggregate(keys=[city#2], functions=[sum(amount#3)])
+- Exchange hashpartitioning(city#2, 4), ENSURE_REQUIREMENTS, [plan_id=66]
+- HashAggregate(keys=[city#2], functions=[partial_sum(amount#3)])
+- Project [city#2, amount#3]
+- Scan ExistingRDD[order_id#0L,customer#1,city#2,amount#3,order_date#4]Filter and select before wide operations
Fewer rows and columns entering a shuffle means less data written and moved. Catalyst often pushes filters down for you, but be explicit with early column selection and filtering.
Quick check: Which operation normally causes a shuffle?
- select of two columns
- groupBy aggregation
- filter on a column
- withColumn with a constant
Answer
groupBy aggregation — Grouping needs all rows of a key together, so data must move between partitions.