ML Engineer MasterClass (October) | 4 seats left

Spark · Partition and order
AmazonAmazon Analytics
Amazon · Window Functions

Partition and order

Define which rows a window sees and the sequence used within each group.

Step 1 of 6 · Learn

Choose window groups

A window adds a value to each original row rather than collapsing rows like groupBy. partitionBy(status) counts all orders of the same status. With no window ordering, count covers the whole partition. These logical groups are not a request to set Spark physical partition counts.

Lesson reference: PySpark and Scala

Choose window groups

A window adds a value to each original row rather than collapsing rows like groupBy. partitionBy(status) counts all orders of the same status. With no window ordering, count covers the whole partition. These logical groups are not a request to set Spark physical partition counts.

PySpark example

from pyspark.sql import Window
from pyspark.sql.functions import (
    col, lit, array, when, explode, explode_outer, posexplode,
    create_map, to_json, struct, from_json, get_json_object, count, row_number
)


result = orders.withColumn(
    "partition_count", count("*").over(Window.partitionBy("status"))
).select("order_id", "customer_id", "status", "partition_count").orderBy("order_id")
result.show(truncate=False)

Scala example

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


val result = orders.withColumn(
    "partition_count", count("*").over(Window.partitionBy("status"))
).select("order_id", "customer_id", "status", "partition_count").orderBy("order_id")
result.show(false)

Order within partitions

Window orderBy defines the sequence used by row_number, restarting at 1 for each status. The example uses total ascending and order_id to break ties. A separate final orderBy controls how the result is displayed; the window order alone does not guarantee display order.

PySpark example

from pyspark.sql import Window
from pyspark.sql.functions import (
    col, lit, array, when, explode, explode_outer, posexplode,
    create_map, to_json, struct, from_json, get_json_object, count, row_number
)


result = orders.withColumn(
    "row_position", row_number().over(
        Window.partitionBy("status").orderBy(col("total").asc(), col("order_id").asc())
    )
).select("order_id", "status", "total", "row_position").orderBy("status", "row_position")
result.show(truncate=False)

Scala example

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


val result = orders.withColumn(
    "row_position", row_number().over(
        Window.partitionBy("status").orderBy(col("total").asc(), col("order_id").asc())
    )
).select("order_id", "status", "total", "row_position").orderBy("status", "row_position")
result.show(false)

Use multiple partition keys

A partition is defined by the whole key combination. Customer 101 has two orders, but they have different statuses. Adding status to the customer partition separates them. count still covers each entire partition because the window has no ordering.

PySpark example

from pyspark.sql import Window
from pyspark.sql.functions import (
    col, lit, array, when, explode, explode_outer, posexplode,
    create_map, to_json, struct, from_json, get_json_object, count, row_number
)


result = orders.withColumn(
    "partition_count", count("*").over(Window.partitionBy("customer_id"))
).select("order_id", "customer_id", "status", "partition_count").orderBy("order_id")
result.show(truncate=False)

Scala example

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


val result = orders.withColumn(
    "partition_count", count("*").over(Window.partitionBy("customer_id"))
).select("order_id", "customer_id", "status", "partition_count").orderBy("order_id")
result.show(false)
example.pyPySpark
Original rowsLogical window groupsOne value per original row
status: DeliveredOrder IDs in this partition1001, 1004, 1005
status: ShippedOrder IDs in this partition1002, 1006
status: CancelledOrder IDs in this partition1003
Source orders6 rows
order_idcustomer_idstatustotalitem_count
1001101Delivered89.52
1002102Shipped1493
1003101Cancelled351
1004103Delivered219.994
1005104Delivered49.991
1006105Shipped1202
Choose window groups
Result6 rows
order_idcustomer_idstatuspartition_count
1001101Delivered3
1002102Shipped2
1003101Cancelled1
1004103Delivered3
1005104Delivered3
1006105Shipped2
Every row remains. The three Delivered orders each receive 3; the two Shipped orders each receive 2.

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.