CertSafari
    Databricks Certified Associate Developer for Apache Spark· Lessons

    Domain 1 · Lesson 5/32

    Spark Partitions, Shuffles and Shuffle Partition Settings

    Configure Spark partitioning in distributed data processing, including shuffles and partitions

    11 min read
    3.12% of exam
    5 sources
    Published 3 Oct 2026
    Docs as of 30 Sep 2026

    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.

    Settings that control how file input is split into partitions (file-based sources only)
    PropertyDefaultWhat it controls
    spark.sql.files.maxPartitionBytes134217728 (128 MB)The maximum number of bytes packed into a single partition when reading files
    spark.sql.files.openCostInBytes4194304 (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.minPartitionNumDefault ParallelismA suggested (not guaranteed) minimum number of split file partitions
    spark.sql.files.maxPartitionNumNoneA 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?

    Sources123

    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?

    Sources34

    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.

    spark.sql.shuffle.partitions as documented by each platform
    PlatformDocumented defaultDescription given
    Apache Spark (performance tuning guide)200Configures the number of partitions to use when shuffling data for joins or aggregations
    Databricks (Spark configuration page)autoThe 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.

    Setting a Spark property on the current SparkSession from Python (Databricks example)python
    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.

    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?

    Sources21

    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:

    AQE partition-coalescing properties
    PropertyDefaultMeaning
    spark.sql.adaptive.coalescePartitions.enabledtrueTurns partition coalescing on or off
    spark.sql.adaptive.advisoryPartitionSizeInBytes64MBTarget size after coalescing. Coalesced partitions will be close to this size but no bigger
    spark.sql.adaptive.coalescePartitions.minPartitionSize1MBCoalesced partitions will be no smaller than this size
    spark.sql.adaptive.coalescePartitions.minPartitionNum2x no. of cluster coresMinimum 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?

    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?

    Sources54

    Exam traps

    Each one states something that sounds right. Open it to see what is actually true.

    1. 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. 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. 1.
      “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. 2.
      “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. 3.
      “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. 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. 5.
      “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

    Continue to page 2 of 2

    Repartition vs Coalesce: Controlling DataFrame Partitions

    Spotted a mistake, or was something unclear? Tell us.