Apache Spark interview questionsQuestion 4 of 5
Apache Spark interview question · Question 4 of 5
What causes a shuffle in Spark?
Short answer
A shuffle happens when a transformation needs rows with the same key to be in the same partition, which usually means moving data between executors. Wide transformations cause it: groupBy and aggregations, most joins, distinct, repartition, window functions with partitionBy, and global orderBy. A shuffle writes data to disk, sends it over the network and starts a new stage, so it is often the most expensive part of a job.
Detailed explanation
| Causes a shuffle | Usually does not |
|---|---|
groupBy().agg(), distinct, dropDuplicates |
select, filter, withColumn |
| Sort-merge and shuffle-hash joins | Broadcast hash join (large side stays put) |
repartition(n) / repartition(col) |
coalesce(n) (merges partitions without a full shuffle) |
Window with partitionBy, global orderBy |
union |
Seeing it
from pyspark.sql import SparkSession, functions as F
spark = SparkSession.builder.master("local[2]").appName("example").getOrCreate()
df = spark.range(0, 100_000).withColumn("k", F.col("id") % 10)
df.groupBy("k").count().explain()
The plan contains Exchange hashpartitioning(k, ...), which is the shuffle, with a partial_count before it: Spark pre-aggregates each partition so only one row per key per partition crosses the network.
Reducing shuffle cost
- Filter and select columns before wide operations.
- Broadcast small join sides.
- Pre-aggregate before joining.
- Avoid unnecessary
repartition, and avoid a globalorderBywhen only per-group order is needed.
Common mistakes
- Saying “joins always shuffle” (broadcast joins do not shuffle the large side).
- Thinking
coalesceandrepartitionare equivalent.
Progress is saved in this browser only. No account needed.