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

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).

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:

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:

  1. AQE skew join handling for joins.
  2. Broadcast the smaller side, if it fits, to avoid shuffling the skewed side.
  3. Filter or separately handle junk keys (nulls, sentinels).
  4. 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:

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

References