Menu

PySpark course · Lesson 7 of 8

PySpark Window Functions: Ranking, Lag and Running Totals

Use PySpark window functions for top-N per group, previous-row comparisons and running totals, and avoid the single-partition trap of windows without partitionBy.

  • Intermediate
  • 3 min read
  • Updated Oct 2026
On this page
  1. The setup
  2. Top row per group
  3. Previous and next rows: lag and lead
  4. Running totals and explicit frames
  5. The same thing in Spark SQL
  6. Performance notes
  7. Common mistakes
  8. Interview relevance
  9. Key takeaway

PySpark window functions work like their SQL counterparts: they compute a value for each row from a group of related rows without collapsing the rows. If you have not used them in SQL, read SQL window functions first; the ideas are identical.

The setup

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

spark = SparkSession.builder.master("local[2]").appName("windows").getOrCreate()

sales = spark.createDataFrame(
    [("Asha","north","2026-01-01",100), ("Asha","north","2026-01-02",150),
     ("Ben","north","2026-01-01",150),  ("Ben","north","2026-01-03",90),
     ("Chen","south","2026-01-01",200), ("Chen","south","2026-01-02",50)],
    ["rep", "region", "sale_date", "amount"],
)

A window is described by a Window specification, then applied with .over(...).

Top row per group

w = Window.partitionBy("region").orderBy(F.desc("amount"), F.asc("rep"))

top = (
    sales.withColumn("rn", F.row_number().over(w))
         .filter("rn = 1")
         .drop("rn")
)
top.show()
+----+------+----------+------+
| rep|region| sale_date|amount|
+----+------+----------+------+
|Asha| north|2026-01-02|   150|
|Chen| south|2026-01-01|   200|
+----+------+----------+------+

Two reps in the north tie at 150. Adding F.asc("rep") as a tiebreaker makes the result deterministic. Without it, which tied row gets rn = 1 can change between runs.

The same rule as SQL applies: use row_number for exactly one row per group, rank to allow ties with gaps, and dense_rank for ties without gaps.

Previous and next rows: lag and lead

by_rep = Window.partitionBy("rep").orderBy("sale_date")

changes = (
    sales.withColumn("previous", F.lag("amount").over(by_rep))
         .withColumn("change", F.col("amount") - F.col("previous"))
)
changes.orderBy("rep", "sale_date").show()

The first row of each partition has no previous row, so previous and change are NULL. Pass a default as the third argument to lag if you prefer a value.

Running totals and explicit frames

running = (
    Window.partitionBy("rep")
          .orderBy("sale_date")
          .rowsBetween(Window.unboundedPreceding, Window.currentRow)
)
sales.withColumn("running_total", F.sum("amount").over(running)).show()

If you give an aggregate window an orderBy and no frame, Spark uses a range-based frame from the start of the partition to the current row, which treats rows with equal ordering values as one unit. Specify rowsBetween when you want a row-by-row running total.

The same thing in Spark SQL

sales.createOrReplaceTempView("sales")
spark.sql("""
    SELECT rep, region, amount,
           ROW_NUMBER() OVER (PARTITION BY region ORDER BY amount DESC, rep) AS rn
    FROM sales
""").show()

DataFrame code and Spark SQL go through the same optimizer and produce equivalent plans, so choose whichever your team finds more readable.

Performance notes

  • A window with partitionBy triggers a shuffle so that all rows of a partition meet on one task, then a sort inside each partition. See partitions, shuffles and skew.
  • If a few partition keys hold most of the rows, those tasks run far longer than the rest. That is data skew.
  • Reusing the same window specification for several columns lets Spark compute them in one pass.

Common mistakes

  1. Using a window with no partitionBy on large data.
  2. Forgetting a tiebreaker, so “top row” results are not reproducible.
  3. Filtering on a window column in the same select. Compute it with withColumn first, then filter.
  4. Assuming the default frame is the whole partition when an ordering is present.

Interview relevance

Window functions come up as “top N per group”, “deduplicate keeping the latest record” and “compare with the previous period”. Interviewers often follow with: what happens if you leave out partitionBy?

Key takeaway

Choose the partition, the order and the frame deliberately. Partition columns decide both the correctness and the cost of a PySpark window.

By Data Career Hub Editorial · Last reviewed Oct 2026 · Examples run on PySpark 4.2 in local mode; the API is the same in Spark 3.x

Progress is saved in this browser only. No account needed.

Search
Filter by type