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)