Apache Spark courseLesson 3 of 4
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.
On this page
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
- 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). - 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")
- Filter or handle the hot key separately. If most of the skew is
NULLkeys that would not match anyway, filter them out before the join, or process the hot key as its own branch. - 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.
- 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
- Setting
spark.sql.shuffle.partitionsonce and never revisiting it as data size changes. - Calling
repartitionwhencoalescewould do, paying for a full shuffle. - Broadcasting a table that is not actually small, causing driver or executor memory errors.
- 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.
Progress is saved in this browser only. No account needed.