ML Engineer MasterClass (October) | 4 seats left

Spark · Repartition and coalesce
AmazonAmazon Analytics
Amazon · How Spark Executes

Repartition and coalesce

Choose how to change partition counts without confusing fewer tasks with faster execution.

Step 1 of 6 · Learn

Reduce without a new shuffle

coalesce can reduce partitions through a narrow dependency, combining existing partitions without a new redistribution shuffle. The example first requests four partitions, then reduces to two. A drastic reduction can reduce parallelism and create uneven work. The initial repartition still shuffles; coalesce does not make the whole pipeline shuffle-free.

Lesson reference: PySpark and Scala

Reduce without a new shuffle

coalesce can reduce partitions through a narrow dependency, combining existing partitions without a new redistribution shuffle. The example first requests four partitions, then reduces to two. A drastic reduction can reduce parallelism and create uneven work. The initial repartition still shuffles; coalesce does not make the whole pipeline shuffle-free.

PySpark example

from pyspark.sql.functions import col, sum, count, lit

result = spark.range(1).select(
    lit(orders.repartition(4).coalesce(2).rdd.getNumPartitions()).alias("partition_count")
)
result.show(truncate=False)

Scala example

import org.apache.spark.sql.functions._

val result = spark.range(1).select(
    lit(orders.repartition(4).coalesce(2).rdd.getNumPartitions).alias("partition_count")
)
result.show(false)

Increase with repartition

coalesce does not increase the number of partitions: asking for four after two leaves two. repartition can increase or decrease the count by redistributing rows. More partitions can increase overhead for tiny datasets; the goal is suitable work sizes, not the largest possible count. Optimizers may simplify consecutive repartitions.

PySpark example

from pyspark.sql.functions import col, sum, count, lit

result = spark.range(1).select(
    lit(orders.repartition(2).coalesce(4).rdd.getNumPartitions()).alias("partition_count")
)
result.show(truncate=False)

Scala example

import org.apache.spark.sql.functions._

val result = spark.range(1).select(
    lit(orders.repartition(2).coalesce(4).rdd.getNumPartitions).alias("partition_count")
)
result.show(false)

Verify data is preserved

Changing physical partitioning should not change order counts or sums. This example checks row_count and total_value and separately reports the requested orders partition count. The aggregate result itself can have a different partition layout. Avoid assuming output-file count, task attempts, or worker count equals this measurement.

PySpark example

from pyspark.sql.functions import col, sum, count, lit

result = orders.repartition(2).agg(
    count("*").alias("row_count"), sum("total").alias("total_value")
).withColumn("partition_count", lit(orders.repartition(2).rdd.getNumPartitions())).select("partition_count", "row_count", "total_value")
result.show(truncate=False)

Scala example

import org.apache.spark.sql.functions._

val result = orders.repartition(2).agg(
    count("*").alias("row_count"), sum("total").alias("total_value")
).withColumn("partition_count", lit(orders.repartition(2).rdd.getNumPartitions)).select("partition_count", "row_count", "total_value")
result.show(false)
example.pyPySpark
1. Initial repartitionRequest four partitions using an exchange.
2. Coalesce to twoCombine existing partitions without a new redistribution shuffle.
3. Trade-offFewer task slots can be used downstream; balance is not guaranteed.
Source orders6 rows
order_idcustomer_idstatustotalitem_count
1001101Delivered89.52
1002102Shipped1493
1003101Cancelled351
1004103Delivered219.994
1005104Delivered49.991
1006105Shipped1202
Reduce without a new shuffle
Result1 rows
partition_count
2
The example ends with two partitions. It makes no promise that each holds three rows.

The example is loaded in the editor. Run it as written, then try a small change.

Runs on the Spark backend. First startup may take a moment.

Run your code to see the result.