Partitions
The unit of parallelism in Spark: where partitions come from and why their number matters.
On this page
Show code in
Every code block on the page follows this.
You will learn
- What a partition is and how it maps to tasks
- How Spark decides the number of partitions when reading files
- When to use repartition and when coalesce
- How to tell if you have too few or too many partitions
What a partition is
A DataFrame is split into partitions, each processed by one task on one core. A DataFrame with 8 partitions on a cluster with 40 cores can only use 8 of them in a stage. A DataFrame with 40,000 tiny partitions keeps the scheduler busy launching tasks that each finish in milliseconds.
Where partition counts come from
| Situation | Number of partitions |
|---|---|
| Reading files (Parquet, ORC, CSV, JSON) | Files are split into chunks of up to spark.sql.files.maxPartitionBytes (128 MB). Small files are packed together, counting each file as at least spark.sql.files.openCostInBytes (4 MB). |
| After a shuffle (join, groupBy, distinct) | spark.sql.shuffle.partitions (200), then coalesced by AQE |
spark.range, parallelize | spark.default.parallelism, usually the total number of cores |
repartition(n) / coalesce(n) | Exactly n (coalesce can only go down) |
So a 10 GB Parquet table read with defaults gives about 80 partitions (10 GB / 128 MB), regardless of how many files it has, as long as the files are large enough to split.
df.rdd.getNumPartitions() # how many right now df.repartition(64) # full shuffle into 64 even partitions df.repartition(64, "customer_id") # shuffle by key df.coalesce(8) # merge down to 8, no full shuffle
-- SQL hints SELECT /*+ REPARTITION(64) */ * FROM orders; SELECT /*+ REPARTITION(64, customer_id) */ * FROM orders; SELECT /*+ COALESCE(8) */ * FROM orders;
repartition vs coalesce
repartition(n)
coalesce(n)
coalesce(1) before a write does not just produce one file: because it avoids a shuffle, Spark runs the entire upstream stage, joins and filters included, in a single task. For one output file after heavy work, repartition(1) is often faster, because the heavy part stays parallel.How many is right?
- Aim for partitions of roughly 100 to 200 MB of data in memory, a common rule of thumb.
- Aim for at least 2 to 4 times as many partitions as cores in a stage, so a slow task does not leave the rest of the cluster idle.
- Check the Spark UI: tasks taking under 100 ms suggest too many partitions; tasks spilling to disk or taking many minutes suggest too few or skewed ones.
Not to be confused with table partitioning
A table partitioned on disk (folders like date=2025-03-01/) is a storage layout for skipping data. In-memory partitions are a runtime concept for parallelism. Reading one folder can produce many in-memory partitions, and many small folders can be packed into a few. The Partitioning done right lesson covers table partitioning.
Common mistakes
coalesce(1) after heavy work
repartition before every step
Ignoring partition count when reading many tiny files
Key takeaways
- One task per partition per stage: partitions set the parallelism.
- File reads split into about 128 MB chunks; shuffles produce 200 partitions before AQE.
- repartition shuffles and balances; coalesce merges without a shuffle.
- Target roughly 100 to 200 MB per partition and several partitions per core.
Check yourself
3 questions1. With defaults, roughly how many partitions does reading a 6.4 GB splittable Parquet table produce?
Show the answer
50. 6.4 GB divided by 128 MB per partition is about 50.
2. Which can increase the number of partitions?
Show the answer
repartition. coalesce only merges partitions; repartition shuffles into any number.
3. Why can coalesce(1) make a whole job slow?
Show the answer
Without a shuffle, the entire upstream stage runs in one task. coalesce is narrow, so it reduces the parallelism of the stage it belongs to.
Go deeper
Jobs, stages and tasksSpark internals
Shuffle partitionsSpark internals
The small file problem
Primary sources: Performance tuning · Configuration