Lesson 20 / 25
Event-Time Windows and Watermarks
Aggregate by event time and bound the state with a watermark for late data.
Event time, not arrival time
Real events arrive late and out of order, so you usually aggregate by event time (when it happened, a column in the data) rather than processing time. groupBy(F.window("event_time", "5 minutes")) creates tumbling windows. Because late events could update old windows forever, Spark would have to keep all state. A watermark (withWatermark("event_time", "10 minutes")) says "I accept data at most 10 minutes late": Spark keeps window state until the watermark passes it and then drops it, bounding memory, and events later than that may be ignored. Choose the delay from how late your data really arrives.
Windowed aggregation with a watermark (illustrative)
events is a streaming DataFrame with event_time and user_id columns. Not run here.
counts = (events
.withWatermark("event_time", "10 minutes")
.groupBy(F.window("event_time", "5 minutes"), "user_id")
.count())
query = (counts.writeStream
.outputMode("append") # emit a window once the watermark passes it
.format("delta").option("path", "s3a://lake/gold/user_counts")
.option("checkpointLocation", "s3a://lake/_chk/user_counts")
.start())Quick check: What does a watermark let Spark do?
- Make events arrive faster
- Drop old window state and bound memory while tolerating limited lateness
- Encrypt the stream
- Disable windows
Answer
Drop old window state and bound memory while tolerating limited lateness — It defines how late data may be, so old state can be cleaned up.