Count event-time windows with late arrivals
“Our event counter receives records out of order. Count events in fixed time windows according to when they occurred, not when they arrived. A separate source advances a watermark. Keep accepting late records until each window's declared lateness allowance expires, then emit that window exactly once.”
Constructed practice question. Prerequisites: maps and explicit time boundaries. Event time is the record timestamp. A watermark is an externally supplied progress estimate; this exercise does not infer it from local wall time or the largest event seen.
| Contract | Required behavior |
|---|---|
| Input | Positive integer width W, nonnegative lateness L, integer timestamps/watermarks |
| Window | Timestamp t belongs to [floor(t/W)×W, start+W); negatives allowed |
| Add | Returns true when accepted, false when its window is already closed |
| Emit | Advancing watermark returns sorted (start,end,count) for nonempty newly closed windows |
| Boundary | Close when end+L <= watermark; equal watermarks are idempotent |
| Failure/scope | Decreasing watermark/invalid numbers raise ValueError; duplicates count, no retractions |
Optional refresher · the underlying tool
Event time belongs to the event; arrival time is when your process sees it. Tumbling windows have fixed boundaries, often [start,end), so an event exactly at 60 goes in the next 60-second window:
window_size = 60
print((59 // window_size) * window_size) # 0
print((60 // window_size) * window_size) # 60
A watermark is a stated belief about how late events can arrive; it does not change an event's timestamp. Work through duplicate IDs and a late arrival after finalization.
A design choice worth saying aloud
Key aggregates by (entity, window_start), not by arrival minute; otherwise late events land in the wrong window. Keep a processed-event ID set only for the retention period needed by the delivery contract. Once a watermark closes a window, define whether later events are dropped, side-output, or corrections before choosing storage.
What the interviewer expects
The interviewer gives you the scenario and the contract above. Explain what a successful call returns, walk through one example below, and name what your state means before choosing a data structure.
Done means: Accept or reject each event against the watermark and emit every newly closed, nonempty window exactly once in start order.
Now predict each output before looking at the reference; invalid input should leave any existing state unchanged unless the contract says otherwise.
Test-case scenarios to settle before coding
01 · Window boundary
- Input / starting state
- width 10; events at 9 and 10
- Expected result
- counts in
[0,10)and[10,20)
What it is testing: Intervals are half-open.
02 · Allowed late
- Input / starting state
- event arrives behind watermark but window not closed
- Expected result
addreturns true and count includes it
What it is testing: Arrival order differs from event time.
03 · Too late
- Input / starting state
end + lateness <= watermark- Expected result
addreturns false
What it is testing: Closed windows never reopen.
04 · Exact close
- Input / starting state
- watermark equals end plus lateness
- Expected result
- emit that nonempty window once
What it is testing: Equality belongs to closed.
05 · Duplicates/empty
- Input / starting state
- same timestamp twice; untouched windows
- Expected result
- duplicates count; empty windows omitted
What it is testing: Events are observations, not unique IDs.
06 · Invalid/atomic
- Input / starting state
- decreasing watermark or bad number
- Expected result
ValueError; watermark/state unchanged
What it is testing: Failed control input cannot move time backward.
For each case, show which branch or state change produces that result.
With W=10, L=5, add events 9 and 10. Advancing to 10 emits nothing; late event 2
still joins [0,10). At watermark 15 emit (0,10,2). Another event 9 is rejected,
while event 10 belongs to [10,20) and may still arrive. Ask who owns watermark
advancement: choosing it too aggressively discards legitimate delayed records.
Implement independently. Predict whether watermark 14.999 would be a valid input (no: integer contract), and whether watermark 15 still accepts timestamp 9 (no).
Solution, finality, and follow-ups
A baseline groups by arrival time and immediately emits on the local clock's ten-unit boundary. A delayed timestamp 2 arriving after timestamp 10 then enters the wrong window or is lost prematurely. Another tempting baseline closes by maximum observed event time: one far-future outlier can close valid earlier windows.
Store counts keyed by window start. On add, derive start with floor division and compare its close deadline against the current watermark before incrementing. On watermark advancement, validate monotonicity, find active windows whose deadline is reached, sort them by start, return their counts, and remove them from state. Missing/empty windows are not materialized or emitted.
stored windows have not yet closed under the latest watermark, and their counts equal accepted arrivals assigned to those exact half-open intervals. Because the watermark never decreases, a removed window can never be recreated: the add-time boundary test rejects it without keeping unbounded tombstones. This is local single-run finality, not durable delivery or deduplication.
| Arrival/progress | Window 0 count | Window 10 count | Emission |
|---|---|---|---|
| add 9; add 10 | 1 | 1 | none |
| watermark 10; add 2 | 2 | 1 | none |
| watermark 15 | removed | 1 | (0,10,2) |
| add 9 | unchanged | 1 | rejected as too late |
Add costs expected O(1). With A active windows and F closing windows, advancement costs O(A + F log(F+1)) time and O(F) auxiliary/output space. State is O(A), which is not automatically bounded by lateness: arbitrary future timestamps or a stalled watermark can retain arbitrarily many windows. A heap can reduce closure scans, but introduces additional indexing and cleanup work.
Follow-up 1 — emit early and accept corrections. Predict the new output if count 1 was emitted at watermark 10 and event 2 arrives afterward. Send a versioned replacement count 2 or delta +1, then a final marker at 15. Consumers must know which update model applies and deduplicate/reconcile retries.
Follow-up 2 — multiple input partitions. A global watermark normally depends on each active partition's progress, often their minimum. Idle partitions need an explicit exclusion/reactivation policy; taking the maximum closes windows ahead of slow partitions. Future-time skew limits and backpressure can bound retained state.
Senior depth tests exact boundaries and arrival permutations. Lead depth names watermark authority, idle-source policy, restart state, and output retry semantics.
Reference: solution.py (download file, source below); tests cover negative timestamps, event order, late rejection, duplicate counts, empty windows, and monotonic progress.
"""Explicit-watermark tumbling counts with a fixed allowed-lateness duration."""
class TumblingCounter:
def __init__(self, width, allowed_lateness=0):
if not isinstance(width, int) or width <= 0:
raise ValueError("positive integer width required")
if not isinstance(allowed_lateness, int) or allowed_lateness < 0:
raise ValueError("nonnegative integer lateness required")
self.width, self.allowed_lateness = width, allowed_lateness
self.watermark = None
self._counts = {}
def add(self, timestamp):
if not isinstance(timestamp, int):
raise ValueError("integer event timestamp required")
start = (timestamp // self.width) * self.width
close_at = start + self.width + self.allowed_lateness
if self.watermark is not None and close_at <= self.watermark:
return False
self._counts[start] = self._counts.get(start, 0) + 1
return True
def advance_watermark(self, watermark):
if not isinstance(watermark, int):
raise ValueError("integer watermark required")
if self.watermark is not None and watermark < self.watermark:
raise ValueError("watermark cannot decrease")
self.watermark = watermark
closed = sorted(start for start in self._counts
if start + self.width + self.allowed_lateness <= watermark)
return [(start, start + self.width, self._counts.pop(start)) for start in closed]
python -m unittest discover -s curriculum/04-scale-and-evolution/01-data-at-scale/problems/40-event-time-windows -p 'test_*.py'