skip to content

Partitioning and Shuffles

Partitioning decides how work is split across the cluster, and a shuffle is what happens when data has to move between those splits. Nearly every Spark performance question — skew, spill, too many small files — resolves to something in this area.

on this pageshow

explore

questions

19

In Spark, what is a partition and what decides how many partitions a DataFrame starts with?

level: juniorimportance: must knowfreq 78%

answer

  1. one of these equals one task
  2. cores can only work on so many
  3. how big a chunk of a file
  4. the 128 MB packing target
  5. a shuffle resets it to 200

basics

~20 s

A partition is Spark's unit of parallelism: one partition is processed by one task on one core. The starting count comes from the input — file splits packed to spark.sql.files.maxPartitionBytes (128 MB by default) — not from a fixed number.

solid answer

~40 s

A **partition** is a slice of the dataset that Spark processes as a single **task**; the number of partitions is therefore the maximum parallelism of a stage. For a file-based read, Spark packs files and file splits into partitions targeting `spark.sql.files.maxPartitionBytes` (default 128 MB), so a 1 GB splittable Parquet dataset lands around 8 partitions, while 5,000 tiny JSON files land near 5,000 partitions. For `spark.range` or `sc.parallelize` with no explicit count, Spark uses `spark.default.parallelism`, which in cluster mode is the total number of executor cores. After a shuffle the count is reset to `spark.sql.shuffle.partitions` (default 200), which Adaptive Query Execution may then coalesce downward at runtime. The practical rule of thumb is a few hundred MB per partition and roughly 2–3 tasks per core, so the cluster stays busy without paying per-task overhead.

code

python · 9 lines
python
df = spark.read.parquet("s3://lake/events/dt=2026-08-20/")
print(df.rdd.getNumPartitions())

# packing target for the scan, in bytes
print(spark.conf.get("spark.sql.files.maxPartitionBytes"))

# post-shuffle partition count
agg = df.groupBy("country").count()
print(agg.rdd.getNumPartitions())

go deeper

for a junior

Be ready to say that a partition is the unit of work behind one task, and to name at least one way to inspect or change the count, such as getNumPartitions and repartition.

for a middle

Explain where the count comes from at each stage: file packing against maxPartitionBytes, default parallelism for generated data, and the reset to 200 after any shuffle.

for a senior

Show how you size partitions for a real job — target bytes per task, tasks per core, and the non-splittable-input trap — and how you read the evidence off the Spark UI stage page.

for a principal

Own the tradeoff between scheduler overhead and parallelism across a shared platform: defaults you set cluster-wide, when per-job overrides are justified, and how upstream file layout decisions govern every downstream job's parallelism.

## What a partition actually is A Spark DataFrame or RDD is a logical description of a dataset; physically it is a list of **partitions**. A partition is a contiguous chunk of rows that lives on one executor and is processed by exactly one **task**. Everything about Spark's parallelism follows from this: a stage that operates on 12 partitions launches 12 tasks, and those tasks run concurrently up to the number of free executor cores. If your cluster has 100 cores and your DataFrame has 4 partitions, 96 cores sit idle no matter how much data you have. Note the vocabulary trap: a Spark *task* is one partition's work inside one stage. It is not a MapReduce task (a whole mapper or reducer) and not a Flink subtask (a parallel instance of an operator). Within Spark, "partition" also has a second, unrelated meaning on the write side — Hive-style directory partitioning via `partitionBy` — which is about file layout, not about in-memory parallelism. ## Where the initial count comes from **File sources.** Spark's file scan builds partitions by packing file splits until it reaches a target size. The knob is `spark.sql.files.maxPartitionBytes`, whose default is 128 MB. Spark also charges each file a notional opening cost (`spark.sql.files.openCostInBytes`, 4 MB by default) so that many tiny files are bundled together rather than each becoming its own task. Two consequences follow directly: - Splittable formats (Parquet, ORC, uncompressed or LZ4/Snappy-block text) can be cut mid-file, so 1 GB of Parquet gives roughly eight 128 MB partitions. - Non-splittable input (a gzip-compressed CSV, for example) cannot be divided, so one 4 GB `.gz` file is exactly one partition and one task — a very common cause of "my 200-core cluster is running one core flat out". **Generated data.** `spark.range(n)` and `sc.parallelize(seq)` with no explicit partition count fall back to `spark.default.parallelism`, which is the total number of executor cores in cluster mode and the number of local threads in `local[n]` mode. **After a shuffle.** Any wide transformation — `groupBy`, `join`, `distinct`, `repartition` — resets the count to `spark.sql.shuffle.partitions`, whose default is 200. This is the single most surprising number for newcomers: a 5-row DataFrame that goes through a `groupBy` comes out with 200 partitions, 198 of them empty. Since Spark 3.2, Adaptive Query Execution is enabled by default (`spark.sql.adaptive.enabled = true`) and its coalescing feature merges those over-provisioned post-shuffle partitions using real runtime statistics, so the *physical* count you observe is often far lower than 200. **Cached and checkpointed data** keep whatever partitioning they had when they were materialized. ## Why the count matters Too few partitions means idle cores, and it also means each task must hold more data in memory — the direct route to spilling to disk or an executor OOM. Too many partitions means the scheduler pays fixed overhead (task serialization, launch, result reporting, one output file per partition on write) for tasks that do a few milliseconds of real work; tens of thousands of trivial tasks can make the driver the bottleneck. The conventional targets are: a few hundred MB of input per partition, and roughly two to three tasks per available CPU core so that stragglers are absorbed by the scheduler rather than leaving cores idle at the end of a stage. These are heuristics, not laws — a task that does expensive per-row work (a UDF calling a model) wants far smaller partitions than a task that just projects two columns. ## How to look at it In any language, `df.rdd.getNumPartitions()` reports the current count. The Spark UI's stage page shows the task count and the distribution of task durations and input sizes, which tells you not just how many partitions you have but whether they are the same size as each other. ## Changing it You raise or rebalance the count with `repartition(n)` (a full shuffle) and lower it cheaply with `coalesce(n)` (no shuffle, but it also caps the upstream stage's parallelism). On the read side you can shift `spark.sql.files.maxPartitionBytes` for the whole job; for a single skewed source, splitting or recompressing the input file is usually the better fix, because no Spark setting can subdivide a gzip file.

  • You read a single 4 GB gzip-compressed CSV and only one task runs. Why, and what do you do?
    Gzip is not splittable, so the whole file becomes one partition and one task regardless of `spark.sql.files.maxPartitionBytes`. No Spark setting fixes it. Either re-land the source as many smaller files or as a splittable format (Parquet, or bzip2/LZ4-framed text), or accept the single-threaded read and immediately `repartition` afterwards so at least the downstream stages run in parallel.
  • After a groupBy on tiny input you see 200 tasks, nearly all empty. Is that a problem, and what handles it?
    It is scheduler waste rather than a correctness problem: `spark.sql.shuffle.partitions` defaults to 200 regardless of data size. With Adaptive Query Execution on — the default since Spark 3.2 — Spark coalesces those post-shuffle partitions using runtime statistics, so far fewer tasks actually launch. Without AQE you would lower the setting for that job.
  • How does a partition relate to an executor?
    An executor is a JVM process with several cores; each core runs one task, and each task processes one partition. So an executor with 4 cores works on up to 4 partitions at a time, and a partition never spans two executors. Partition count sets available parallelism; executor cores set how much of it is realised concurrently.

Partitions are the boxes a moving job is split into: how fast the move goes depends on how many boxes there are and how evenly they are packed, not on how much stuff you own.

saying these in an interview costs you the question

  • Says the partition count is fixed at 200 for every DataFrame
  • Thinks more partitions is always faster, ignoring per-task overhead
  • Believes a gzip file can be split across tasks by tuning a config
  • Confuses in-memory partitions with partitionBy write directories
  • Thinks one partition equals one executor rather than one task

context

open as a page

In Spark, how does repartition differ from coalesce, and when is coalesce the wrong choice?

level: middleimportance: must knowfreq 85%

basics

~20 s

repartition performs a full shuffle, so it can raise or lower the partition count and rebalances rows evenly. coalesce only merges existing partitions without a shuffle, can only lower the count, and caps the parallelism of the whole upstream stage.

open as a page

In Spark, what files does a map task write during a sort-based shuffle?

level: middleimportance: must knowfreq 68%

basics

~20 s

Each Spark map task writes two local files: one data file whose records are grouped and ordered by destination reduce partition, plus a small index file giving each partition's byte offset inside that data file.

open as a page

In a Spark stage, which task metrics separate data skew from a slow executor?

level: middleimportance: must knowfreq 70%

basics

~20 s

Compare the stage's per-task quantiles in the Spark UI. Skew shows one task reading many times the median task's shuffle bytes or records; a slow executor shows normal input sizes with long durations concentrated on one host.

open as a page

A Spark stage fails with FetchFailedException — what happened, and what does Spark do next?

level: seniorimportance: must knowfreq 62%

basics

~20 s

A reduce task could not read shuffle blocks from a peer after retries, almost always because the executor holding those files died or the node is unreachable. Spark marks the lost map output as missing, recomputes it, and retries the stage.

open as a page

In a Spark join, how does salting a hot key spread its rows across many tasks?

level: seniorimportance: must knowfreq 58%

basics

~20 s

You append a random 0..N-1 salt to the join key on the big side and replicate every row of the small side N times, once per salt value. The composite key now hashes into N partitions instead of one.

open as a page

In Spark, how does HashPartitioner assign a key, and when do you want RangePartitioner instead?

level: middleimportance: should knowfreq 55%

basics

~20 s

HashPartitioner sends a key to partition nonNegativeMod(key.hashCode, numPartitions) — same key, same partition, but a hot key makes one partition huge. RangePartitioner samples the data, derives sorted range boundaries, and is what makes globally ordered output possible.

open as a page

In Spark, what makes a transformation narrow or wide, and why does the difference matter?

level: middleimportance: should knowfreq 62%

basics

~20 s

A transformation is narrow when each output partition draws from a bounded set of input partitions — map, filter, mapPartitions — so it runs in place. It is wide when an output partition may need rows from every input partition, which forces a shuffle: groupBy, join, distinct, repartition.

open as a page

In Spark, how does a reduce task learn where to fetch its shuffle blocks?

level: middleimportance: should knowfreq 50%

basics

~20 s

Every finished Spark map task reports a MapStatus to the driver's MapOutputTracker, recording which executor holds its output and how large each partition is. A reduce task asks the tracker for its partition's locations, then fetches those byte ranges directly.

open as a page

How do you choose a value for spark.sql.shuffle.partitions on a large job?

level: middleimportance: should knowfreq 75%

basics

~20 s

Size it from the shuffled data volume, targeting post-shuffle partitions of a few hundred megabytes each, and keep the count a comfortable multiple of total executor cores. The default of 200 is a fixed number that ignores your data size entirely.

open as a page

In a Spark join, why does broadcasting the small side remove the shuffle and the skew?

level: middleimportance: should knowfreq 62%

basics

~20 s

A broadcast hash join ships one whole small table to every executor and probes it from inside each existing partition of the large table. No rows are redistributed by key, so a dominant key value never concentrates into a single task.

open as a page

What does spark.sql.adaptive.coalescePartitions.enabled do to a stage's task count?

level: middleimportance: should knowfreq 68%

basics

~20 s

With Adaptive Query Execution on, Spark reads the finished map stage's real output sizes and merges contiguous small shuffle partitions into fewer, larger reducer tasks. It can only reduce the post-shuffle partition count, never raise it.

open as a page

A Spark job writes 40,000 tiny Parquet files a day — what caused it and how do you fix it?

level: seniorimportance: should knowfreq 50%

basics

~20 s

File count is partitions multiplied by the distinct partitionBy values each task touches: 200 shuffle partitions across 200 directories can emit 40,000 files. Fix it by repartitioning on the same columns you partition by before writing, so one task owns each directory.

open as a page

What makes a Spark shuffle task spill records to disk?

level: seniorimportance: should knowfreq 55%

basics

~20 s

A shuffle task spills when its sorter or aggregation map cannot acquire more execution memory for the records it is holding. Spark sorts what it has, writes it to local disk as a run, and merges the runs later.

open as a page

In Spark AQE, when does spark.sql.adaptive.skewJoin.enabled split a shuffle partition?

level: seniorimportance: should knowfreq 56%

basics

~20 s

AQE splits a shuffle partition when it is both more than five times the median partition size and larger than 256 MB. It cuts the oversized partition into smaller pieces and replicates the matching partition from the other side so the pieces join in parallel.

open as a page

For a nightly Spark job with a recurring hot key, how do you choose between AQE, salting, and isolating the key?

level: principalimportance: should knowfreq 38%

basics

~20 s

Start with the free options: leave adaptive skew handling on and broadcast the small side if you can. Salt when a known hot key recurs and the operator is an aggregation or an unbroadcastable join. Isolate the key into its own path when one or two values dominate everything.

open as a page

In Spark AQE, how can a sort-merge join become a broadcast hash join mid-query?

level: seniorimportance: nice to knowfreq 28%

basics

~20 s

Adaptive Query Execution re-plans at shuffle boundaries. Once a join side's shuffle has been written, Spark knows its exact size; if that is under the adaptive broadcast threshold, it replaces the planned sort-merge join with a broadcast hash join.

open as a page

How do you choose the partitionBy columns for a Spark-written table many teams will query?

level: principalimportance: nice to knowfreq 32%

basics

~20 s

Partition by the low-cardinality columns that most queries filter on — usually a date — sized so each directory holds hundreds of megabytes. High-cardinality columns explode into tiny files and slow directory listing; use clustering inside a partition for finer pruning.

open as a page

For a shuffle-heavy Spark platform on elastic infrastructure, where should shuffle files live?

level: principalimportance: nice to knowfreq 30%

basics

~20 s

Choose by how disposable your compute is. Executor-local disk is fastest but dies with the executor; a per-node external shuffle service outlives executors; block migration on decommission covers planned removal; a disaggregated shuffle service decouples shuffle storage from compute entirely.

open as a page