CertSafari
    Databricks Certified Associate Developer for Apache Spark· Lessons

    Domain 4 · Lesson 22/32

    Spark repartition vs coalesce: controlling partitions and shuffles

    Implement performance tuning strategies & optimize cluster utilization, including partitioning, repartitioning, coalescing, identifying data skew, and reducing shuffling

    9 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

    • 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?

    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.

    Sources123

    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()

    Sources3

    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.

    coalesce(1) merges the three partitions of a range into onepython
    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.

    repartition() and coalesce() compared
    Aspectrepartition()coalesce()
    ArgumentsnumPartitions and/or partitioning columnsnumPartitions only
    ShuffleYes, adds a shuffle stepNo, narrow dependency
    Can increase partitionsYesNo, stays at the current number
    Upstream parallelism on a drastic reductionPreservedMay collapse to fewer nodes

    Checkpoint 3 of 5· Check yourself

    A DataFrame has 1000 partitions and you call df.coalesce(2000). What happens?

    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?

    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.

    Partitioning hint forms in Spark SQLsql
    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?

    Sources5

    Exam traps

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

    1. 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.

      Covered in coalesce(): fewer partitions without a shuffle

    2. 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.

      Covered in coalesce(): fewer partitions without a shuffle

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

    Continue to page 2 of 2

    Spark data skew and shuffle tuning with Adaptive Query Execution

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