A Spark partition is a slice of a dataset that one task processes on one core. Partitioning therefore decides almost everything about a job's physical behaviour. The count sets how much of the cluster can work at once. The size sets how much memory each task needs and whether it spills to disk. The distribution decides whether the stage finishes when the average task does or when the slowest one does.
Spark chooses partition counts at three boundaries: when it reads files, when it shuffles for a join or aggregation, and when it writes output. Each boundary has its own rules and its own settings. This article works through each one with the arithmetic, then covers the operators that change partitioning explicitly, skew, and output layout. Adaptive query execution, the shuffle mechanism and partition pruning have their own pages; here they appear only where they change partition counts.
What a partition is and why the count matters
Within a stage, Spark runs one task per partition. A stage with 800 partitions on a cluster with 200 cores runs in four waves. A stage with 150 partitions on the same cluster leaves 50 cores idle for its whole duration. Too few partitions waste the cluster and make each task large: it may exceed execution memory and spill, or hit out-of-memory errors. Too many make each task tiny, so scheduling overhead, per-task startup and, for shuffles, the number of shuffle blocks (map tasks times reduce partitions) dominate.
A useful rule of thumb is to size partitions by bytes, not by count: somewhere around 100 to 200 MB of input per task for scans, and an advisory 64 to 128 MB per shuffle partition. Then check that the count is at least two to three times the core count, so waves finish evenly. Narrow transformations such as filter, select and map keep the partition count of their parent. Only a shuffle, or an explicit operator, changes it. See stages and tasks for how the stage boundary itself is drawn.
Boundary one: how file reads are split
For file sources such as Parquet, ORC and JSON, three settings decide read partitions. spark.sql.files.maxPartitionBytes (default 128 MB) caps the bytes per partition. spark.sql.files.openCostInBytes (default 4 MB) is an estimated cost of opening a file, expressed as bytes that could have been read in the same time. spark.sql.files.minPartitionNum, which defaults to the leaf-node default parallelism (normally the total core count), is a suggested minimum. Spark computes a split size as the smaller of maxPartitionBytes and the larger of openCostInBytes and bytes-per-core. Total bytes include one open cost per file. It then cuts splittable files into chunks of that size and packs chunks into partitions greedily, largest first. Each chunk added to a partition is charged its size plus the open cost.
# How Spark sizes file read partitions (mirrors FilePartition.maxSplitBytes).
MB = 1024 * 1024
def max_split_bytes(file_sizes, max_partition_bytes=128 * MB,
open_cost=4 * MB, min_partition_num=200):
total = sum(size + open_cost for size in file_sizes)
bytes_per_core = total // min_partition_num
return min(max_partition_bytes, max(open_cost, bytes_per_core))
def pack(file_sizes, split, open_cost=4 * MB):
"""Greedy packing of (splittable) file chunks into read partitions."""
chunks = []
for size in file_sizes:
off = 0
while off < size:
chunks.append(min(split, size - off))
off += split
chunks.sort(reverse=True)
parts, current = 0, 0
for c in chunks:
if current and current + c > split:
parts, current = parts + 1, 0
current += c + open_cost
return parts + (1 if current else 0)
small = [1 * MB] * 1000
s = max_split_bytes(small)
print(s // MB, "MB split ->", pack(small, s), "partitions") # 25 MB -> 200
big = [10 * 1024 * MB] * 10
s = max_split_bytes(big)
print(s // MB, "MB split ->", pack(big, s), "partitions") # 128 MB -> 800Two worked cases on a 200-core cluster. With 1,000 files of 1 MB, total cost is 1,000 × 5 MB = 5,000 MB, bytes-per-core is 25 MB, and the split size becomes 25 MB. Each file costs 5 MB, so five files fit per partition, giving 200 partitions: one wave, each task reading only 5 MB. The open cost exists to stop Spark from packing hundreds of tiny files into one task. With 10 files of 10 GB, bytes-per-core is about 500 MB. The 128 MB cap wins, and each file splits into 80 chunks, giving 800 partitions and four waves.
Splitting assumes the format allows it. Parquet and ORC split on row-group or stripe boundaries, so a chunk that contains no row-group start reads almost nothing. A gzip-compressed CSV or JSON file cannot be split at all: one file is one partition regardless of size, which is a common reason for a single task taking an hour. Partition pruning and filter pushdown reduce which files and row groups are read, but not how the survivors are split; see partition pruning.
Boundary two: shuffle partitions
A wide operation such as groupBy, a sort-merge join, distinct or a window redistributes rows so that equal keys meet in the same task. Map-side tasks write their output bucketed by hash(key) modulo N, and N reduce tasks fetch their bucket from every map task. The mechanics are in Spark shuffle internals. N is spark.sql.shuffle.partitions, with a default of 200. That default is wrong for most jobs: too many for a 2 GB aggregation, far too few for a 2 TB join, where each reducer would take 10 GB.
Sizing N statically is simple arithmetic. Take the shuffle write size from the Spark UI for a representative run and divide by a target partition size. A 500 GB shuffle at 128 MB per partition needs about 4,000 partitions. With adaptive query execution enabled, which is the default in current Spark, you can set N high and let Spark merge contiguous small partitions at runtime, because AQE's coalescing only merges partitions. There is a trap in the defaults. spark.sql.adaptive.coalescePartitions.parallelismFirst is true, which makes coalescing ignore the 64 MB advisory size and respect only a 1 MB minimum, to maximise parallelism. On a busy shared cluster, the documentation recommends setting it to false.
# Shuffle sizing on a busy shared cluster
spark.conf.set("spark.sql.adaptive.enabled", "true")
spark.conf.set("spark.sql.adaptive.coalescePartitions.enabled", "true")
# Default true: ignores the 64 MB advisory size when coalescing and keeps
# many 1 MB partitions to maximise parallelism. The docs recommend false
# on a busy cluster so partitions coalesce toward the advisory size.
spark.conf.set("spark.sql.adaptive.coalescePartitions.parallelismFirst", "false")
spark.conf.set("spark.sql.adaptive.advisoryPartitionSizeInBytes", str(128 * 1024 * 1024))
# Start high and let AQE merge down; too few initial partitions cannot be split later
spark.conf.set("spark.sql.shuffle.partitions", "2000")When one side of a join fits under spark.sql.autoBroadcastJoinThreshold (default 10 MB), Spark can broadcast it and avoid shuffling the large side at all, so partition counts on that side stay as they were read; see broadcast joins.
Operators that change partitioning explicitly
Five operators let you control partitioning directly, and they differ in whether they shuffle and in how rows are assigned.
| Operator | Shuffles? | Assignment | Use for |
|---|---|---|---|
| repartition(n) | Yes | Round robin, even sizes | Rebalancing after a filter left uneven partitions |
| repartition(n, col) | Yes | Hash of col mod n | Co-locating keys before a write or repeated joins |
| repartitionByRange(n, col) | Yes | Sampled range boundaries | Globally ordered output, range-clustered files |
| coalesce(n) | No | Merges existing partitions | Reducing file count at the end of a job |
| REBALANCE hint | Yes | AQE splits skew, merges small | Even output files without manual counts |
coalesce deserves a warning. It avoids a shuffle by merging partitions locally, but it does so in the same stage as the work before it. A job that reads 800 partitions, does expensive parsing and then calls coalesce(10) runs the parsing on only 10 tasks, because the stage now has 10 partitions from its start. If the upstream work is heavy, repartition(10) is faster despite the shuffle: the heavy work keeps 800 tasks and only the final redistribution narrows. Range partitioning samples the data to choose boundaries. That costs an extra pass, and it produces skew when one value is very frequent, since every row with that value must land in the same range.
Skew: when one partition is the whole job
Hash partitioning spreads distinct keys evenly, but not rows. If one key holds 30% of the rows, one task holds 30% of the stage, and the stage takes as long as that task. In the Spark UI, skew shows as a large gap between median and maximum task duration, or between median and maximum shuffle read size, in the stage's summary metrics. Typical culprits are null keys and placeholder values.
AQE handles skew in sort-merge joins. It treats a partition as skewed if it is both larger than spark.sql.adaptive.skewJoin.skewedPartitionFactor (5.0) times the median and larger than skewedPartitionThresholdInBytes (256 MB). It then splits that partition and replicates the matching partition from the other side; adaptive query execution covers the mechanics. It does not fix skew in aggregations or in joins it cannot identify. For those, filter out junk keys such as nulls first, then salt the hot keys.
from pyspark.sql import functions as F
SALT = 16
# Big side: spread each key over SALT buckets at random
clicks_s = clicks.withColumn("salt", (F.rand(seed=7) * SALT).cast("int"))
# Small side: replicate each row once per bucket so every salted key finds its match
ads_s = ads.crossJoin(spark.range(SALT).withColumnRenamed("id", "salt"))
joined = clicks_s.join(ads_s, ["ad_id", "salt"]).drop("salt")Salting multiplies the small side by the salt factor, so apply it only to hot keys when the small side is large. For aggregations, use a two-phase approach: aggregate by (key, salt), then aggregate the partial results by key.
Boundary three: write layout and the small-file problem
Each write task writes at least one file for each output directory it has rows for. With partitionBy on a column that has D distinct values, and T write tasks whose rows are spread across all values, the job can produce up to T × D files. At 2,000 tasks and 365 dates, that is 730,000 files, most of them tiny. Readers pay an open cost for each one, and object-store listing slows down.
The fix is to shuffle by the partition column before writing, so that each directory's rows meet in one task, and cap file size with maxRecordsPerFile for large directories:
(events
.repartition("event_date") # all rows for a date meet in one task
.sortWithinPartitions("event_date", "user_id")
.write
.option("maxRecordsPerFile", 5_000_000) # cap file size for the one big date
.partitionBy("event_date")
.mode("overwrite")
.option("partitionOverwriteMode", "dynamic") # replace only dates present
.parquet("s3://lake/events/"))Choose partitionBy columns with modest cardinality that queries filter on, such as a date. Never partition by a user id. Sorting within partitions clusters values so that Parquet min/max statistics make row-group skipping effective. When overwriting partitioned tables, set spark.sql.sources.partitionOverwriteMode to dynamic to replace only the partitions present in the output. With the default, static mode, an overwrite without a partition spec first deletes every existing partition of the table, including dates the job never touched. bucketBy on a table stored in the metastore is the other layout tool: it pre-hashes data into a fixed number of buckets so that later joins on the bucket key can skip the shuffle. It only helps if both sides use the same bucket count and key.
Worked example: a daily events job
A job reads a day of click events, 480 GB of Parquet in 3,800 files, on 400 cores. It joins them to a 40 GB ads table and writes aggregates partitioned by country. Read: bytes-per-core is about 1.2 GB, so the 128 MB cap applies and the scan produces about 3,850 partitions, roughly ten waves of 128 MB tasks. Shuffle: the join shuffles about 350 GB after column pruning. With the default of 200 partitions, each reducer would take 1.75 GB and spill heavily. Setting shuffle partitions to 3,000 and parallelismFirst to false gives about 120 MB per partition, and AQE merges the tail of small ones.
The Spark UI then shows one join task at 40 minutes against a median of 20 seconds. The cause is ad_id 0, a placeholder for organic clicks that holds 18% of rows. Those rows never match an ad, so filtering them before the join and adding them back afterwards removes the straggler. Write: the output has 220 countries. Without a repartition, 3,000 tasks could write up to 660,000 files. With repartition by country and maxRecordsPerFile, it writes a few hundred files; the largest country gets several files, though still from one task, which the REBALANCE hint would spread.
Failure modes
- Default 200 shuffle partitions on a large join: multi-gigabyte reducers, spill, executor out-of-memory, and retries that fail the same way.
- Unsplittable compressed input: one gzip file becomes one task that runs for the whole job.
- coalesce collapsing upstream work: an expensive stage runs on a handful of cores.
- Small-file explosion: partitionBy without a matching repartition writes tasks × values files.
- Skew hidden by averages: the stage duration is set by the maximum task, which average metrics hide.
- Over-partitioned small jobs: thousands of millisecond tasks where scheduler overhead exceeds the work.
What to do next
- For your heaviest job, record input partitions, shuffle write size and task duration percentiles per stage from the Spark UI.
- Set spark.sql.shuffle.partitions from shuffle bytes divided by about 128 MB, and set coalescePartitions.parallelismFirst to false on shared clusters.
- Look for gzip or other unsplittable inputs and convert them to Parquet at ingestion.
- Compare maximum and median task time in each stage; filter junk keys and salt any key above a few percent of rows.
- Replace late coalesce calls that follow heavy work with repartition, and repartition by the partitionBy column before every partitioned write.
- Count output files per partition directory after each run, and alert when the count grows.