A late event can connect two activity sessions that were already reported separately; session boundaries therefore need an explicit correction contract.
Session-window gaps and late bridges
Define the gap before assigning IDs
For each account, consecutive events belong to one session when the event-time gap is at most eleven minutes. Equality is intentional. The records at 09:00 and 09:22 start separate sessions; an event at 09:11 bridges them. Process arrival order does not decide membership. A user-facing timeout may differ from this analytical gap, so name the clock and key in the contract. Event-time handling explains why ingestion order alone produces inconsistent groups.
Keep enough state to find a bridge
A session operator must retain intervals or their event membership while an accepted late event could extend an endpoint or connect neighboring intervals. Closing the first result at its current end is not the same as forgetting it. Measure watermark progress for every source partition, and publish the latest time at which an event can still merge a session. Retention policy sets the boundary for later replay.
Publish revisions, not independent totals
Suppose two provisional sessions have totals of 74 and 51 units. A bridging event worth 9 units produces one corrected total of 134. Emitting 134 while keeping 74 and 51 active makes a dashboard count 259. Retire both former result identities, then publish the merged result. A sink that lacks retract or upsert semantics requires a versioned replacement table and a serving pointer rather than append-only rows.
Treat boundary changes as identity changes
Using the current start timestamp as a permanent session ID fails when an earlier late event extends the start. Give each published session a revision and track which earlier outputs it supersedes. A stable identity can be derived from immutable event membership for a particular revision, but a merge still creates a new identity and must remove the two old ones. Replay-safe output identity provides a related correction pattern.
Test both late and too-late paths
Deliver the 09:11 bridge after both initial results have fired but before allowed lateness expires; expect one active merged row. Deliver the same event after state expires; expect a replay request or a documented rejection, never a silent third session. Add duplicate event IDs, equal-gap boundaries, a stalled partition and a hot account. Compare the final table against a batch recomputation from unique source events.
Implementation
from datetime import datetime, timedelta, timezone
def group_sessions(activity_events, gap):
ordered = sorted(activity_events, key=lambda event: (event[1], event[0]))
sessions = []
for event_id, event_at, amount in ordered:
if not sessions or event_at - sessions[-1]["end"] > gap:
sessions.append({"start": event_at, "end": event_at,
"event_ids": [event_id], "amount": amount})
else:
session = sessions[-1]
session["end"] = event_at
session["event_ids"].append(event_id)
session["amount"] += amount
return sessions
base = datetime(2026, 10, 6, 9, tzinfo=timezone.utc)
activity = [("activity-47", base, 74),
("activity-83", base + timedelta(minutes=22), 51)]
assert len(group_sessions(activity, timedelta(minutes=11))) == 2
activity.append(("activity-61", base + timedelta(minutes=11), 9))
merged = group_sessions(activity, timedelta(minutes=11))
assert len(merged) == 1 and merged[0]["amount"] == 134Performance and operating cost
Sorting N events for a bounded replay costs O(N log N) time and O(N) space. An indexed streaming implementation can inspect neighboring sessions under one key instead, but the retained state, timers and correction records increase with active sessions and permitted lateness. A longer grace interval improves repair coverage while raising checkpoint size and delaying finality.
Common Mistakes
- Do not use processing-time arrival order as the event-time session boundary.
- Do not append a merged total while its two superseded totals remain active.
- Do not assume a late bridge can be repaired after all contributing state has been discarded.
