Lesson 8 / 25

Window Functions

Rank rows within groups and compute running totals without collapsing the data.

Aggregate without losing rows

A window function computes a value for each row using a "window" of related rows, defined by Window.partitionBy(...) (the group), orderBy(...) (the order inside it) and optionally a frame. Unlike groupBy, it keeps every row. Typical uses: row_number() / rank() / dense_rank() for "top N per group", lag() / lead() to compare with the previous or next row, and sum().over(...) for running totals. Windows with a partitionBy shuffle data by the partition key, and a window with no partitionBy pulls everything into a single partition, which does not scale.

Top order per customer, run

I ran this on Apache Spark 4.0.0 (PySpark, local mode, in the official Docker image). desc_nulls_last puts the null amount last, so Kiran's only order is still rank 1 with a null amount.

w = Window.partitionBy("customer").orderBy(F.col("amount").desc_nulls_last())
orders.withColumn("rank", F.row_number().over(w)).filter("rank = 1").select("customer", "amount").orderBy("customer").show()

Output:

+--------+------+
|customer|amount|
+--------+------+
|    asha| 200.0|
|   kiran|  NULL|
|   meera|  50.0|
|    ravi| 300.0|
+--------+------+

Running total, run

I ran this on Apache Spark 4.0.0 (PySpark, local mode, in the official Docker image). Each row shows the cumulative amount so far in date and order order; the null was filled with 0.0 first, so the last two rows stay at 750.0.

w2 = Window.orderBy("order_date", "order_id")
orders.na.fill({"amount":0.0}).withColumn("running", F.sum("amount").over(w2)).select("order_id", "running").orderBy("order_id").show()

Output:

+--------+-------+
|order_id|running|
+--------+-------+
|       1|  120.0|
|       2|  200.0|
|       3|  400.0|
|       4|  450.0|
|       5|  750.0|
|       6|  750.0|
+--------+-------+

Quick check: What is the key difference between groupBy and a window function?

  • There is no difference
  • groupBy keeps every row, windows collapse them
  • A window function keeps every input row
  • Windows only work on strings
Answer

A window function keeps every input row — groupBy collapses each group to one row; window functions add a computed column to all rows.