skip to content

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

level: middleimportance: must knowfreq 85%

answer

  1. one of them moves data, one does not
  2. which direction can the count go?
  3. no shuffle means no stage boundary
  4. the cheap one throttles what comes before it
  5. coalesce(1) before a write is the classic trap

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.

solid answer

~40 s

`repartition(n)` is a **wide** transformation: it shuffles every row, distributing round-robin (or by hash when you pass columns), so it can increase or decrease the partition count and it produces evenly sized partitions. `coalesce(n)` is a **narrow** transformation: it merges neighbouring partitions in place with no network transfer, so it is cheap, but it can only reduce the count and the resulting partitions are as uneven as the ones it glued together. The trap is that `coalesce` has no stage boundary, so it propagates upstream — `df.filter(...).coalesce(1).write...` runs the *filter itself* with one task, not just the write. When you want one output file but full parallel compute, use `repartition(1)`, which puts a shuffle between the work and the write. Use `coalesce` after a heavy filter that merely left many small, already-balanced partitions.

code

python · 11 lines
python
# Trap: one task scans, filters AND writes
(spark.read.parquet("s3://lake/events/")
      .filter("country = 'PT'")
      .coalesce(1)
      .write.mode("overwrite").parquet("s3://lake/out/"))

# Fix: shuffle boundary keeps the scan and filter parallel
(spark.read.parquet("s3://lake/events/")
      .filter("country = 'PT'")
      .repartition(1)
      .write.mode("overwrite").parquet("s3://lake/out/"))

go deeper

for a junior

Know the headline: repartition shuffles and can go up or down, coalesce only merges downward without a shuffle. Be able to say why coalesce is cheaper.

for a middle

Explain the mechanics — narrow versus wide, why no stage boundary means the parallelism cap travels upstream, and why coalesce cannot rebalance uneven partitions.

for a senior

Diagnose the real failure: a job that mysteriously runs with one task after a coalesce, and the tradeoff of paying a shuffle via repartition to restore upstream parallelism.

for a principal

Own the policy question of how output file sizing is handled across many pipelines: manual repartition conventions versus adaptive coalescing and rebalance hints, and the cost each imposes on a shared cluster.

## Two ways to change the partition count Both `repartition` and `coalesce` change how many partitions a DataFrame or RDD has, and that is where the resemblance stops. The difference is whether a shuffle happens, and a shuffle is what buys you the ability to move a row from any partition to any other. ## repartition: full shuffle, any count, even sizes `repartition(n)` writes every row to a shuffle file and reads it back on the target side. In the RDD API it is literally defined as `coalesce(n, shuffle = true)`. Because every row is redistributed, it can: - **increase** the partition count as well as decrease it, - produce **evenly sized** output partitions — with no columns given, Spark uses round-robin placement from a random starting offset, which balances rows regardless of how skewed the input partitions were. There are column-taking forms too. `repartition(n, "customer_id")` hash-partitions by that column, so all rows with the same value land in the same partition — this is what you want before a write partitioned by the same column, or before a `mapPartitions` that must see a whole group at once. Spark SQL's hash here is its own Murmur3-based `hash()` function applied modulo `n`, not the RDD `HashPartitioner`'s `key.hashCode`. `repartition("col")` with no number uses `spark.sql.shuffle.partitions` (200 by default). The cost is real: serialization, disk writes on the map side, network fetches on the reduce side, and a stage boundary. You pay it deliberately. ## coalesce: narrow merge, downward only `coalesce(n)` builds a new partition list where each output partition is the concatenation of several input partitions, preferring inputs already co-located on the same executor. No data crosses the network, no shuffle files are written, no stage boundary is created. It is essentially free. Two limitations follow from that design: 1. **It cannot increase the count.** `coalesce(200)` on a 50-partition DataFrame silently returns 50 partitions. There is no error; the request is simply ignored, because raising the count would require redistributing rows, which requires a shuffle. (`coalesce(n, shuffle = true)` opts back into one, and is just `repartition`.) 2. **It cannot rebalance.** If one input partition holds 90% of the rows, the merged partition containing it still holds 90% of the rows. ## The upstream-propagation trap This is the part interviews are really probing. Because `coalesce` introduces no stage boundary, the reduced parallelism applies to the *entire narrow chain above it* in the same stage. Consider: ```python (spark.read.parquet("s3://lake/events/") # 2,000 partitions .filter("country = 'PT'") .coalesce(1) .write.parquet("s3://lake/out/")) ``` A reasonable reader expects 2,000 tasks to scan and filter, then one task to write. What actually happens is one task that scans all 2,000 files, filters them, and writes — because scan, filter and write are all narrow and therefore fused into a single stage whose parallelism `coalesce(1)` has set to one. On a large input this turns a five-minute job into an hours-long one, or into an executor OOM. The fix is `repartition(1)`: the shuffle boundary lets the scan and filter run with full parallelism, and only the post-shuffle side runs single-threaded. You pay a shuffle to buy back the upstream parallelism, and on a heavily-filtered dataset that is a bargain. ## When each one is right Use **coalesce** when a stage has already produced many small, roughly balanced partitions and you just want fewer output files — classically after a highly selective filter, reducing 2,000 near-empty partitions to 50 — and when the reduction is modest so the upstream parallelism cap still leaves enough tasks. Use **repartition** when you need to *increase* parallelism, when partitions are unevenly sized, when you must group rows by a key before a write or a partition-local computation, or when the target count is drastically lower than the upstream count and you cannot afford to throttle the upstream stage. ## Interactions worth knowing Adaptive Query Execution (on by default since Spark 3.2) already coalesces over-provisioned *post-shuffle* partitions using runtime statistics, so a lot of the manual `coalesce`-after-`groupBy` tuning that older guides recommend is now unnecessary. Note also that an explicit `repartition(n)` historically blocked AQE's coalescing of that exchange, since you asked for a specific number on purpose; `spark.sql.adaptive.optimizeSkewsInRebalancePartitions` and the `REBALANCE` hint exist for the "give me evenly sized output, you pick the count" case. Finally, both operations are about *in-memory* partitions and therefore about task parallelism and output file count. Neither is related to `df.write.partitionBy(...)`, which creates directories in the output path.

  • Why does repartition(1) before writing a single file beat coalesce(1)?
    `coalesce(1)` creates no stage boundary, so the single-task limit propagates up the whole narrow chain — the scan and every transformation also run with one task. `repartition(1)` inserts a shuffle: everything above it keeps full parallelism and only the final write is single-threaded. You pay one shuffle to avoid serialising the entire pipeline.
  • What does repartition(n, "user_id") give you that repartition(n) does not?
    Co-location by key. The column form hash-partitions on `user_id`, so every row for a user lands in one partition — required before a partition-local aggregation, a `mapPartitions` that needs whole groups, or a write partitioned by the same column. The plain form places rows round-robin, which balances sizes but scatters each key across all partitions.
  • What happens if you call coalesce(500) on a DataFrame with 100 partitions?
    Nothing: you get 100 partitions back. `coalesce` can only merge, never split, and it fails silently rather than raising. Increasing the count requires redistributing rows, which requires a shuffle — that is `repartition(500)`, or equivalently `coalesce(500, shuffle = true)` on the RDD API.

repartition is unpacking every box and repacking them all evenly; coalesce is stacking existing boxes onto fewer pallets — cheap, but the heavy box stays heavy.

saying these in an interview costs you the question

  • Says coalesce and repartition are interchangeable, one is just faster
  • Thinks coalesce can increase the partition count
  • Claims coalesce(1) only affects the write step
  • Believes repartition guarantees equal partitions even when partitioning by a skewed column
  • Confuses either operation with write-side partitionBy directories

context