lag / lead
Read the previous or next row in a window to compute changes, gaps and streaks.
On this page
Show code in
Every code block on the page follows this.
You will learn
- How lag and lead read neighbouring rows
- How to compute day-over-day changes and gaps
- What the default value does at the edges
- The gaps-and-islands trick for streaks
lag(col, n) returns the value from n rows before the current row in the window; lead looks n rows ahead. At the edge of a partition they return null, or the default you pass.What it does
Inside each window partition, ordered by the window's order, lag("close") gives the previous row's close and lead("close") the next row's. They need an ordered window; a frame clause does not apply to them.
Step by step
prices (ACME)
| day | close |
|---|---|
| 03-03 | 100 |
| 03-04 | 104 |
| 03-05 | 101 |
| 03-06 | 108 |
Output
| day | close | prev | change |
|---|---|---|---|
| 03-03 | 100 | null | null |
| 03-04 | 104 | 100 | 4 |
| 03-05 | 101 | 104 | -3 |
| 03-06 | 108 | 101 | 7 |
The first day has no previous row, so prev is null and so is the change computed from it.
Run the example
from pyspark.sql import functions as F from pyspark.sql.window import Window w = Window.partitionBy("ticker").orderBy("day") result = (prices .withColumn("prev", F.lag("close").over(w)) .withColumn("change", F.col("close") - F.col("prev")) .withColumn("next_close", F.lead("close").over(w)))
SELECT *, LAG(close) OVER w AS prev, close - LAG(close) OVER w AS change, LEAD(close) OVER w AS next_close FROM prices WINDOW w AS (PARTITION BY ticker ORDER BY day)
Switch to PySpark to edit and run this example in your browser.
Offset and default
| Call | Returns |
|---|---|
F.lag("close") | Previous row, or null on the first row |
F.lag("close", 2) | Two rows back |
F.lag("close", 1, 0.0) | Previous row, or 0.0 when there is none |
F.lead("close") | Next row, or null on the last row |
The default only fills the edges, where no row exists. If the previous row exists and its value is null, lag returns that null.
Gaps between events
With dates, lag lets you measure the time since the previous event: F.datediff("day", F.lag("day").over(w)). A gap above a threshold starts a new session, which is the basis of sessionisation.
Streaks: gaps and islands
To find runs of consecutive days, compare each row with the previous one and start a new group whenever the run breaks. A running sum of those break flags numbers the groups:
- 1
is_break = 1when the gap to the previous day is not exactly 1, else0(use lag). - 2
streak_id = sum(is_break)over the window ordered by day, from the start up to the current row. - 3Group by
streak_idto get each streak's start, end and length.
The same shape solves "consecutive months of growth", "login streaks" and "status unchanged since". Interviewers love it because it combines lag with a running total.
Under the hood
lag and lead run in the same shuffle-and-sort as other window functions over the same spec. After sorting, Spark keeps a small buffer of rows, so the offset itself is nearly free. The cost is the shuffle, and the risk is skew: one ticker with most of the rows makes one slow task.
Common mistakes
Ordering by a string date in the wrong format
Forgetting partitionBy
Assuming the default replaces nulls
Key takeaways
- lag reads a previous row, lead a following one, within each ordered partition.
- At the partition edges they return null or the default you pass.
- Day-over-day change is
col - lag(col). - lag plus a running sum of break flags finds streaks (gaps and islands).
Check yourself
3 questions1. What does F.lag("v", 1, 0) return when the previous row exists but its v is null?
Show the answer
null. The default is used only when there is no row at the offset. An existing null is returned as null.
2. Why must the lag window be partitioned by ticker?
Show the answer
Otherwise the first row of one ticker reads the last row of another. Without partitioning, rows from different tickers are neighbours in one ordered sequence.
3. Which pattern finds runs of consecutive days?
Show the answer
lag to flag breaks, then a running sum of the flags. Flags mark where a run breaks; their running sum gives every run its own id.
Practice it
Interview problems that use lag / lead: write the PySpark, run it, and get graded on hidden tests.
Go deeper
Primary sources: functions.lag · functions.lead