Skip to content
Great engineers know 6 min read · Window 3 practice problems ↓

Gaps and islands

Streaks, sessions and status runs: one window pattern behind many interview questions.

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:

QuestionAn island isA new island starts when
Longest login streak per userConsecutive days with a loginThe day is more than 1 after the previous login day
Sessionize a clickstreamEvents with no long pauseMore than 30 minutes since the previous event
How long was each machine in each status?Consecutive readings with the same statusThe status differs from the previous reading
Merge overlapping bookingsBookings that overlap or touchThe start is after the latest end seen so far

Flag, then running sum

One method solves all of them, in three window steps:

  1. 1
    Look back. Use lag over a window partitioned by the entity and ordered by time, to put the previous row's value next to each row.
  2. 2
    Flag the starts. Set new_run to 1 where a new island begins (no previous row, or the gap rule is met), else 0.
  3. 3
    Number 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

weekprev_weeknew_runrun_id
1null11
2101
3201
5312
6502

Run the example

PySparkSpark SQL · Runs of consecutive training weeks per athlete

Input: workouts

athleteweek
ana1
ana2
ana3
ana5
ana6
bo2
bo4
bo5
bo6
bo7
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.

PySparkSpark SQL · Athletes who trained at least 3 weeks in a row

Input: workouts

athleteweek
ana1
ana2
ana3
ana5
ana6
bo2
bo4
bo5
bo6
bo7
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.

Watch out: the shortcut only works when the gap rule is "exactly the next integer" and there are no duplicates. Two workouts in the same week break it, so deduplicate to one row per athlete per week first. For dates, subtract with date_sub. For time gaps (30 minutes), status changes or overlaps, use the flag-and-sum method.

The same method, other questions

Sessions

Flag when 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

Flag when 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

Flag when 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."

Deduplicate to one row per user per day. Then either flag starts with lag (a new streak when the gap from the previous day is more than 1) and number them with a running sum, or use day minus row_number as the group key. Group by user and group to get each streak's length, then take the max per user. Mention the deduplication, the first-row case, and that the flag method generalises to sessions and status runs.

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?"
The default frame for an ordered window is RANGE up to the current row, which includes all rows tied on the order column. With duplicate timestamps, tied rows would get the same running total even if one of them starts a new island. rowsBetween makes it row by row.
Scenario"Merge overlapping booking intervals per room into continuous busy periods."
Order by start. A row starts a new island if its start is after the maximum end of all earlier rows (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."
Only if consecutive means exactly one unit apart. For sessions with a 30-minute timeout, gaps vary, so the difference is not constant within a session. Use the flag-and-sum method.

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 questions

In 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.

Solve: Three-Day Login Streak →

Keep going

Track complete
You reached the end of PySpark functions. Pick another track →

Related lessons

Previous: last(ignorenulls)

Primary sources: functions.lag · Window functions (SQL)