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

Session-window gaps and late bridges

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

A late event can connect two activity sessions that were already reported separately; session boundaries therefore need an explicit correction contract.

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

python
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"] == 134

Performance 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.

Read next

ai-data
data-engineering
Storage details