A watermark states how far event time has probably advanced; retention defines when an old window can no longer accept ordinary corrections.
Watermark retention and late corrections
Use event time as the grouping key
A card authorization created on day 47 may reach the processor on day 49. Grouping by arrival day changes the question the metric answers. Keep event time, arrival time and source identity separately. Event-time policy defines when a window first emits, but that first result need not be final.
Choose allowed lateness from evidence
Measure the distribution of arrival delay by source and incident class. A watermark lag of 36 hours could cover routine network delays but not a three-day source outage. State the promise: routine updates before the cutoff, explicit correction after it. A larger delay reduces ordinary restatements but holds more open state and delays final reporting.
Retain enough to make corrections coherent
Window state, deduplication IDs and source positions need compatible horizons. If window state lasts 60 hours but event IDs expire after 24, a retransmitted record can be counted again while the window is still open. Retention planning must include recovery, replay and legal boundaries, not only the dashboard delay.
Publish revisions with identity
A correction after a provisional result should replace the same account-window key with a higher revision or publish a versioned snapshot. Appending a second total without a revision contract makes readers sum both. Keep prior and revised totals, reason, source event ID and publication time in the audit record. Send very late events to a repair path rather than discarding them silently.
Test the edge of the cutoff
Build fixtures for an on-time event, one just before the watermark, one just after it and a duplicate. Assert the expected revision number and exact cents. Then restore from a checkpoint before the late record and replay; the published revision must be the same. A test that only checks aggregate totals misses duplicate revision rows.
Implementation
from collections import defaultdict
records = [
{"id": "auth-47", "account": "acct-47", "event_hour": 47, "arrival_hour": 48, "cents": 3100},
{"id": "auth-48", "account": "acct-47", "event_hour": 47, "arrival_hour": 50, "cents": 700},
]
def window_total(events, cutoff_hour):
accepted = [row for row in events if row["arrival_hour"] <= cutoff_hour]
if len({row["id"] for row in accepted}) != len(accepted):
raise ValueError("duplicate authorization identity")
totals = defaultdict(int)
for row in accepted:
totals[(row["account"], row["event_hour"])] += row["cents"]
return dict(totals)
assert window_total(records, 49) == {("acct-47", 47): 3100}
assert window_total(records, 50) == {("acct-47", 47): 3800}Performance and operating cost
A keyed pass costs O(N) expected time and O(K + N) memory here, including duplicate validation. In production, open windows and retained identities scale with throughput multiplied by the retention horizon. Extending the watermark by one day can materially increase state, checkpoint bytes and recovery time; measure each before changing the promise.
Common Mistakes
- Do not treat a first window emission as immutable unless the contract says so.
- Do not expire deduplication state before a window can still be corrected.
- Do not append a revised total as a second additive fact.
Read next
- Event time and late arrivals: close windows with an explicit correction policy
- Keyed state, checkpoints and recovery
- Compaction, retention and the replay horizon
- Project: release a recoverable streaming risk ledger
- Replay manifests and audit trails
Continue the workflow: Correction-aware quality alerts.
Continue the workflow: Materialized view refresh strategies.
