Count event-time windows with late arrivalsLESSON 10.07 · 7 OF 20 IN CHAPTER
PART C / Data systems at scale
Step 160 of 252
LESSON 10.07 · 7 OF 20 IN CHAPTERTry it, then open the solution

Count event-time windows with late arrivals

THE PROBLEM

“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
add returns 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
add returns 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.

Diagram: Test-case scenarios to settle before coding

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.

Diagram: Test-case scenarios to settle before coding

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.

solution.py · solution.py
"""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'