Skip to content
LakeBench
ProblemsCommunityPricing
Sign inStart practicing
Back
  1. Home
  2. Interview prep
  3. Window functions in PySpark

PySpark · DataFrame API in Practice

Window functions in PySpark

Mediumpyspark-72
window-functionsshufflerow-numberlag

Question

How do you write window functions in PySpark, and what triggers a shuffle?

Solution

In PySpark you build a window specification with Window.partitionBy(...).orderBy(...), then apply a function over it with .over(spec). Every distinct partitionBy needs a shuffle and a sort, so windows are one of the more expensive operations.

Writing one

from pyspark.sql import Window
from pyspark.sql import functions as F

w = Window.partitionBy("customer_id").orderBy(F.col("order_ts").desc())

latest = (orders
    .withColumn("rn", F.row_number().over(w))
    .filter("rn = 1"))

w_run = (Window.partitionBy("customer_id").orderBy("order_ts")
         .rowsBetween(Window.unboundedPreceding, Window.currentRow))

running = orders.withColumn("running_total", F.sum("amount").over(w_run))
prev = orders.withColumn("prev_amount", F.lag("amount", 1).over(
    Window.partitionBy("customer_id").orderBy("order_ts")))

The same frame rules as SQL apply. With an orderBy and no explicit frame, the default is a range frame from the start to the current row, so rows with equal ordering values count as one group. Write rowsBetween when you want a row-by-row running total.

What triggers the shuffle

Spark must bring all rows of one partition key to the same task, then sort them by the order key. That is an exchange (shuffle) followed by a sort. You can see it in explain(): Exchange hashpartitioning(customer_id) then Sort, then Window.

  • If several window functions use the same partitionBy and orderBy, Spark can compute them in one pass with one shuffle.
  • If they use different partition keys, each distinct key needs its own shuffle. Five windows with five different keys means five shuffles.
  • A window with no partitionBy moves all rows into one partition. Spark warns about it, and on large data it is a guaranteed slow single task or an out-of-memory error.

Skew

If one customer_id has 100 million rows, that whole group must be sorted inside a single task. A skewed window key is a very common reason for a window stage with one slow task. Options: pick a finer partition key if the logic allows it, filter junk keys such as NULL before the window, or restructure the logic to avoid a full-group window.

Alternatives

For "latest row per key", an aggregation with max(struct(order_ts, ...)) can be cheaper than row_number, because it uses a partial aggregation before the shuffle. It is worth knowing if the window version is too slow.

🎯 Put this concept into practice

Solidify this answer with real hands-on interview drills in the browser studio.

Open related drill →
PreviousNext