Event time là lúc sự kiện thực sự xảy ra (timestamp ghi trên thiết bị); processing time là lúc hệ thống nhận được nó. Hai mốc này lệch nhau vì mạng chập chờn, app mobile gửi bù khi có mạng lại, hoặc consumer bị dừng một lúc.
Tính "số đơn mỗi 10 phút" theo processing time thì đơn lúc 12:04 nhưng tới lúc 12:11 bị đếm sang cửa sổ 12:10. Muốn số đúng nghiệp vụ phải gom theo event time — nhưng khi đó engine phải giữ state của các cửa sổ cũ vì không biết còn dữ liệu trễ tới nữa hay không.
Watermark giải bài toán đó: dữ liệu trễ quá X so với event time lớn nhất đã thấy thì bỏ qua. Nhờ vậy engine biết khi nào đóng cửa sổ và giải phóng state.
from pyspark.sql.functions import window
counts = (events
.withWatermark("event_time", "15 minutes")
.groupBy(window("event_time", "10 minutes"), "shop_id")
.count())Ở đây Spark đảm bảo không bỏ dòng nào trễ dưới 15 phút; dòng trễ hơn có thể bị bỏ.
Chọn ngưỡng là một đánh đổi: ngưỡng lớn thì số đúng hơn nhưng state nặng và kết quả cuối ra chậm; ngưỡng nhỏ thì nhanh, nhẹ nhưng mất dữ liệu trễ.
Dữ liệu trễ hơn watermark: cách làm thực tế là để stream phục vụ số nhanh và gần đúng, còn job batch hằng đêm tính lại theo event time từ dữ liệu thô và ghi đè partition của ngày, nên số chốt cuối cùng vẫn đúng.
Lưu ý: ở output mode append, kết quả một cửa sổ chỉ được ghi ra sau khi watermark vượt qua cuối cửa sổ, nên độ trễ đầu ra xấp xỉ bằng ngưỡng watermark. Stakeholder cần biết điều này trước khi kỳ vọng dashboard "thời gian thực".