Skip to content
AITroveRead. Build. Understand.
Make this comfortable

Watermark retention and late corrections

Last updated: 6 Oct 20265 min read
tutorial
AdvancedBy AITrove Editorial

A watermark states how far event time has probably advanced; retention defines when an old window can no longer accept ordinary 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

python
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

Continue the workflow: Correction-aware quality alerts.

Continue the workflow: Materialized view refresh strategies.

ai-data
data-engineering
Storage details