What you will be able to do
- Name the settings that decide how many partitions Spark creates when it reads files and when it shuffles
- Explain what repartitioning means and why changing the partition count of an existing DataFrame is a separate step from the settings that first set it
- Choose between repartition() and coalesce() based on whether a shuffle is acceptable and whether the partition count must go up or down
- Spot the drastic-coalesce trap that leaves most of the cluster idle
- Use the COALESCE, REPARTITION, REPARTITION_BY_RANGE and REBALANCE SQL hints, and predict which one wins when several are given
Key concept
Shuffle — A shuffle redistributes rows across partitions and nodes. Joins, aggregations and repartition() all need one. Most partition tuning is a trade-off: you pay for a shuffle to get parallelism, or you avoid one and accept the partitioning you already have.
1.Where partition counts come from
A partition is the unit of parallel work in Spark. One task processes one partition, so the number of partitions sets how many cores a stage can keep busy. Before you change that number, you need to know which settings produced it. Two moments matter: when Spark reads files, and when it shuffles.
When Spark reads file-based sources such as Parquet, JSON and ORC, spark.sql.files.maxPartitionBytes caps the amount of data packed into each input partition. Its default is 128 MB. A related setting, spark.sql.files.openCostInBytes, is the estimated cost of opening a file. Spark uses it when it packs several small files into one partition.
Checkpoint 1 of 5· Check yourself
With default settings, about how much data does Spark pack into a single partition when it reads a Parquet table?
spark.sql.files.maxPartitionBytes defaults to 134217728 bytes (128 MB). It only applies to file-based sources. 64 MB is the AQE advisory size used after a shuffle, which is a different setting.
“spark.sql.files.maxPartitionBytes | 134217728 (128 MB) | The maximum number of bytes to pack into a single partition when reading files.”Source: spark.apache.org
Joins and aggregations have to bring rows with the same key together, so they always shuffle. Their output partition count is set by spark.sql.shuffle.partitions, which defaults to 200. On Databricks you can also set it to auto, which turns on auto-optimized shuffle: Databricks then chooses the number from the query plan and the input size. One limit applies to Structured Streaming: you can't change this setting between restarts of a query that uses the same checkpoint location.
Repartitioning is what you do when the partition count or layout you inherited is not the one you want. Instead of changing a setting, you ask the DataFrame for a new partitioning. repartition() takes either a target number of partitions or a partitioning column, and the result is hash partitioned. If you leave out the number, Spark uses the default partition count. The next two sections cover repartition() and coalesce() in detail.
2.repartition(): any partition count, at the cost of a shuffle
repartition() returns a new DataFrame that is hash partitioned. You can pass a target number of partitions, one or more partitioning columns, or both. If you pass a Column as the first argument, Spark treats it as the first partitioning column. If you leave out the number, Spark uses the default partition count. Every row may move to a different partition, so repartition() always triggers a full shuffle. In return, you can raise the partition count as well as lower it, and rows with the same column values end up in the same partition.
Checkpoint 2 of 5· Fill the gap
This sample spreads the rows into 7 hash partitions keyed on the age column. Which method fills the blank?
df. ? (7, "age").select(
sf.spark_partition_id().alias("partition")
).distinct().sort("partition").show()Only repartition() accepts a partition count plus partitioning columns, and its output is hash partitioned. coalesce() takes nothing but a partition number.
Source: docs.databricks.comSources3
3.coalesce(): fewer partitions without a shuffle
coalesce(numPartitions) takes only a number. It creates a narrow dependency: each new partition takes over several existing ones, so no data moves through a shuffle. If you go from 1000 partitions to 100, each new partition claims 10 of the old ones. This makes coalesce() the cheap way to reduce the partition count, for example before a write that would otherwise produce many small files.
from pyspark.sql import functions as sf
spark.range(0, 10, 1, 3).coalesce(1).select(
sf.spark_partition_id().alias("partition")
).distinct().sort("partition").show()There are two limits. First, coalesce() can't add partitions. If you ask for more than you already have, the DataFrame keeps its current count. Second, a drastic coalesce hurts cluster utilization. Because there is no shuffle, the upstream work runs inside the reduced partitions. After coalesce(1), that work may run on a single node. If you call repartition(1) instead, Spark adds a shuffle, but the upstream partitions still run in parallel before the data is gathered into one.
| Aspect | repartition() | coalesce() |
|---|---|---|
| Arguments | numPartitions and/or partitioning columns | numPartitions only |
| Shuffle | Yes, adds a shuffle step | No, narrow dependency |
| Can increase partitions | Yes | No, stays at the current number |
| Upstream parallelism on a drastic reduction | Preserved | May collapse to fewer nodes |
Checkpoint 3 of 5· Check yourself
A DataFrame has 1000 partitions and you call df.coalesce(2000). What happens?
coalesce() can only merge partitions. If you request more partitions than exist, the count stays where it is. To increase it, use repartition().
“If a larger number of partitions is requested, it will stay at the current number of partitions.”Source: docs.databricks.com
Checkpoint 4 of 5· Exam question
A DataFrame with 1000 partitions is filtered down to a small result set. An engineer wants the final write to produce exactly 8 output files while spending as little execution time as possible on the operation itself: ```python result = spark.read.table("sales.transactions").filter(col("region") == "EMEA") result.write.mode("overwrite").parquet("/mnt/output/emea_sales") ``` Which change achieves this most efficiently?
Correct answer: A — Insert `result = result.coalesce(8)` before the write, since it merges existing partitions on the same executors and avoids a full shuffle across the cluster.
- A. `coalesce()` reduces the number of partitions by combining existing ones on the same executors, so it can shrink partition count without a full shuffle when the target is fewer partitions than the current count. This makes it the cheapest way to hit exactly 8 output files here.
- B. `repartition()` always performs a full shuffle to redistribute rows across the new partition count, regardless of direction, so it is more expensive than needed when only reducing the number of partitions.
- C. `spark.sql.shuffle.partitions` only controls the number of partitions produced by shuffle operations such as joins and aggregations; it has no effect on the partition count of a plain filter-then-write DataFrame operation.
- D. Hash partitioning on a single low-cardinality column like `region` does not guarantee balanced output; rows with the same region hash to the same partition, so file sizes can still be very uneven, and this still forces a full shuffle.
- E. Caching a DataFrame stores its existing partitions in memory or on disk; it does not change the partition count at all, so the write would still emit the original number of files.
Sources4
4.The same controls in SQL: partitioning hints
SQL users get the same controls through hints. COALESCE, REPARTITION and REPARTITION_BY_RANGE are equivalent to the coalesce, repartition and repartitionByRange Dataset APIs. COALESCE takes only a partition number. REPARTITION takes a number, columns, both or neither. REPARTITION_BY_RANGE requires columns, and the number is optional.
SELECT /*+ COALESCE(3) */ * FROM t;
SELECT /*+ REPARTITION(3) */ * FROM t;
SELECT /*+ REPARTITION(c) */ * FROM t;
SELECT /*+ REPARTITION(3, c) */ * FROM t;
SELECT /*+ REPARTITION */ * FROM t;
SELECT /*+ REPARTITION_BY_RANGE(c) */ * FROM t;
SELECT /*+ REPARTITION_BY_RANGE(3, c) */ * FROM t;
SELECT /*+ REBALANCE */ * FROM t;REBALANCE is different. It is meant for queries whose results you write to a table, so that every output partition is a reasonable size: neither too small nor too big. It tries to partition by any columns you give it. If some partitions are skewed, Spark splits them. REBALANCE depends on AQE and is ignored when AQE is disabled.
Checkpoint 5 of 5· Check yourself
A query uses /*+ REPARTITION(100), COALESCE(500), REPARTITION_BY_RANGE(3, c) */. Which partitioning ends up in the physical plan?
Each hint adds a node to the logical plan, but the optimizer keeps only the leftmost one. The physical plan in the docs shows Exchange RoundRobinPartitioning(100).
“When multiple partitioning hints are specified, multiple nodes are inserted into the logical plan, but the leftmost hint is picked by the optimizer.”Source: docs.databricks.com
Sources5
Exam traps
Each one states something that sounds right. Open it to see what is actually true.
1.coalesce(n) with a larger n increases the partition count, just without a shuffle.Why is that wrong?
coalesce() only merges partitions. If you ask for more than exist, the DataFrame keeps its current count. Only repartition() can increase it.
2.Because coalesce() avoids a shuffle, coalesce(1) is always faster than repartition(1).Why is that wrong?
Without a shuffle, the upstream work runs inside the reduced partitions and may land on a single node. repartition(1) adds a shuffle but keeps the upstream stages running in parallel.
3.The REBALANCE hint works the same way whether or not AQE is turned on.Why is that wrong?
REBALANCE relies on AQE to split skewed partitions and size the output, and Spark ignores the hint when AQE is disabled.
Covered in The same controls in SQL: partitioning hints
Sources
Every claim above is drawn from one of these pages, quoted as it was written on the date shown.
- 1.https://docs.databricks.com/aws/en/optimizations/aqeOfficial docs
“The default number of partitions to use when shuffling data for joins or aggregations.”
↩︎ Where partition counts come from“Setting the value auto enables auto-optimized shuffle, which automatically determines this number based on the query plan and the query input data size.”
↩︎ Where partition counts come from - 2.https://spark.apache.org/docs/latest/sql-performance-tuning.htmlSecondary source
“The maximum number of bytes to pack into a single partition when reading files.”
↩︎ Where partition counts come from“spark.sql.files.maxPartitionBytes | 134217728 (128 MB) | The maximum number of bytes to pack into a single partition when reading files.”
↩︎ Checkpoint“spark.sql.shuffle.partitions | 200 | Configures the number of partitions to use when shuffling data for joins or aggregations.”
↩︎ Prediction - 3.
“can be an int to specify the target number of partitions or a Column.”
↩︎ Where partition counts come from“Returns a new DataFrame partitioned by the given partitioning expressions. The resulting DataFrame is hash partitioned.”
↩︎ repartition(): any partition count, at the cost of a shuffle - 4.
“this operation results in a narrow dependency, e.g. if you go from 1000 partitions to 100 partitions, there will not be a shuffle”
↩︎ coalesce(): fewer partitions without a shuffle“This will add a shuffle step, but means the current upstream partitions will be executed in parallel”
↩︎ Key concept“If a larger number of partitions is requested, it will stay at the current number of partitions.”
↩︎ Exam trap 1“this may result in your computation taking place on fewer nodes than you like”
↩︎ Exam trap 2“If a larger number of partitions is requested, it will stay at the current number of partitions.”
↩︎ Checkpoint - 5.
“COALESCE, REPARTITION, and REPARTITION_BY_RANGE hints are supported and are equivalent to coalesce, repartition, and repartitionByRange Dataset APIs, respectively.”
↩︎ The same controls in SQL: partitioning hints“if there are skews, Spark will split the skewed partitions, to make these partitions not too big.”
↩︎ The same controls in SQL: partitioning hints“This hint is ignored if AQE is not enabled.”
↩︎ Exam trap 3“When multiple partitioning hints are specified, multiple nodes are inserted into the logical plan, but the leftmost hint is picked by the optimizer.”
↩︎ Checkpoint