skip to content

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