Gaps and islands
Streaks, sessions and status runs: one window pattern behind many interview questions.
On this page
Show code in
Every code block on the page follows this.
Read first: lag / lead · row_number / rank (skip if you know them)
The pattern
An island is a run of consecutive rows that belong together: days in a login streak, events in one session, readings while a machine was in the same status. A gap is where one island ends and the next begins. Many interview questions are this pattern in disguise:
| Question | An island is | A new island starts when |
|---|---|---|
| Longest login streak per user | Consecutive days with a login | The day is more than 1 after the previous login day |
| Sessionize a clickstream | Events with no long pause | More than 30 minutes since the previous event |
| How long was each machine in each status? | Consecutive readings with the same status | The status differs from the previous reading |
| Merge overlapping bookings | Bookings that overlap or touch | The start is after the latest end seen so far |
Flag, then running sum
One method solves all of them, in three window steps:
- 1Look back. Use
lagover a window partitioned by the entity and ordered by time, to put the previous row's value next to each row. - 2Flag the starts. Set
new_runto 1 where a new island begins (no previous row, or the gap rule is met), else 0. - 3Number the islands. A running
sum(new_run)over the same window gives every row its island number: it only goes up at a start.
Then group by the entity and island number to get each island's start, end and length.
workouts (ana)
| week |
|---|
| 1 |
| 2 |
| 3 |
| 5 |
| 6 |
after the three steps
| week | prev_week | new_run | run_id |
|---|---|---|---|
| 1 | null | 1 | 1 |
| 2 | 1 | 0 | 1 |
| 3 | 2 | 0 | 1 |
| 5 | 3 | 1 | 2 |
| 6 | 5 | 0 | 2 |
Run the example
Input: workouts
| athlete | week |
|---|---|
| ana | 1 |
| ana | 2 |
| ana | 3 |
| ana | 5 |
| ana | 6 |
| bo | 2 |
| bo | 4 |
| bo | 5 |
| bo | 6 |
| bo | 7 |
from pyspark.sql import functions as F from pyspark.sql.window import Window w = Window.partitionBy("athlete").orderBy("week") running = w.rowsBetween(Window.unboundedPreceding, Window.currentRow) flagged = (workouts .withColumn("prev_week", F.lag("week").over(w)) .withColumn("new_run", F.when(F.col("prev_week").isNull() | (F.col("week") - F.col("prev_week") > 1), 1).otherwise(0)) .withColumn("run_id", F.sum("new_run").over(running))) result = (flagged.groupBy("athlete", "run_id") .agg(F.min("week").alias("first_week"), F.max("week").alias("last_week"), F.count("*").alias("weeks")) .orderBy("athlete", "run_id"))
WITH prev AS ( SELECT athlete, week, LAG(week) OVER (PARTITION BY athlete ORDER BY week) AS prev_week FROM workouts ), flagged AS ( SELECT athlete, week, CASE WHEN prev_week IS NULL OR week - prev_week > 1 THEN 1 ELSE 0 END AS new_run FROM prev ), numbered AS ( SELECT athlete, week, SUM(new_run) OVER (PARTITION BY athlete ORDER BY week ROWS BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW) AS run_id FROM flagged ) SELECT athlete, run_id, MIN(week) AS first_week, MAX(week) AS last_week, COUNT(*) AS weeks FROM numbered GROUP BY athlete, run_id ORDER BY athlete, run_id
Switch to PySpark to edit and run this example in your browser.
Ana trained three weeks in a row (weeks 1 to 3), then two; Bo's longest run is four weeks. For "the longest run per athlete", aggregate once more with max("weeks").
The row_number shortcut
For consecutive integers or dates, there is a well-known trick: inside a run, week and row_number() both go up by 1, so their difference stays constant. That difference identifies the island.
Input: workouts
| athlete | week |
|---|---|
| ana | 1 |
| ana | 2 |
| ana | 3 |
| ana | 5 |
| ana | 6 |
| bo | 2 |
| bo | 4 |
| bo | 5 |
| bo | 6 |
| bo | 7 |
from pyspark.sql import functions as F from pyspark.sql.window import Window w = Window.partitionBy("athlete").orderBy("week") result = (workouts .withColumn("block", F.col("week") - F.row_number().over(w)) .groupBy("athlete", "block") .agg(F.min("week").alias("first_week"), F.count("*").alias("weeks")) .filter(F.col("weeks") >= 3) .orderBy("athlete", "first_week"))
SELECT athlete, MIN(week) AS first_week, COUNT(*) AS weeks FROM (SELECT athlete, week, week - ROW_NUMBER() OVER (PARTITION BY athlete ORDER BY week) AS block FROM workouts) GROUP BY athlete, block HAVING COUNT(*) >= 3 ORDER BY athlete, first_week
Switch to PySpark to edit and run this example in your browser.
date_sub. For time gaps (30 minutes), status changes or overlaps, use the flag-and-sum method.The same method, other questions
Sessions
unix_timestamp(ts) - unix_timestamp(lag(ts)) > 1800 or there is no previous event. The running sum is the session number per user.Status runs
status != lag(status) or there is no previous row. Use eqNullSafe if status can be null, so null to null is not a change.Overlapping intervals
start > max(end) over all earlier rows (a frame ending at 1 preceding). Then each island is one merged interval: min(start), max(end).What it costs Optional deep dive
All windows here share the same partitionBy and orderBy, so Spark sorts once and evaluates them together: one shuffleshuffle: Moving rows between machines so that all rows with the same key end up together. Needed by joins, groupBy and sorting, and usually the most expensive step of a job. Learn more → by user, one sort, then the group by. A skewedskew: When one key has far more rows than the others, so one task does most of the work while the rest wait. Learn more → entity (one bot with millions of events) puts all its rows in one tasktask: The work for one partition in one stage, run on one CPU core. Learn more →, so check for hot keys first.
Common mistakes
- Not deduplicating before the row_number trick — Two rows on the same day make the difference change inside a streak.
- Forgetting the first row — lag is null on the first row of each partitionpartition: A chunk of a DataFrame's rows. Spark processes each partition as one task, so partitions decide how much work runs in parallel. Learn more →; it must start an island, or it gets run_id 0 and merges wrongly.
- Using the default frame for the running sum — With ties in the order column, the default range frame sums tied rows together. Use rowsBetween.
- Comparing timestamps as strings — Convert to seconds (unix_timestamp) before subtracting.
Interview prep
THE QUESTION
"Find each user's longest streak of consecutive login days."
Avoid saying: a self-join of each day with the next two days. It works for exactly three days and does not generalise to "the longest streak".
What the interviewer asks next. Answer out loud first, then open the strong answer.
Follow-up"Why does the running sum need an explicit rowsBetween frame?"
Scenario"Merge overlapping booking intervals per room into continuous busy periods."
max(end) over a frame ending at 1 preceding), or it is the first row. Running sum of that flag gives the group; then min(start) and max(end) per group. Comparing only with the previous row's end is the classic bug: an earlier long booking can cover later ones.Trap"day - row_number() works for timestamps if you convert them to seconds."
What you learned
- What the gaps-and-islands pattern is, and how to recognise it in a question
- The flag-and-running-sum method, step by step
- The row_number subtraction shortcut, and when it works
- How sessionization and status runs are the same pattern
Key takeaways
- Gaps-and-islands groups consecutive rows that belong together.
- Flag island starts with lag, then a running sum numbers the islands.
- value - row_number() is a shortcut for consecutive integers or dates without duplicates.
- Sessions, status runs and overlapping intervals use the same method.
Check yourself
3 questionsIn the flag-and-sum method, what does the running sum of new_run give?
Show the answer
An island number that only increases at a start. The sum only goes up on rows flagged as a start, so every row of an island shares the same number.
Why can week - row_number() fail?
Show the answer
Duplicate weeks or gaps of more than one break the constant difference. The trick assumes each row advances by exactly 1.
For sessionization with a 30-minute timeout, what starts a new island?
Show the answer
More than 30 minutes since the previous event, or no previous event. That is the gap rule; the running sum then numbers sessions.
Practice it
Interview problems that use this: write the PySpark, run it, and get graded on hidden tests.
Keep going
Track completeYou reached the end of PySpark functions. Pick another track →
Related lessons
lag / leadPySpark functions · 4 min read
row_number / rankPySpark functions · 3 min read
last(ignorenulls)
Previous: last(ignorenulls)
Primary sources: functions.lag · Window functions (SQL)