Spark's coalesce and repartition both change how many partitions a DataFrame has, and most people learn them as a pair of rules: coalesce to go down, repartition to go up or to rebalance. The rules are not wrong, but they leave out the part that causes slow jobs. Coalesce does not add a stage boundary, so it reaches backwards and changes how many tasks run the work before it. Repartition adds a shuffle, which costs I/O but protects everything upstream.
This article explains both operators from the physical plan up, works through a job where the wrong choice costs five times the runtime, and covers balance, hash and range variants, retries, file layout on write and how adaptive query execution changes the picture. It assumes you know what a partition and a shuffle are; Spark partitioning covers that ground and the shuffle covers the exchange in detail.
Two operators, one difference: whether there is a shuffle
A shuffle is a stage boundary: the tasks before it write their output into files bucketed by target partition, and the tasks after it fetch those files over the network. It costs disk, network and serialisation, and it decouples the two sides. The number of tasks before the boundary is independent of the number after.
repartition(n) always shuffles. In the DataFrame API it deals rows round-robin into n partitions, giving near-equal row counts, and it can increase or decrease the count. coalesce(n) never shuffles. It creates a narrow dependency in which each new partition is the union of several existing ones; the Spark documentation describes going from 1,000 partitions to 100 as each new partition claiming 10 of the current ones. It can only decrease the count: asking for more partitions than exist leaves the count unchanged, silently.
Because coalesce adds no boundary, the new partitions are computed by the same stage that computes the old ones. That single fact explains both its speed when used well and its damage when used badly.
What each one puts in the plan
You can see the difference in explain(). Coalesce shows up as a Coalesce node directly above the operators it shares a stage with. Repartition shows up as an Exchange with a partitioning scheme and, in recent versions, a shuffle origin that records why the exchange exists.
-- clicks.coalesce(64).explain() (abridged; exact text varies by version)
Coalesce 64
+- *(1) Filter (type = click)
+- *(1) Project [...]
+- FileScan parquet [...]
-- clicks.repartition(64).explain()
Exchange RoundRobinPartitioning(64), REPARTITION_BY_NUM
+- *(1) Filter (type = click)
+- *(1) Project [...]
+- FileScan parquet [...]
-- clicks.repartition("event_date").explain()
Exchange hashpartitioning(event_date, 200), REPARTITION_BY_COL
+- ...The shuffle origin matters later, when adaptive execution decides which exchanges it may adjust: an exchange the user requested with an explicit number is kept at that number. Note also the default in the third plan: repartition(col) without a number uses spark.sql.shuffle.partitions, which defaults to 200, regardless of how much data you have. With adaptive execution enabled the plan is wrapped in an AdaptiveSparkPlan node and the final layout appears only after the query runs.
The trap: coalesce pulls its parallelism upstream
Consider a daily job on a 400-core cluster. It reads 400 GB of Parquet, which with the default spark.sql.files.maxPartitionBytes of 128 MB becomes about 3,200 read tasks. Each task parses a JSON column and filters to click events, about 30 seconds of CPU per task. The filter keeps 2 percent, about 8 GB, and the team wants about 64 output files of roughly 128 MB, so they add coalesce(64) before the write.
Without the coalesce, the read stage runs 3,200 tasks on 400 cores: 8 waves of 30 seconds, about 4 minutes. With coalesce(64), the scan, parse and filter all belong to the same stage as the coalesce, which now has 64 tasks. Each task reads its 50 assigned splits one after another: 50 times 30 seconds is 25 minutes, and 336 cores sit idle the whole time.
With repartition(64) instead, the read stage keeps its 3,200 tasks and its 4 minutes, writes about 8 GB of shuffle data, and a second stage of 64 tasks reads and writes it. On typical cluster networks moving 8 GB takes well under a minute. Total: about 5 minutes against 25. The shuffle you avoided was cheap; the parallelism you gave up was not.
The extreme case is coalesce(1) to produce a single output file, which runs the entire upstream stage, back to the last shuffle or the scan, in one task. The Spark documentation warns about exactly this drastic coalesce and recommends repartition instead. The general rule: coalesce is safe when the work since the last stage boundary is cheap relative to the shuffle you would otherwise pay, and the input partitions are reasonably even.
Balance: merged neighbours versus dealt cards
Coalesce merges existing partitions without looking at their sizes. If the inputs are uneven, say a filter emptied some partitions and left others full, the merged partitions inherit that unevenness, and one slow task can set the stage time. For RDDs, the default coalescer also prefers to group partitions by data locality, which can make groups uneven in count as well as size.
Round-robin repartition deals rows like cards, so partitions end up with near-equal row counts. Equal rows are not always equal bytes, since rows vary in width, but they are usually close enough. That balance is the reason to pay for a shuffle even when the count is not changing: repartition(n) after a skewed filter or join evens out the next stage's tasks.
Check balance rather than assuming it. Group by spark_partition_id() and count, as in the code below, or read the task duration and input size distributions for the stage in the Spark UI. A maximum several times the median is the signature of imbalance.
Repartition by columns and by range
repartition(n, col) and repartition(col) hash-partition by the given columns, so all rows with the same key land in the same partition. That is what you want before a write partitioned by the same column, or to pre-partition data for repeated joins or aggregations on that key. It also imports the key's skew: if one key holds 30 percent of rows, one partition holds 30 percent of rows. And when the key has few distinct values, hash collisions leave some partitions empty and give others two keys; 30 dates hashed into 32 partitions will not produce 30 equal partitions.
repartitionByRange(n, col) samples the data to choose range boundaries, then assigns each row to the partition covering its value. Partitions hold contiguous key ranges, which suits writes where readers filter by range, or a global sort that follows. Sampling makes the boundaries approximate and adds a pass over the data. Neither variant fixes a single hot key; that needs salting, splitting the hot key across several partitions with an added random column, or adaptive skew handling.
Retries and determinism
Round-robin repartition has a subtle correctness issue. Which row goes to which partition depends on the order rows arrive in each input task. If a task after a shuffle is retried, for example because an executor was lost, and its input arrives in a different order, rows can be dealt differently the second time. Downstream tasks that already ran with the first dealing and tasks that rerun with the second can then disagree, producing duplicated or missing rows.
Spark addressed this (SPARK-23207) by sorting rows locally before the round-robin deal, controlled by spark.sql.execution.sortBeforeRepartition, which is on by default. The sort costs CPU and memory; turning it off trades that safety for speed. Later reports showed the sort is not a complete fix when the upstream values themselves are not deterministic, for example floating-point aggregates whose result depends on summation order. Hash and range partitioning do not have this problem because a row's destination depends only on its values. If exact once-only output matters, prefer partitioning on a key over round-robin.
Writing files: the reason most people reach for these
Most uses of these operators are about output files. Each write task produces at least one file per output directory it touches, so the number of files is roughly tasks times distinct partition values per task. A job with 200 tasks writing a table partitioned by date across 30 dates can create up to 6,000 small files, which slows every later read and stresses object-store listing.
from pyspark.sql import SparkSession, functions as F
spark = SparkSession.builder.getOrCreate()
events = spark.read.parquet("s3://bucket/raw/events/") # ~3,200 read splits
clicks = (events
.withColumn("payload", F.from_json("raw", "struct<page:string,ms:long>"))
.filter(F.col("type") == "click")) # keeps about 2 percent
# Anti-pattern for this job: parsing and filtering now run in only 64 tasks.
# clicks.coalesce(64).write.mode("overwrite").parquet("s3://bucket/clicks/")
# Keep the scan wide, shuffle only the survivors (about 8 GB) into 64 even partitions.
clicks.repartition(64).write.mode("overwrite").parquet("s3://bucket/clicks/")
def target_partitions(bytes_estimate, target_file_bytes=256 * 1024**2):
return max(1, -(-bytes_estimate // target_file_bytes)) # ceiling division
# Date-partitioned output: hashing on the partition column puts each date in one task.
(clicks
.repartition(target_partitions(8 * 1024**3), "event_date")
.write.partitionBy("event_date")
.option("maxRecordsPerFile", 5_000_000) # caps a hot date's file size
.mode("overwrite")
.parquet("s3://bucket/clicks_by_date/"))
# Inspect before you trust: the plan, the count and the row balance.
clicks.repartition(64).explain()
(clicks.repartition(64)
.groupBy(F.spark_partition_id().alias("pid")).count()
.orderBy(F.desc("count")).show(5))Hash-repartitioning on the same column used in partitionBy sends each date to exactly one task, so each date directory gets one file, or several if maxRecordsPerFile splits it. The cost is that a hot date is written by a single task. Here the file count follows the number of dates and n only sets tasks and hash collisions; for unpartitioned output, size n from data volume: estimated bytes divided by a target file size of roughly 128 MB to 1 GB.
How AQE changes the choice
With adaptive query execution enabled, the default since Spark 3.2, Spark measures shuffle output at runtime and can merge small adjacent shuffle partitions, controlled by spark.sql.adaptive.coalescePartitions.enabled (default true). The target size is spark.sql.adaptive.advisoryPartitionSizeInBytes (default 64 MB), but by default spark.sql.adaptive.coalescePartitions.parallelismFirst is true, which makes Spark ignore that target and respect only a 1 MB minimum, favouring parallelism. The Spark documentation recommends setting it to false on busy clusters to avoid many small tasks.
AQE's coalescing applies to exchanges it is allowed to resize. A repartition(n) with an explicit number is kept at n, because the number is treated as a user requirement. A repartition by columns without a number can be resized; check the plan for an AQEShuffleRead node to confirm what happened in your version.
Spark 3.2 also added the REBALANCE hint, written in SQL as /*+ REBALANCE(col) */. It asks for a shuffle whose purpose is reasonable partition sizes: AQE may merge small partitions and also split skewed ones, which it does not do for a plain repartition. It takes effect only when AQE is enabled. For output written partitioned by a column with skewed values, REBALANCE on that column is often the best default. More on runtime re-optimization is in Spark AQE.
The RDD API
On RDDs, repartition(n) is implemented as coalesce(n, shuffle=True). Without the shuffle flag, rdd.coalesce(n) behaves like the DataFrame version: narrow, decrease-only, and asking for more partitions does nothing. With the flag, it can increase the count and redistributes data. Code that mixes the APIs, for example converting to an RDD for a custom function and back, should check partition counts at each step with getNumPartitions().
Choosing: a decision table and failure modes
| Situation | Use | Why |
|---|---|---|
| Fewer, even partitions after cheap work | coalesce(n) | No shuffle; upstream cost is small |
| Fewer partitions after expensive work | repartition(n) | Keeps upstream parallelism |
| More partitions, or rebalance after skew | repartition(n) | Coalesce cannot increase or balance |
| Write partitioned by column | repartition(col) or REBALANCE(col) | One task per value; REBALANCE splits hot values |
| Range-filtered reads or global sort next | repartitionByRange | Contiguous key ranges |
| Single output file | repartition(1) on small output | Avoids a one-task upstream |
Failure modes to recognise: a stage whose task count dropped from thousands to a handful after someone added coalesce (the upstream trap); a final stage with one task far longer than the rest (key skew from hash repartitioning); thousands of tiny files per date (no repartition before a partitioned write); a repartition(col) that produced 200 partitions on a small dataset (the default shuffle partition count); and coalesce(n) that did nothing because n exceeded the current count. Each shows up in the Spark UI's stage page before it shows up in a bill.
What to do next
- Find every coalesce in your jobs and check, in the Spark UI, how many tasks the stage containing it now runs.
- Where upstream work is heavy, replace coalesce with repartition and compare stage times.
- Size output partitions from data volume and a target file size, not from fixed numbers copied between jobs.
- Before partitioned writes, repartition or REBALANCE on the partition column and set maxRecordsPerFile.
- With AQE on, decide whether parallelismFirst should be false for your cluster, and look for AQEShuffleRead in plans.
- Measure partition balance with spark_partition_id() counts on a sample before trusting a layout.
- Prefer key-based partitioning where retried tasks must produce identical output.