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

Interval-join state horizons and watermark skew

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

A stream join needs a finite event-time search interval and a retention rule that cannot evict a match while the other input may still deliver it.

Write the match interval precisely

A payment at time P can match a risk signal from P minus 47 minutes through P plus 8 minutes. Put both bounds in the contract and say whether endpoints are inclusive. A key match alone has no natural stopping point; state grows until an explicit horizon or cleanup policy removes it. A stream-to-table version lookup has a different state model and should not be substituted for a two-stream interval join without checking semantics.

Use both progress signals

An event-time watermark from the payment stream cannot prove the risk stream has finished sending older signals. Track progress per input and per partition. The join may release a stored payment only after the risk-side watermark passes the latest matching risk time plus any agreed lateness. A stalled partition can hold state for longer than the configured time interval; measure this operational lag separately from the logical matching horizon.

Separate lateness from TTL

A state TTL is a memory and storage control, often measured from processing-time access or write. It is not necessarily the event-time rule that determines whether a record can still join. If TTL is shorter than interval plus allowed lateness and observed skew, a valid late counterpart can arrive after the first side has vanished. The job may emit a false unmatched result while all source events are present. Late-correction policy determines what to do after the accepted horizon closes.

Test asymmetric arrival

Send one payment first and its matching signal just before the risk watermark passes the upper bound. Then reverse arrival order. Also hold one risk partition idle while other partitions advance, and send a counterpart outside the interval that must never match. Inspect state size, late-drop count and unmatched-output count. A single happy-path interleaving cannot prove the cleanup timer is safe.

Bound the state budget

Estimate state bytes from active keys, records per key, serialized bytes and checkpoint overhead. A busy merchant may hold many payments under one key; partition-level averages hide that pressure. Apply explicit per-key or workload admission limits where product semantics allow, and send over-limit input to a controlled recovery path rather than silently truncating it. Alert when watermark lag makes the planned state budget unattainable.

Implementation

python
from datetime import datetime, timedelta, timezone

payment_time = datetime(2026, 10, 6, 10, 47, tzinfo=timezone.utc)
earliest_signal = payment_time - timedelta(minutes=47)
latest_signal = payment_time + timedelta(minutes=8)

def signal_matches(payment_at, signal_at):
    return (payment_at - timedelta(minutes=47)
            <= signal_at <= payment_at + timedelta(minutes=8))

def may_evict_payment(payment_at, risk_watermark, allowed_lateness):
    return risk_watermark > payment_at + timedelta(minutes=8) + allowed_lateness

assert signal_matches(payment_time, earliest_signal)
assert signal_matches(payment_time, latest_signal)
assert not signal_matches(payment_time, latest_signal + timedelta(seconds=1))
assert not may_evict_payment(payment_time, latest_signal, timedelta(minutes=3))
assert may_evict_payment(payment_time, latest_signal + timedelta(minutes=4), timedelta(minutes=3))

Performance and operating cost

An indexed keyed join probes matching records in expected O(M) time for M candidates under one key, while active state grows with the records retained across both inputs. Checkpoint bytes and restore time grow with that state. A shorter TTL saves storage but can invalidate accepted matches; a longer one raises checkpoint and backpressure cost. Measure both against an explicit lateness contract.

Common Mistakes

  • Do not use one stream watermark as proof that the other side is complete.
  • Do not confuse processing-time TTL with the event-time match horizon.
  • Do not silently drop a hot key because its state exceeded an average-key budget.

Read next

ai-data
data-engineering
Storage details