Menu

Apache Spark course · Lesson 3 of 4

Spark Partitions, Shuffles and Data Skew

Learn how Spark splits data into partitions, why shuffles are expensive, how to recognise data skew and which fixes (AQE, broadcast, salting) apply.

  • Advanced
  • 4 min read
  • Updated Oct 2026
On this page
  1. Partitions
  2. Narrow and wide transformations
  3. Why shuffles are expensive
  4. What data skew looks like
  5. Fixes, in the order to try them
  6. Reading a plan
  7. Common mistakes
  8. Interview relevance
  9. Key takeaway

Most Spark performance problems come down to three connected ideas: how data is split into partitions, how much data has to move in a shuffle, and whether the work is spread evenly or concentrated by skew.

Partitions

A partition is a chunk of a dataset processed by one task. Spark runs one task per partition, and a core runs one task at a time. So:

  • Too few partitions leave cores idle and make each task large.
  • Too many tiny partitions waste time on scheduling overhead.
  • Very uneven partitions make a few tasks run long while the rest wait.

You can inspect and change partitioning:

df.rdd.getNumPartitions()        # how many partitions now
df.repartition(200, "customer_id")  # full shuffle, hash-partition by column
df.coalesce(20)                  # reduce partitions without a full shuffle

repartition can increase or decrease the count and always shuffles. coalesce only decreases, merging neighbouring partitions without a full shuffle, so it is cheaper but can leave partitions uneven.

Narrow and wide transformations

A narrow transformation (filter, select, map) lets each output partition depend on one input partition, so no data moves between machines.

A wide transformation (groupBy, most joins, distinct, repartition, window functions with partitionBy) needs rows with the same key to end up together. Spark must redistribute data across the network: a shuffle.

Why shuffles are expensive

A shuffle writes intermediate data to disk on the map side, moves it across the network and reads it on the reduce side. It also marks a stage boundary: Spark cannot start the next stage until the previous one finishes. Reducing the amount of data shuffled, or the number of shuffles, is usually the biggest tuning lever. Filter and select early so less data enters the shuffle.

The number of partitions after a shuffle is controlled by spark.sql.shuffle.partitions, which defaults to 200. That default suits neither tiny nor huge jobs, which is one reason Adaptive Query Execution exists.

What data skew looks like

Skew means some keys have far more rows than others. After a shuffle by key, one partition receives the “hot” key and its task runs for much longer than the rest.

Typical signs in the Spark UI:

  • In a stage, most tasks finish quickly and a few take many times longer.
  • The maximum task duration and shuffle-read size are far above the median.
  • Occasional out-of-memory errors or spills to disk on a single task.

Common causes are NULL or default keys (for example -1 or "unknown"), a very popular customer or product, and joins on low-cardinality columns.

Fixes, in the order to try them

  1. Let Adaptive Query Execution (AQE) help. With spark.sql.adaptive.enabled (on by default in recent versions), Spark can coalesce small shuffle partitions, switch a join to a broadcast join when one side turns out to be small, and split skewed partitions in sort-merge joins (spark.sql.adaptive.skewJoin.enabled, also on by default).
  2. Broadcast the small side of a join. If one side is small (default threshold for automatic broadcast is 10 MB, set by spark.sql.autoBroadcastJoinThreshold), broadcasting it avoids shuffling the large side at all:
from pyspark.sql import functions as F
result = big.join(F.broadcast(small), "key")
  1. Filter or handle the hot key separately. If most of the skew is NULL keys that would not match anyway, filter them out before the join, or process the hot key as its own branch.
  2. Salt the key. Add a random suffix to the hot key on the large side and replicate the small side across all suffixes, so one hot key becomes several balanced ones. It works on any Spark version but makes the code more complex.
  3. Re-think the data model. Sometimes the real fix is pre-aggregating, or choosing a better partitioning key upstream.

Reading a plan

df.explain() shows Exchange nodes, which are shuffles. A BroadcastHashJoin with a BroadcastExchange means the small side was broadcast; a SortMergeJoin with Exchange hashpartitioning on both sides means both sides were shuffled.

Common mistakes

  1. Setting spark.sql.shuffle.partitions once and never revisiting it as data size changes.
  2. Calling repartition when coalesce would do, paying for a full shuffle.
  3. Broadcasting a table that is not actually small, causing driver or executor memory errors.
  4. Optimising before looking at the Spark UI.

Interview relevance

Expect: “What causes a shuffle?”, “What is data skew and how do you fix it?” and “Broadcast join versus sort-merge join”. A strong answer names the symptom (a few slow tasks), the diagnostic (Spark UI) and an ordered set of fixes.

Key takeaway

Fewer, smaller shuffles and evenly sized partitions make Spark jobs fast. Measure first, let AQE do what it can, then apply targeted fixes.

By Data Career Hub Editorial · Last reviewed Oct 2026 · Configuration defaults checked on Spark 4.2; Adaptive Query Execution is enabled by default since Spark 3.2

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

Search
Filter by type