ML Engineer MasterClass (October) | 4 seats left

Spark · Narrow and wide transformations
AmazonAmazon Analytics
Amazon · How Spark Executes

Narrow and wide transformations

Identify local row operations and cross-partition exchanges in a transformation pipeline.

Step 1 of 6 · Learn

Calculate within each partition

A projection such as select with total + 10 can be computed from each input row locally. It does not require other partitions to send rows. This is a narrow operation. The final orderBy used for a repeatable teaching result is a separate global sort and can require a shuffle, so the entire example is not shuffle-free.

Lesson reference: PySpark and Scala

Calculate within each partition

A projection such as select with total + 10 can be computed from each input row locally. It does not require other partitions to send rows. This is a narrow operation. The final orderBy used for a repeatable teaching result is a separate global sort and can require a shuffle, so the entire example is not shuffle-free.

PySpark example

from pyspark.sql import Window
from pyspark.sql.functions import col

result = orders.select(
    col("order_id"), (col("total") + 10).alias("adjusted_total")
).orderBy("order_id")
result.show(truncate=False)

Scala example

import org.apache.spark.sql.expressions.Window
import org.apache.spark.sql.functions._
import spark.implicits._

val result = orders.select(
    col("order_id"), (col("total") + 10).alias("adjusted_total")
).orderBy("order_id")
result.show(false)

Bring matching keys together

Grouping rows by status usually requires partial results from different partitions to meet. Spark can aggregate locally first, exchange by key, and then combine the partial sums. This is a wide operation when redistribution is required. Existing partitioning and adaptive execution can change the physical plan; a groupBy call is not proof of a fixed task count.

PySpark example

from pyspark.sql import Window
from pyspark.sql.functions import col, sum

result = orders.groupBy("status").agg(
    sum("total").alias("total_value")
).select(col("status").cast("string").alias("group_key"), col("total_value")).orderBy("group_key")
result.show(truncate=False)

Scala example

import org.apache.spark.sql.expressions.Window
import org.apache.spark.sql.functions._
import spark.implicits._

val result = orders.groupBy("status").agg(
    sum("total").alias("total_value")
).select(col("status").cast("string").alias("group_key"), col("total_value")).orderBy("group_key")
result.show(false)

Reduce rows before an exchange

A row-level filter can run before grouped aggregation, reducing the input to partial aggregation and the exchange. Here the total threshold applies to individual orders, not to the final group sum. Spark may combine compatible narrow operations in one stage. The diagram describes the data dependency, not an exact stage count for every execution.

PySpark example

from pyspark.sql import Window
from pyspark.sql.functions import col, sum

result = orders.filter(col("total") >= 100).groupBy("status").agg(
    sum("total").alias("total_value")
).select(col("status").alias("group_key"), col("total_value")).orderBy("group_key")
result.show(truncate=False)

Scala example

import org.apache.spark.sql.expressions.Window
import org.apache.spark.sql.functions._
import spark.implicits._

val result = orders.filter(col("total") >= 100).groupBy("status").agg(
    sum("total").alias("total_value")
).select(col("status").alias("group_key"), col("total_value")).orderBy("group_key")
result.show(false)
example.pyPySpark
1. Input partitionRead the rows already assigned to this partition.
2. Local expressionSelect, calculate, or filter without needing other partitions.
3. Display sortThe final global orderBy is separate and may shuffle rows.
Source orders6 rows
order_idcustomer_idstatustotalitem_count
1001101Delivered89.52
1002102Shipped1493
1003101Cancelled351
1004103Delivered219.994
1005104Delivered49.991
1006105Shipped1202
Calculate within each partition
Result6 rows
order_idadjusted_total
100199.5
1002159
100345
1004229.99
100559.99
1006130
Each adjusted value needs only that row’s total. Display ordering is handled afterward.

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.