Spark Partitioning & Shuffles
In Apache Spark, a partition is the unit of parallelism: each partition of a DataFrame is processed by one task on one executor core. How data is split into partitions, in memory during a job and on disk in storage, determines how much parallelism you get, how much data moves across the network, and whether a job finishes in minutes or runs for hours with one straggling task.
Most Spark performance problems trace back to partitioning: shuffles that move too much data, skewed keys that overload a single task, too few partitions that underuse the cluster, or too many small files that make every read slow. Understanding partitioning is the foundation of Spark performance tuning.
TL;DR
- Partitions = tasks. Aim for partitions of roughly 100–200 MB in memory, and enough of them to keep all cores busy.
- Shuffles (from joins, groupBy, distinct, repartition) redistribute data by key across the network. They're the most expensive operation.
repartition(n, cols)shuffles to a new layout.coalesce(n)merges partitions without a full shuffle (only to reduce).- AQE coalesces shuffle partitions and splits skewed ones automatically. Enable it and tune the targets.
- Skew: a few hot keys make giant partitions. Fix it with AQE skew joins, broadcasting, or salting.
- On disk:
partitionBylow-cardinality columns (date), avoid tiny files, and use clustering / Z-order in table formats.
Quick Example
Core Concepts
In-Memory Partitions
When Spark reads files, it creates input partitions based on file splits (spark.sql.files.maxPartitionBytes, default 128 MB). After a shuffle, the number of partitions is spark.sql.shuffle.partitions (default 200), or whatever AQE coalesces it to. Parallelism is bounded by min(partitions, total executor cores).
- Too few partitions: idle cores, huge tasks, memory spills, and out-of-memory errors.
- Too many partitions: scheduling overhead, tiny tasks, and many small output files.
Shuffles and Stages
A shuffle happens when rows must be regrouped by key (wide transformations). Map tasks write shuffle files partitioned by key hash, and reduce tasks fetch their partition from every map output across the network. Shuffles mark stage boundaries in the Spark UI. They cost disk I/O, network, serialization, and sorting, so minimizing shuffle volume (by filtering and projecting before shuffles, or broadcasting small tables) is the biggest optimization lever.
repartition vs coalesce
Beware: coalesce(1) before a write can pull the entire upstream stage into a single task, since coalesce propagates upstream without a shuffle boundary. Use repartition(1) when you truly need one file from a heavy computation, or better, avoid single-file outputs.
Adaptive Query Execution (AQE)
AQE (on by default since Spark 3.2) uses runtime shuffle statistics to:
- Coalesce small shuffle partitions toward a target size, which makes a fixed
shuffle.partitionsless critical. - Split skewed partitions in sort-merge joins into smaller sub-tasks.
- Switch join strategies (for example, to broadcast) when a side turns out to be small.
Data Skew
Skew occurs when some keys have vastly more rows than others (null keys, "unknown" values, a mega-customer). Symptoms: most tasks finish in seconds while a few run for an hour, spill, or fail. Remedies:
- AQE skew join handling for joins.
- Broadcast the smaller side, if it fits, to avoid shuffling the skewed side.
- Filter or separately handle junk keys (nulls, sentinels).
- Salting: add a random suffix to hot keys to spread them across partitions, aggregate in two phases, or replicate the other join side across salt values.
On-Disk Partitioning
df.write.partitionBy("event_date") creates directory partitions (event_date=2026-09-26/). Queries filtering on the partition column skip irrelevant directories (partition pruning). Guidelines:
- Partition by low-cardinality, commonly filtered columns (date, region). Never user IDs.
- Target partitions of at least ~1 GB. Over-partitioning creates millions of small files.
- Bucketing (
bucketBy) pre-shuffles tables by key for repeated joins, but it's fiddly. Table-format features are usually better.
Table Formats: Hidden Partitioning and Clustering
Apache Iceberg supports hidden partitioning (for example days(ts), bucket(16, id)) and partition evolution without rewriting data. Delta Lake offers liquid clustering and Z-ordering, which co-locate related data within files, so min/max statistics skip more files, even for high-cardinality columns. Regular compaction (OPTIMIZE, rewrite_data_files) fixes small files. See data lakehouse.
Best Practices
Size Partitions by Data, Not Defaults
Let AQE coalesce, set an advisory partition size (64–256 MB), and for very large shuffles raise shuffle.partitions so that the initial partitions aren't enormous. Check task durations and spill in the Spark UI.
Reduce Before You Shuffle
Filter, project, and pre-aggregate before joins and groupBys. Shuffling 10 columns instead of 80 cuts network and disk use proportionally.
Control Output File Sizes
Repartition by the write partition columns before writing (one shuffle, fewer files), use maxRecordsPerFile, and schedule compaction for tables written incrementally.
Investigate the Slowest Task
Stage summary metrics (max vs median task time, shuffle read size) reveal skew immediately. Fix skew rather than adding more executors.
Common Mistakes
Partitioning Tables by High-Cardinality Columns
partitionBy("user_id") creates millions of directories and tiny files, which crushes listing and metadata performance. Use bucketing or clustering instead.
Ignoring Null Key Skew
Joining on a column where 30% of rows are null sends all of them to one partition. Filter nulls out before the join (they won't match anyway), and union them back if needed.
coalesce(1) on a Heavy Pipeline
It serializes the whole final stage onto one core. Write in parallel, or repartition(1) only after the expensive work is done.
FAQ
How many partitions should a Spark job have?
Enough that each partition is roughly 100–200 MB, and at least 2–4 times the number of executor cores, so work is balanced. With AQE, set a generous initial shuffle partition count and let AQE coalesce down to the advisory size.
What is the difference between repartition and partitionBy?
repartition controls in-memory partitions during computation (it's a shuffle). partitionBy on a writer controls directory layout in storage. They're often combined: repartition("date").write.partitionBy("date") produces few files per date directory.
What is data skew in Spark?
An uneven distribution of rows across partitions, usually caused by a few very frequent key values, so some tasks process far more data than others and dominate the job's runtime. It's addressed with AQE skew handling, broadcasting, filtering hot or null keys, or salting.
What is the small files problem?
Having many tiny files (kilobytes to a few megabytes) in a table. Every file adds metadata and open overhead, which slows planning and reading. It's caused by over-partitioning and frequent small writes, and fixed with compaction, better write partitioning, and optimized writes in table formats.
Related Topics
- Apache Spark — Pillar overview
- Spark DataFrames — Transformations and plans
- Spark Performance Tuning — Memory, joins, and configuration
- Apache Iceberg — Hidden partitioning and compaction
- Data Lakehouse — Table layout in the lake
- Kafka Partitions — Partitioning in streaming systems