After a checkpoint restore or late correction, the same logical match must have a stable identity so the sink can replace it instead of counting it twice.
Join replay output identity and recovery
Identify the relationship
A match between payment revision 2 and signal revision 4 is a different result from the same payment joined to an earlier signal revision. Define a pair identity from both source IDs, both revisions and the join-rule version. Do not use only processing timestamp or task attempt ID. A sink can then upsert the identical result after replay while retaining a deliberate new result when one side changes. Event identity supplies the source-level foundation.
Checkpoint both inputs and state
A successful checkpoint must tie the payment offset, risk offset, keyed join state and sink commit to one recovery boundary. If a worker restores old state but starts one source too far ahead, it can permanently miss a match; if it restarts too far behind without deduplication, it can emit the match twice. Record positions in the replay manifest and verify the sink generation that acknowledged them.
Give unmatched output a lifecycle
An outer join may emit a provisional payment-without-signal result. If an accepted late signal arrives before the horizon, retract or replace the provisional output, then publish the matched result under a distinct logical identity. A sink that only appends cannot represent this correction safely. State whether consumers see provisional rows, final rows or both, and test those semantics under replay.
Rebuild beyond retained state
A correction older than the live join TTL cannot be reconstructed by increasing TTL after the fact. Replay both raw streams into a separate candidate from the oldest required source positions, or compute the affected interval in a bounded batch. Validate keys and totals against the current result, then switch the serving generation atomically. Backfill publication avoids mixing the candidate with live output.
Inject recovery faults
Crash after writing output but before checkpoint acknowledgement. Restart and prove that the sink still contains one logical match. Then expire state and introduce a valid historical correction; the live path must route it to replay instead of inventing an unmatched result. Record source positions, state snapshot, output IDs, provisional retractions and final counts for each attempt.
Implementation
import hashlib
def match_identity(payment_id, payment_revision, signal_id, signal_revision, rule_version):
pieces = (payment_id, str(payment_revision), signal_id,
str(signal_revision), rule_version)
payload = "|".join(f"{len(piece)}:{piece}" for piece in pieces)
return hashlib.sha256(payload.encode()).hexdigest()
matched_id = match_identity("pay-47", 2, "risk-83", 4, "interval-v3")
replayed_id = match_identity("pay-47", 2, "risk-83", 4, "interval-v3")
revised_id = match_identity("pay-47", 3, "risk-83", 4, "interval-v3")
assert matched_id == replayed_id
assert revised_id != matched_idPerformance and operating cost
Computing one identity is O(B) time for B encoded identifier bytes and O(B) temporary space. The sink needs an index or merge key to deduplicate outputs; stateful retractions and retained source streams add storage. A full replay costs both source scans and candidate writes, so reserve it for intervals that can no longer be repaired from live retained state.
Common Mistakes
- Do not build a result key from the worker attempt or wall-clock timestamp.
- Do not append a new match while leaving an obsolete provisional result active.
- Do not assume expired join state can be recovered from a newer checkpoint alone.
Read next
- Interval-join state horizons and watermark skew
- Event envelopes, identity and tombstones
- Keyed state, checkpoints and recovery
- Backfills: rebuild history without exposing a half-written result
- Project: recover a payment-risk interval join
Continue the workflow: Session revisions and sink retractions.
