Implement a bounded two-stream join that keeps accepted late matches, prevents replay duplicates and routes expired corrections through a candidate rebuild.
Project: recover a payment-risk interval join
Specify the workload
Create 470 payment events and risk signals keyed by merchant. Use a 47-minute look-back and 8-minute look-ahead, with three minutes of accepted lateness. Put several events under one hot merchant, a signal exactly on each interval endpoint and a counterpart one second outside. Define whether unmatched payments are provisional or only final. The horizon lesson supplies the eviction rule.
Build with observable state
Track watermarks separately for payment and risk inputs, keyed record counts, bytes, checkpoint size and late drops. Buffer each side only while the opposite side can still send an accepted match. Simulate an idle risk partition that delays cleanup. Reject an implementation that evicts from a fixed wall-clock timer merely because the payment record has aged.
Exercise replay identity
Assign stable pair IDs using source IDs, revisions and join-rule version. Crash after the sink receives one match but before checkpoint completion. Restore the recorded source positions and state, replay, and show exactly one current sink row for that logical pair. If a signal revision changes, publish the new pair and retire the old one under the documented correction policy. Replay identity is verified at the sink.
Repair an old correction
Expire the live state, then deliver a correction for a payment older than the horizon. Send it to a dated replay queue rather than forcing it through the current join. Rebuild both inputs for the affected interval into a candidate, compare pair keys and merchant totals, then atomically publish the corrected generation. Keep the last-good generation readable until validation passes.
Submit measurable proof
Provide endpoint tests, asymmetric-arrival trace, idle-partition trace, hot-key state peak, checkpoint metadata, crash/replay output IDs, provisional retractions and historical-rebuild diff. Include one intentionally too-short TTL run and show the missing valid match it causes. A green job status without these result checks is insufficient.
Implementation
from datetime import datetime, timedelta, timezone
payments = {"pay-47": ("merchant-29", datetime(2026, 10, 6, 10, 47, tzinfo=timezone.utc))}
signals = {"risk-83": ("merchant-29", datetime(2026, 10, 6, 10, 52, tzinfo=timezone.utc)),
"risk-84": ("merchant-61", datetime(2026, 10, 6, 10, 51, tzinfo=timezone.utc))}
def candidate_pairs(payment_events, risk_events):
return {(payment_id, signal_id)
for payment_id, (payment_merchant, payment_at) in payment_events.items()
for signal_id, (signal_merchant, signal_at) in risk_events.items()
if payment_merchant == signal_merchant
and payment_at - timedelta(minutes=47) <= signal_at
<= payment_at + timedelta(minutes=8)}
first_attempt = candidate_pairs(payments, signals)
replayed_attempt = candidate_pairs(payments, signals)
sink_pairs = first_attempt | replayed_attempt
assert sink_pairs == {("pay-47", "risk-83")}
assert len(sink_pairs) == 1Performance and operating cost
The reference scans P × S pairs; a keyed time-indexed operator narrows probes to eligible records under each key but retains both sides until their horizons close. State bytes, checkpoint duration and sink merge cost rise with active keys and skew. Historical rebuilds require extra scans and candidate storage, so measure frequency separately from the live path.
Common Mistakes
- Do not test only payment-first arrival order.
- Do not accept two sink rows for one pair after a crash and replay.
- Do not silently discard a correction simply because live state has expired.
