Menu

Apache Spark interview question · Question 4 of 5

What causes a shuffle in Spark?

  • Medium
  • conceptual
  • ~6 min
  • High relevance
  • 2 min read
  • Updated Oct 2026

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.

On this page
  1. Detailed explanation
  2. Seeing it
  3. Reducing shuffle cost
  4. Common mistakes

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

  1. Filter and select columns before wide operations.
  2. Broadcast small join sides.
  3. Pre-aggregate before joining.
  4. Avoid unnecessary repartition, and avoid a global orderBy when only per-group order is needed.

Common mistakes

  1. Saying “joins always shuffle” (broadcast joins do not shuffle the large side).
  2. Thinking coalesce and repartition are equivalent.

By Data Career Hub Editorial · Last reviewed Oct 2026 · Examples run on PySpark 4.2 in local mode; behaviour notes say where Spark 3.x differs

Progress is saved in this browser only. No account needed.

Search
Filter by type