◐ Off-By-One · answer catalog

temporal-join-retraction-frontier

1 answer(s)pythonpython3

The failure mode is eager eviction driven by an over-trusted watermark. A temporal (interval) join needs both

📦 Source in repository (JSON)

Answer

Bounded-Memory Temporal Join with Retractions: Diagnosis and Fix

Root-cause analysis

The failure mode is eager eviction driven by an over-trusted watermark. A temporal (interval) join needs both members of every pair. The naive design advances a single global watermark, emits/evicts an event as soon as the watermark exceeds t + W, and drops anything at or below the watermark.

Three independent bugs fall out of that:

  1. Unsound cutoff. With any allowed lateness L > 0, an event below the watermark can still arrive and match. Evicting at wm > t + W loses those matches. The correct safe point is t + W + L < min(wm_A, wm_B), because the opposite side must also have moved past the match window minus the lateness slack.
  2. Global watermark head-of-line blocking. One stalled key or source pins the global watermark. State for every key accumulates, and a hard cap then forces lossy drops. The frontier must be per key (and the cap enforced without silent drops).
  3. Duplicate double-counting. Deduplicating with an ever-growing id set is unbounded and contradicting the cap. In fact dedup state only needs ids of currently buffered (unsealed) events: any later copy lands below the seal threshold and is provably redundant.

Two consequences the problem statement is testing:

Model and invariants

Events: (side ∈ {A,B}, key, timestamp, event_id). Join output: all pairs with equal key and |t_A − t_B| ≤ W.

State per key:

state meaning
wm[k][side] max timestamp observed for (k, side) — a lower bound on future timestamps
buf[k][A], buf[k][B] unsealed events, kept sorted by timestamp
seen[k] ids of events currently in buf (dedup)

Invariants that make the result exact:

  1. Frontier: F[k] = min(wm[k][A], wm[k][B]).
  2. Seal (eviction) rule: event x may be evicted only when t_x + W + L < F[k]. At that instant every possible partner on the other side has already arrived (watermark honest, lateness ≤ L), and is still buffered, so emitting x against the opposite buffer is complete.
  3. Unique emission: a pair (a,b) is emitted exactly once — when the first of its two members is evicted. Candidates are popped in increasing timestamp order, so the partner is guaranteed present.
  4. Duplicate rule: id in seen[k] ⇒ skip. Id not in seen[k] but t < F[k] − L − W ⇒ the event is below the seal threshold; under the lateness bound it is a redundant copy (or a contract violation) and is dropped without affecting the final result.
  5. Hard cap is not a drop rule. When size ≥ cap, the event goes to the spill queue (disk) and is replayed when capacity frees. Dropping unsealed events is only ever safe if you accept approximate results.

Bounded memory guarantee: at any moment per key only events with timestamps in (F[k]−W−L, F[k]] are retained, so cap ≥ max_k (#events in a W+L horizon) suffices for zero spilling. A single hot/stalled key can still consume the cap; that key’s source must be back-pressured (or spilled), not the data dropped.

Exact fix (code)

Self-contained, tested implementation. Save as bounded_temporal_join.py and run.

"""
Bounded-memory temporal (interval) join over out-of-order streams with
per-key watermarks, bounded lateness, duplicate suppression, hard state cap
with disk spill, and a retraction protocol for speculative outer joins.
"""
from bisect import insort
from collections import defaultdict, deque

NEG = float('-inf')
POS = float('inf')


class IntervalJoin:
    """
    Inner interval join: (a,b) emitted iff key(a)==key(b) and |t_a - t_b| <= W.
    An event x is sealed once  t_x + W + L < min(wm[key]['A'], wm[key]['B']).
    Only sealed events are evicted; every pair is emitted exactly once.
    """
    def __init__(self, W, lateness=0, max_state=POS, spill_enabled=True):
        self.W, self.L = W, lateness
        self.cap, self.spill_enabled = max_state, spill_enabled
        self.buf = defaultdict(lambda: {'A': [], 'B': []})
        self.wm = {}
        self.seen = defaultdict(set)          # ids of currently buffered events
        self.size = 0
        self.spill = deque()                  # disk-backed pending queue
        self.emitted = []
        self.n_spilled = self.n_late_dropped = self.n_dup = 0

    def _ensure(self, k):
        if k not in self.wm:
            self.wm[k] = {'A': NEG, 'B': NEG}
        return self.buf[k]

    def _frontier(self, k):
        return min(self.wm[k]['A'], self.wm[k]['B'])

    def _seal_threshold(self, k):
        return self._frontier(k) - self.L - self.W

    def _insert(self, side, k, t, eid, payload):
        b = self._ensure(k)
        if t > self.wm[k][side]:
            self.wm[k][side] = t
        insort(b[side], (t, eid, payload))
        self.seen[k].add(eid)
        self.size += 1

    def _advance(self, k):
        b = self.buf[k]
        while True:
            thr = self._seal_threshold(k)
            side = None
            for s in ('A', 'B'):                       # pop globally smallest t
                if b[s] and b[s][0][0] < thr:
                    if side is None or b[s][0][0] < b[side][0][0]:
                        side = s
            if side is None:
                break
            t, eid, _ = b[side].pop(0)
            self.size -= 1
            self.seen[k].discard(eid)
            other = 'B' if side == 'A' else 'A'
            for (yt, yeid, _) in b[other]:
                if yt < t - self.W:
                    continue
                if yt > t + self.W:
                    break
                self.emitted.append((k, eid, yeid) if side == 'A'
                                    else (k, yeid, eid))

    def _drain_spill(self):
        while self.spill and self.size < self.cap:
            side, k, t, eid, payload = self.spill.popleft()
            if eid in self.seen[k]:
                self.n_dup += 1
            elif t < self._seal_threshold(k):
                self.n_late_dropped += 1
            else:
                self._insert(side, k, t, eid, payload)
                self._advance(k)

    def push(self, side, k, t, eid, payload=None):
        self._ensure(k)
        if eid in self.seen[k]:
            self.n_dup += 1
            return 'dup'
        if t < self._seal_threshold(k):
            self.n_late_dropped += 1
            return 'too-late'
        if self.size >= self.cap:
            if not self.spill_enabled:
                self.n_late_dropped += 1
                return 'dropped-cap'
            self.spill.append((side, k, t, eid, payload))
            self.n_spilled += 1
            return 'spilled'
        self._insert(side, k, t, eid, payload)
        self._advance(k)
        self._drain_spill()
        return 'ok'

    def close(self):
        """End of stream: drain disk ignoring cap, then flush all watermarks."""
        while self.spill:
            side, k, t, eid, payload = self.spill.popleft()
            if eid not in self.seen[k]:
                self._insert(side, k, t, eid, payload)
                self._advance(k)
        for k in list(self.wm):
            self.wm[k]['A'] = self.wm[k]['B'] = POS
            self._advance(k)
        return sorted(self.emitted)


class RetractingLeftOuter:
    """
    A-left outer interval join with low-latency speculative emission.
    A provisional (a, None) row is emitted at the SOFT frontier t_a+W < wm_B,
    retracted if a late B within L matches, and its retraction state is dropped
    at the HARD frontier  t_a+W+L < min(wm_A, wm_B).
    """
    def __init__(self, W, L):
        self.W, self.L = W, L
        self.bufA, self.bufB = [], []
        self.wmA = self.wmB = NEG
        self.prov = {}                       # aid -> {'t','rows','final'}
        self.seenA, self.seenB = set(), set()
        self.log = []                        # ('assert'|'retract', row)
        self.n_retract = 0

    def _matches(self, a):
        return [(b, be) for (b, be) in self.bufB if abs(b - a[0]) <= self.W]

    def _emit(self, kind, row):
        self.log.append((kind, row))
        if kind == 'retract':
            self.n_retract += 1

    def _soft_seal(self, a):
        aid = a[1]
        ms = self._matches(a)
        rows = [(aid, be) for (_, be) in ms] or [(aid, None)]
        for r in rows:
            self._emit('assert', r)
        self.prov[aid] = {'t': a[0], 'rows': rows}

    def _refresh(self, aid):
        a = next((x for x in self.bufA if x[1] == aid), (self.prov[aid]['t'], aid))
        ms = self._matches(a)
        newrows = [(aid, be) for (_, be) in ms] or [(aid, None)]
        for r in self.prov[aid]['rows']:
            self._emit('retract', r)
        for r in newrows:
            self._emit('assert', r)
        self.prov[aid]['rows'] = newrows

    def _flush_A(self):
        F = min(self.wmA, self.wmB)
        for a in list(self.bufA):
            t, eid = a
            hard = t < F - self.L - self.W
            soft = t < self.wmB - self.W
            if eid not in self.prov and (hard or soft):
                self._soft_seal(a)
            if hard:
                self.prov.pop(eid, None)     # rows final; drop retraction state
                self.bufA.remove(a)

    def push(self, side, t, eid):
        if side == 'A':
            if eid in self.seenA:
                return
            self.seenA.add(eid)
            self.bufA.append((t, eid))
            self.wmA = max(self.wmA, t)
            self._flush_A()
        else:
            if eid in self.seenB:
                return
            self.seenB.add(eid)
            self.bufB.append((t, eid))
            self.wmB = max(self.wmB, t)
            for aid, st in list(self.prov.items()):
                if abs(st['t'] - t) <= self.W:
                    self._refresh(aid)
            self._flush_A()

    def close(self):
        self.wmA = self.wmB = POS
        self._flush_A()
        for a in list(self.bufA):
            if a[1] not in self.prov:
                self._soft_seal(a)
        final = set()
        for kind, row in self.log:
            final.add(row) if kind == 'assert' else final.discard(row)
        return final


# ----------------------------- verification -----------------------------
def brute_force(events, W):
    """Exact batch answer, dedup by (side,key,id)."""
    seen, A, B = set(), defaultdict(list), defaultdict(list)
    for side, k, t, eid in events:
        if (side, k, eid) in seen:
            continue
        seen.add((side, k, eid))
        (A if side == 'A' else B)[k].append((t, eid))
    return {(k, ea, eb)
            for k in set(A) | set(B)
            for ta, ea in A[k] for tb, eb in B[k]
            if abs(ta - tb) <= W}


if __name__ == '__main__':
    import random

    def actual_lateness(events):
        wm = defaultdict(lambda: {'A': NEG, 'B': NEG})
        worst = 0.0
        for s, k, t, _ in events:
            f = min(wm[k]['A'], wm[k]['B'])
            if f > NEG:
                worst = max(worst, f - t)
            wm[k][s] = max(wm[k][s], t)
        return max(worst, 0.0)

    # 1. exactness under arbitrary permutations
    bad = 0
    for seed in range(500):
        rng = random.Random(seed)
        W = rng.choice([0, 1, 3, 5])
        evs = [(rng.choice('AB'), rng.randrange(3), rng.randrange(40), f'e{i}')
               for i in range(rng.randrange(5, 60))]
        order = list(range(len(evs))); rng.shuffle(order)
        sh = [evs[i] for i in order]
        j = IntervalJoin(W, int(actual_lateness(sh)) + 1)
        for s, k, t, e in sh:
            j.push(s, k, t, e)
        bad += set(j.close()) != brute_force(sh, W)
    assert bad == 0
    print('inner exact: 500/500 random permutations match brute force')

    # 2. duplicates + hard cap with spill
    bad = 0
    for seed in range(300):
        rng = random.Random(seed)
        W = rng.choice([1, 2, 4])
        evs = [(rng.choice('AB'), rng.randrange(3), rng.randrange(20), f'e{i}')
               for i in range(rng.randrange(5, 50))]
        evs += [evs[i] for i in rng.sample(range(len(evs)), min(8, len(evs)))]
        rng.shuffle(evs)
        j = IntervalJoin(W, int(actual_lateness(evs)) + 1, max_state=6)
        for s, k, t, e in evs:
            j.push(s, k, t, e)
        bad += set(j.close()) != brute_force(evs, W)
    assert bad == 0
    print('duplicates + hard cap/disk-spill: 300/300 match brute force')

    # 3. adversarial skew: B stalls, A floods; drop-on-cap loses data
    evs = [('A', 'k', t, f'A{t}') for t in range(50)] + [('B', 'k', 25, 'B25')]
    random.Random(1).shuffle(evs)
    L = int(actual_lateness(evs)) + 1
    j = IntervalJoin(1, L, max_state=4)
    for s, k, t, e in evs:
        j.push(s, k, t, e)
    exact = set(j.close())
    j2 = IntervalJoin(1, L, max_state=4, spill_enabled=False)
    for s, k, t, e in evs:
        j2.push(s, k, t, e)
    lossy = set(j2.close())
    assert exact == brute_force(evs, 1) and lossy != exact
    print(f'adversarial skew: exact={len(exact)} via {j.n_spilled} spilled; '
          f'drop-on-cap loses {len(exact) - len(lossy)}')

    # 4. retraction for speculative outer join
    evs = [('A', 5, 'A5'), ('B', 10, 'B10'), ('B', 5, 'B5')]
    r = RetractingLeftOuter(1, 3)
    for s, t, e in evs:
        r.push(s, t, e)
    assert r.close() == {('A5', 'B5')} and r.n_retract >= 1
    print('retraction log:', r.log)
    print('ALL CHECKS PASSED')

Expected output:

inner exact: 500/500 random permutations match brute force
duplicates + hard cap/disk-spill: 300/300 match brute force
adversarial skew: exact=3 via 47 spilled; drop-on-cap loses 3
retraction log: [('assert', ('A5', None)), ('retract', ('A5', None)), ('assert', ('A5', 'B5'))]
ALL CHECKS PASSED

Adversarial analysis

Watermark skew. With F[k] = min(wm_A, wm_B), a stalling side correctly freezes sealing: no matches can be retired while the other side is behind. But that same rule means a single hot/stalled key accumulates state up to cap. Per-key watermarks fix cross-key head-of-line blocking (a stall in key k1 no longer pins k2); they do not fix within-key growth. The cap must therefore back-pressure that key’s source or spill. The skew test shows 47 events spilled while the result stayed exact; the spill_enabled=False variant silently lost all 3 pairs.

Duplicates. seen[k] holds only buffered ids. Once an event seals, any later copy has t < F − L − W, hits the too-late branch, and is a no-op. Dedup memory is O(buffered), not O(all history). Caveat: if copies can arrive arbitrarily late (no lateness bound), you need either an unbounded id set or an idempotent sink keyed by pair id — the latter keeps memory bounded and is the recommended production choice.

Impossibility. Take a key, an A event at t=0, and an adversary that delivers the matching B event at an arbitrary future time. To be exact you must keep A forever, so any finite cap without spill/backpressure must either lose the match or drop a later event. Hence: exactness ⇔ bounded lateness L (or unbounded/disk state). This is why the fix replaces "drop at cap" with "spill at cap".

Correctness argument for the seal rule. When t_x + W + L < F[k], the honest watermark and lateness bound guarantee no future opposite event can fall in [t_x − W, t_x + W]; every such partner has arrived and, because it is not yet evictable by the same bound, is still in the opposite buffer. Emitting then is complete and final. Popping candidates in increasing t ensures each pair is emitted exactly once: at the first member’s eviction the other member is still buffered, and after eviction it cannot be emitted again.

Verification summary

check result
500 random permutations vs. batch brute force (inner join) 0 mismatches
300 cases with injected duplicates and cap=6 (spill on) 0 mismatches
adversarial skew, stalled side + flooded side exact via 47 spilled; drop-on-cap loses 3 pairs
forced speculative null→pair retraction (outer join) correct final {('A5','B5')}, 1 retraction
duplicate handling seen bounded to buffered ids only

Run: python3 bounded_temporal_join.py → ALL CHECKS PASSED.

Evidence & signatures

# Evidence
- Problem class: temporal-join-retraction-frontier
- Model: openrouter/deepseek/deepseek-v4.1-flash
- Solved: 2026-10-06T04:02:07.933Z
- Verification: solution produced by pi in sandbox; see signatures.json
{"description": "Design a bounded-memory temporal join over two out-of-order event streams with per-key watermarks, late-data retractions, and a hard cap on retained state. Specify the minimal state and eviction rules that preserve exact results under arbitrary arrival order, and analyze adversarial watermark skew and duplicate events.", "environment": "python3", "language": "python", "model": "openrouter/deepseek/deepseek-v4.1-flash", "problem_class": "temporal-join-retraction-frontier", "provider": "openrouter", "solved_at": "2026-10-06T04:02:07.936Z", "version": "3.11"}
Generated from the verified corpus · MIT licensedBack to the catalog