Order each user's events by time, flag the events that start a new session, and take a running sum of those flags. The running sum becomes the session number.
The three steps
WITH flagged AS (
SELECT user_id, event_ts,
CASE
WHEN LAG(event_ts) OVER (PARTITION BY user_id ORDER BY event_ts) IS NULL
OR event_ts - LAG(event_ts) OVER (PARTITION BY user_id ORDER BY event_ts)
> INTERVAL '30 minutes'
THEN 1 ELSE 0
END AS new_session
FROM clicks
),
numbered AS (
SELECT *,
SUM(new_session) OVER (
PARTITION BY user_id ORDER BY event_ts
ROWS BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW
) AS session_num
FROM flagged
)
SELECT user_id,
user_id || '-' || session_num AS session_id,
MIN(event_ts) AS session_start,
MAX(event_ts) AS session_end,
COUNT(*) AS events
FROM numbered
GROUP BY user_id, session_num;Interval arithmetic differs by engine (TIMESTAMP_DIFF in BigQuery, DATEDIFF in Snowflake). The structure stays the same.
Walk-through
event_ts gap new_session session_num 09:00 first 1 1 09:10 10 min 0 1 09:25 15 min 0 1 10:05 40 min 1 2 10:20 15 min 0 2
The 40-minute gap starts session 2. The running sum goes 1, 1, 1, 2, 2.
Details that matter
- The first event of each user has no previous event, so
LAGis NULL. That must count as a new session, which is why the CASE checks for it. - The session number is only unique per user. Combine it with
user_idto get a global id. - Events with identical timestamps need a tie-breaker in ORDER BY, or the order is random.
- "30 minutes of inactivity" is measured between events. A session can run for hours if the user keeps clicking.
Same idea in Spark
Use the same lag, when and a cumulative sum over a Window.partitionBy("user_id").orderBy("event_ts"). It is the same logic in the DataFrame API. For very large data in Spark Structured Streaming, session windows exist as a built-in, session_window, where you provide the gap duration.