Skip to content
Everyone knows 4 min read · spark.sql.files.maxPartitionBytes

Partitions

The unit of parallelism in Spark: where partitions come from and why their number matters.

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

Read first

Comfortable with these? Read on.

TL;DR A partition is a chunk of a DataFrame's rows. Each stage runs one task per partition, so partitions are Spark's unit of parallelism. Too few and cores sit idle or tasks run out of memory; too many and scheduling overhead dominates.

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

SituationNumber 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, parallelizespark.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.

PySparkSpark SQL · Inspect and change partitions
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)

Full shuffle. Produces n partitions of roughly equal size. Can increase or decrease. Use it to fix skewed or too few partitions, or to cluster by a key before writing.

coalesce(n)

Merges existing partitions without a shuffle. Only decreases. Cheap, but partitions can end up uneven, and it reduces the parallelism of the whole preceding stage.
Watch out: 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

Runs the whole upstream stage in one task. Use repartition(1), or accept several files.

repartition before every step

Each call is a full shuffle. Only repartition with a reason.

Ignoring partition count when reading many tiny files

Even with packing, millions of files mean millions of file opens. Compact them.

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 questions

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

Primary sources: Performance tuning · Configuration