What you will be able to do
- Explain why the number of partitions sets how much work Spark can do in parallel
- Name the file-read settings that decide how input is split into partitions
- Recognise which operations cause a shuffle, and why a shuffle costs so much
- Set spark.sql.shuffle.partitions for a session and know how its default differs between open-source Spark and Databricks
- Describe how adaptive query execution (AQE) coalesces post-shuffle partitions and detects skewed ones
Key concept
Partition as the unit of parallel work — Spark splits a dataset into partitions, and each task works on exactly one of them. So the number and size of partitions decide how much of the cluster a stage can use, and a shuffle is the costly step that rearranges rows into a new set of partitions.
1.Where partitions come from: splitting input files
A DataFrame is never processed as one block. Spark divides it into partitions, and during computation a single task works on a single partition. So partitioning is the setting that governs parallelism. With too few partitions, cores sit idle. With too many tiny ones, each task does almost no useful work.
The first partitioning decision happens when Spark reads files. For file-based sources such as Parquet, JSON and ORC, Spark packs bytes into partitions up to a maximum size. Both open-source Spark and Databricks document a default of 134217728 bytes (128 MB) for spark.sql.files.maxPartitionBytes.
| Property | Default | What it controls |
|---|---|---|
| spark.sql.files.maxPartitionBytes | 134217728 (128 MB) | The maximum number of bytes packed into a single partition when reading files |
| spark.sql.files.openCostInBytes | 4194304 (4 MB) | The estimated cost of opening a file, counted as bytes that could be scanned in the same time. Used when putting several files into one partition |
| spark.sql.files.minPartitionNum | Default Parallelism | A suggested (not guaranteed) minimum number of split file partitions |
| spark.sql.files.maxPartitionNum | None | A suggested (not guaranteed) maximum. If set, Spark rescales partitions so the count is close to this value |
The answer to the prediction is the open cost. Spark charges every file an estimated 4 MB opening cost, so many small files do not all end up in one partition just because their combined bytes are small. The documentation advises that over-estimating this cost is better. Note the wording on the two partition-count settings: both are *suggested (not guaranteed)*. Spark treats them as targets, not exact numbers.
Checkpoint 1 of 6· Check yourself
A job reads Parquet files and you want each input partition to hold less data than it does now. Which setting do you lower?
maxPartitionBytes caps how many bytes go into one partition when Spark reads files. shuffle.partitions and the AQE advisory size only affect partitions produced after a shuffle.
“The maximum number of bytes to pack into a single partition when reading files.”Source: docs.databricks.com
2.What a shuffle is and what triggers one
Some operations need rows that live in different partitions to end up together. An aggregation per key is the standard example: the values for one key may be on many partitions, or even on many machines, but they have to be in one place before the result can be computed. The Spark RDD guide explains that certain operations trigger an event called the shuffle, which redistributes data so that it is grouped differently across partitions. Spark has to read from all partitions and send matching rows to the same destination. That is an all-to-all operation, and the guide calls the shuffle complex and costly because it usually copies data across executors and machines.
In DataFrame and SQL work, the shuffle appears in the plan as an *exchange*. Databricks' adaptive query execution documentation says a query contains an exchange usually when there is a join, an aggregate or a window. Shuffles also put pressure on memory. Databricks describes execution memory as the memory used for shuffles, joins, sorts and aggregations. When that memory runs low, Spark spills data to disk, and spill is most common during shuffling.
Not every operation that changes partitioning shuffles. Narrow operations work on each partition where it already sits. The repartition and coalesce lesson covers an important case: reducing partitions with coalesce, which the API reference describes as a narrow dependency with no shuffle.
Checkpoint 2 of 6· Check yourself
Which statement best describes a shuffle?
A shuffle moves rows between partitions, and usually between machines, so that related rows end up together. That movement is what makes it expensive.
“This typically involves copying data across executors and machines, making the shuffle a complex and costly operation.”Source: spark.apache.org
3.Configuring spark.sql.shuffle.partitions
File-read settings decide the partitions at the start of a query. After a shuffle, a separate setting takes over: spark.sql.shuffle.partitions, which sets the number of partitions used when data is shuffled for joins or aggregations. Every post-shuffle stage therefore gets that many tasks unless something else changes the number. The default depends on where you run Spark, and exam questions test this difference.
| Platform | Documented default | Description given |
|---|---|---|
| Apache Spark (performance tuning guide) | 200 | Configures the number of partitions to use when shuffling data for joins or aggregations |
| Databricks (Spark configuration page) | auto | The default number of partitions to use when shuffling data for joins or aggregations |
To change a property for the current notebook only, set it on the session. Databricks shows the syntax for a property in SQL as SET <property> = <value> and in Python as follows. The example uses an ANSI setting, and the same call works for any session-settable key.
spark.conf.set("spark.sql.ansi.enabled", "true")Scope matters. A value set in a notebook applies only to that notebook's SparkSession. A value set in the compute configuration applies to every notebook and job on that compute. Databricks also generally recommends against configuring most Spark properties, because legacy settings carried over from open-source Spark can override newer default behaviour that is tuned for the platform.
Checkpoint 3 of 6· Match them up
Match each way of setting a Spark property to its scope or effect
Tap a term, then the definition that fits it.
Databricks scopes a notebook-level conf to that notebook's session and a compute-level conf to everything on the compute. It warns that legacy confs can override optimised defaults.
“Within a notebook | Only the SparkSession for the current notebook.”Source: docs.databricks.com
Checkpoint 4 of 6· Exam question
A Spark job performs a wide join between two large tables where one join key accounts for roughly 40% of all rows, so one shuffle partition holds far more data than the rest and its task runs much longer than the others. Which approach directly targets this specific partition imbalance?
Correct answer: A — Enable AQE skew join handling by leaving `spark.sql.adaptive.skewJoin.enabled` set to `true`, letting Spark split the oversized shuffle partition into smaller tasks at runtime.
- A. AQE's skew join optimization detects a partition that is disproportionately larger than the median and splits it (replicating the matching side as needed) into several roughly evenly sized tasks. This is the mechanism built specifically to fix a single oversized shuffle partition without touching the rest of the job.
- B. Raising the total shuffle partition count spreads out the non-skewed keys further but does nothing for a key that already dominates the data — that key's rows still land on one logical partition regardless of how many total partitions exist. It does not address the root cause of the imbalance.
- C. Coalescing after the shuffle merges partitions together locally, which would combine the already-oversized partition with others and make the imbalance worse, not better. Coalesce reduces partition count; it has no mechanism to split a large partition apart.
- D. Changing the storage level affects how a DataFrame is cached for reuse, not how work is distributed across a shuffle. Spilling to disk might avoid an out-of-memory crash but the single task still has to process all of the skewed key's rows, so the runtime imbalance remains.
4.Adaptive query execution: resizing partitions after a shuffle
A fixed shuffle partition count is a guess made before the data is seen. Adaptive query execution (AQE) re-optimizes the query while it runs. At the end of each shuffle it has accurate statistics, and it uses them to pick a better post-shuffle partition size and number. AQE is enabled by default. It applies to non-streaming queries that contain at least one exchange or subquery. Two of its four main features concern partitions:
- Partition coalescing. AQE combines small partitions into reasonably sized ones after a shuffle exchange. - Skew handling. In sort merge joins and shuffle hash joins, AQE splits skewed tasks, and replicates them if needed, into tasks of roughly equal size.
The coalescing behaviour is controlled by these properties:
| Property | Default | Meaning |
|---|---|---|
| spark.sql.adaptive.coalescePartitions.enabled | true | Turns partition coalescing on or off |
| spark.sql.adaptive.advisoryPartitionSizeInBytes | 64MB | Target size after coalescing. Coalesced partitions will be close to this size but no bigger |
| spark.sql.adaptive.coalescePartitions.minPartitionSize | 1MB | Coalesced partitions will be no smaller than this size |
| spark.sql.adaptive.coalescePartitions.minPartitionNum | 2x no. of cluster cores | Minimum number of partitions after coalescing. Not recommended, because setting it overrides minPartitionSize |
Skew is the opposite problem: one partition is far larger than the others. AQE treats a partition as skewed only when two conditions are both true. Its size must be more than spark.sql.adaptive.skewJoin.skewedPartitionFactor (default 5) times the median partition size, and it must also be larger than spark.sql.adaptive.skewJoin.skewedPartitionThresholdInBytes (default 256MB). You can also spot skew in the Spark UI. Databricks suggests comparing a stage's Max task duration with its 75th percentile. If Max is 50% higher than the 75th percentile, the stage may be suffering from skew.
Checkpoint 5 of 6· Check yourself
With default AQE settings, the median partition after a shuffle is 40 MB. One partition is 300 MB. Is it treated as skewed?
300 MB is more than 5 × the 40 MB median (200 MB) and more than the 256MB threshold. Both conditions hold, so the partition is skewed. Skew join handling is on by default.
“A partition is considered skewed when both (partition size > skewedPartitionFactor * median partition size) and (partition size > skewedPartitionThresholdInBytes) are true.”Source: docs.databricks.com
Checkpoint 6 of 6· Exam question
A DataFrame currently has 400 partitions after several transformations. An engineer needs to write it out as a single file for a downstream system that expects one file per day, and wants to avoid an expensive full shuffle if possible. Which call accomplishes this?
Correct answer: A — `df.coalesce(1)` reduces the partition count by merging existing partitions locally, avoiding a full shuffle while still producing one output file.
- A. Coalesce reduces the number of partitions by combining existing partitions on the same executor wherever possible, so it can produce a single output partition without redistributing every row across the network. This makes it the lowest-cost way to reach one output file when reducing partition count.
- B. This call does produce exactly one partition and one file, but it forces every row to move across the cluster in a full shuffle, which is unnecessary network and CPU cost when the goal is simply fewer output files.
- C. Hash-partitioning by a column into a single partition still requires Spark to shuffle every row to compute the target partition, so this also pays the full shuffle cost the scenario is trying to avoid, and the column argument is pointless when there is only one partition to hash into.
- D. Partitioning the writer output by a column creates one subdirectory and file set per distinct value of that column, so it produces many files rather than the single combined file the downstream system expects.
Exam traps
Each one states something that sounds right. Open it to see what is actually true.
1.spark.sql.shuffle.partitions is always 200, so on Databricks you should set it explicitly to get sensible shuffle parallelism.Why is that wrong?
200 is the default documented for open-source Apache Spark. The Databricks configuration page lists the default as auto, and Databricks generally recommends against overriding most Spark properties.
Covered in Configuring spark.sql.shuffle.partitions
2.To keep AQE from coalescing too aggressively, set spark.sql.adaptive.coalescePartitions.minPartitionNum.Why is that wrong?
Databricks marks minPartitionNum as not recommended, because setting it overrides the size-based minPartitionSize control.
Covered in Adaptive query execution: resizing partitions after a shuffle
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/spark/confOfficial docs
“The maximum number of bytes to pack into a single partition when reading files.”
↩︎ Where partitions come from: splitting input files“Databricks generally recommends against configuring most Spark properties.”
↩︎ Configuring spark.sql.shuffle.partitions“spark.sql.shuffle.partitions | auto | The default number of partitions to use when shuffling data for joins or aggregations.”
↩︎ Exam trap 1“Within a notebook | Only the SparkSession for the current notebook.”
↩︎ Checkpoint - 2.https://spark.apache.org/docs/latest/sql-performance-tuning.htmlSecondary source
“The suggested (not guaranteed) minimum number of split file partitions.”
↩︎ Where partitions come from: splitting input files“Configures the number of partitions to use when shuffling data for joins or aggregations.”
↩︎ Configuring spark.sql.shuffle.partitions - 3.https://spark.apache.org/docs/latest/rdd-programming-guide.htmlSecondary source
“During computations, a single task will operate on a single partition”
↩︎ Where partitions come from: splitting input files“Certain operations within Spark trigger an event known as the shuffle.”
↩︎ What a shuffle is and what triggers one“During computations, a single task will operate on a single partition”
↩︎ Key concept“This typically involves copying data across executors and machines, making the shuffle a complex and costly operation.”
↩︎ Checkpoint - 4.
“This memory is used for computations such as shuffles, joins, sorts, and aggregations.”
↩︎ What a shuffle is and what triggers one“It is most common during data shuffling.”
↩︎ What a shuffle is and what triggers one“If the Max duration is 50% more than the 75th percentile, you may be suffering from skew.”
↩︎ Adaptive query execution: resizing partitions after a shuffle - 5.https://docs.databricks.com/aws/en/optimizations/aqeOfficial docs
“Dynamically coalesces partitions (combine small partitions into reasonably sized partitions) after shuffle exchange.”
↩︎ Adaptive query execution: resizing partitions after a shuffle“The coalesced partition sizes will be close to but no bigger than this target size.”
↩︎ Adaptive query execution: resizing partitions after a shuffle“Not recommended, because setting explicitly overrides spark.sql.adaptive.coalescePartitions.minPartitionSize.”
↩︎ Exam trap 2“Very small tasks have worse I/O throughput and tend to suffer more from scheduling overhead and task setup overhead.”
↩︎ Prediction“A partition is considered skewed when both (partition size > skewedPartitionFactor * median partition size) and (partition size > skewedPartitionThresholdInBytes) are true.”
↩︎ Checkpoint