Build an event-time session table that repairs late bridges, retracts old results and survives checkpoint replay without double counting.
Project: recover late customer activity sessions
Create a controlled event feed
Generate 470 unique activity events across 29 accounts with an eleven-minute inactivity gap and three minutes of accepted lateness. Include two sessions separated by 22 minutes and a bridge at their shared eleven-minute boundary. Inject a duplicate event ID, one event just beyond the gap, a late event after finality and an idle partition. Record source offsets and event timestamps independently of arrival times.
Implement the live session operator
Key by account and merge event-time intervals when neighboring gaps are at most eleven minutes. Track one watermark per input partition, timer state and active session counts. Publish provisional results only if the sink supports later retractions. Keep a clear rule for when a session is final, and demonstrate that a stalled partition can delay that declaration without changing which records belong together. Gap semantics drive the expected answer.
Add an atomic correction sink
After the bridge arrives, retract both former session revisions and publish one merged revision. Crash after preparing the correction but before checkpoint acknowledgement, then replay the same offsets. The current table must contain one merged session with the exact unique-event total. Test a second crash after the sink commit, where retry must detect the already applied generation. Revision identity makes those attempts comparable.
Handle the expired correction
Send an event beyond accepted lateness into a dated repair queue. Rebuild the affected account interval from retained raw events into a candidate generation, compare membership and totals, and switch the serving pointer only after validation. Record a deliberate failed candidate and show the last-good table remains available. Do not change the live watermark backward to force the event into an expired window.
Submit proof from both views
Deliver the arrival-order trace, final event-time grouping, retraction log, generation pointer, checkpoint positions, current-table uniqueness check and batch reconciliation. Report state peak for the hot account and the count of side-output late records. Include a negative run that appends the merged row without retiring the originals; its overcount should be visible in the test result.
Implementation
from datetime import datetime, timedelta, timezone
base = datetime(2026, 10, 6, 9, tzinfo=timezone.utc)
arrival_order = [
("activity-47", base, 74),
("activity-83", base + timedelta(minutes=22), 51),
("activity-61", base + timedelta(minutes=11), 9),
("activity-61", base + timedelta(minutes=11), 9),
]
unique_events = {event_id: (event_at, amount)
for event_id, event_at, amount in arrival_order}
ordered = sorted(unique_events.values())
groups = []
for event_at, amount in ordered:
if not groups or event_at - groups[-1]["end"] > timedelta(minutes=11):
groups.append({"end": event_at, "amount": amount})
else:
groups[-1]["end"] = event_at
groups[-1]["amount"] += amount
assert len(groups) == 1
assert groups[0]["amount"] == sum(value[1] for value in unique_events.values()) == 134Performance and operating cost
The reference recomputes N accepted events in O(N log N) time and O(N) space. A keyed streaming operator reduces repeated full sorting at the cost of retained per-account session state and correction metadata. Measure state bytes, checkpoint duration and sink commits under the hot-account fixture; replay cost depends on the oldest repairable raw-event offset.
Common Mistakes
- Do not count the duplicate bridge twice.
- Do not let a crash expose both old and merged revisions as current.
- Do not hide a too-late event by silently moving the watermark backward.
