Làm ba bước bằng window function: (1) dùng LAG lấy thời điểm sự kiện trước của cùng user, (2) đánh dấu 1 nếu khoảng cách lớn hơn 30 phút hoặc đó là sự kiện đầu tiên, (3) cộng dồn cờ theo thời gian, kết quả chính là số thứ tự session.
WITH flagged AS (
SELECT user_id, event_id, event_time,
CASE
WHEN LAG(event_time) OVER w IS NULL
OR event_time - LAG(event_time) OVER w > INTERVAL '30 minutes'
THEN 1 ELSE 0
END AS is_new_session
FROM events
WINDOW w AS (PARTITION BY user_id ORDER BY event_time, event_id)
),
numbered AS (
SELECT user_id, event_time,
SUM(is_new_session) OVER (PARTITION BY user_id ORDER BY event_time, event_id
ROWS UNBOUNDED PRECEDING) AS session_no
FROM flagged
)
SELECT user_id, session_no,
MIN(event_time) AS session_start,
MAX(event_time) AS session_end,
COUNT(*) AS events
FROM numbered
GROUP BY user_id, session_no;Cú pháp trên là PostgreSQL. BigQuery dùng TIMESTAMP_DIFF(event_time, prev_time, SECOND) > 1800 (không dùng MINUTE vì hàm cắt phần lẻ: khoảng nghỉ 30 phút 40 giây ra 30, không vượt ngưỡng); Spark có sẵn hàm session_window cho cả batch lẫn streaming. event_id trong ORDER BY để hai sự kiện trùng timestamp luôn xếp cùng một thứ tự giữa các lần chạy.
Phần khó là khi đưa vào pipeline chạy hằng ngày:
- Session vắt qua nửa đêm: job ngày D chỉ đọc partition ngày D sẽ cắt đôi session 23:50–00:20. Cách thường dùng là đọc thêm đuôi ngày D-1 (ít nhất 30 phút cuối, hoặc các session còn mở lưu từ lần chạy trước) rồi ghép.
- Dữ liệu đến trễ: sự kiện của hôm qua tới hôm nay làm thay đổi session cũ, nên job cần tính lại vài ngày gần nhất, hoặc chấp nhận sai số nhỏ và ghi rõ.
- Session ID ổn định: tạo bằng hash của user_id và session_start thay vì số thứ tự, để bảng khác join vào không bị lệch khi chạy lại.
Lưu ý: tổng dồn phải khai ROWS UNBOUNDED PRECEDING. Khi có ORDER BY mà không khai frame, mặc định là RANGE: các dòng bằng nhau trên cột sắp xếp được cộng cùng lúc, và kết quả khác đi khi có giá trị trùng.