Join Replay Misalignment
also called One-Sided Replay, Split-Watermark Join
The failure that follows from replaying one input of a stateful join while the other stays live - event times no longer overlap, window state is never cleaned, and the output is wrong in both directions with no error raised.
A team fixes a parser bug and resets only the click consumer's offsets to three days ago. The job is joining clicks to impressions with a two-hour window, it restarts cleanly, throughput is high, and the lag on the replayed partitions falls exactly as intended. The attribution numbers collapse to near zero, then overshoot.
A stateful join is a function of event time, so replaying one side is not a replay: it is a misconfiguration that the runtime cannot distinguish from a very late source.
Why it matters
Two mechanisms combine, and both are counter-intuitive.
First, an operator's watermark is the minimum across its inputs. When one source jumps backwards three days, the whole operator's notion of the present goes with it. Records from the live side are buffered rather than matched, and the replayed records find nothing to match because the counterpart data left state long ago.
Second, state cleanup is driven by that same watermark. No timers fire, so nothing is evicted. The buffer grows at the full arrival rate of the live side until the task exhausts its state budget and restarts, which re-reads the same offsets and repeats the cycle. At 20k events/s and 300 bytes a record, a three-day replay accumulates roughly 1.5 TB of state that nothing will evict, against a two-hour window whose steady-state footprint is about 40 GB.
The result is output that is wrong in both directions — under-reported during the replay and over-reported in the catch-up burst when the watermark jumps forward — which means the error cannot be corrected by averaging or by re-running the same thing again.
Implementation patterns
- Replay both sides from the same event-time point, not the same offset. Offsets are not comparable across topics; event time is.
- Backfill in a separate job with its own sources, state and output table, then swap or merge results. The live job's watermark is never disturbed.
- Configure watermark alignment and source idleness so one input cannot drift arbitrarily far from another, converting a silent misalignment into visible back-pressure.
- Alert on the spread between inputs' maximum event times, which moves first and is not derivable from record lag.
- Model reference data as a compacted table rather than a stream wherever possible; a stream-table join has one watermark, so one-sided replay of the stream is safe.
Industry example
The pattern is generic to every event-time join engine, and the honest framing is a documented mechanism rather than a company anecdote: Flink's watermark semantics (minimum across inputs, with timers and state TTL driven by it) are what make one-sided replay destructive, which is exactly why watermark alignment was added as a configurable guard. A pipeline attributing clicks to impressions at a large ad-supported platform is the canonical place it is discovered, usually during the first backfill after the join has been in production long enough for someone to find a parser bug worth fixing. Flink added watermark alignment for sources in its 1.15 release in 2022 for exactly this class of problem.
Failure scenarios
- Near-zero output during the replay, read as "the fix worked and there was less traffic".
- State growth to out-of-memory, then a restart loop that repeats the whole sequence.
- A burst of late-fired windows when the watermark jumps, double-counting against whatever was emitted before.
- Permanent loss: the records the replayed side needed have already aged out of state, so the correct answer is no longer computable from the running job.
Trade-offs
Avoiding this costs isolation. A separate backfill job means duplicated topology, a second output destination, and a reconciliation step to merge results — real work, on a schedule nobody planned for. Watermark alignment costs throughput, because the fast source is deliberately held back to the slow one. Both are cheaper than a revenue metric wrong in both directions for a week.
When not to use it
The concept does not apply where there is no shared event-time clock: stateless maps, filters and per-record enrichment against an external store can all have one input replayed freely, and adding alignment machinery there buys nothing. Equally, if the join window is seconds rather than hours, a full replay of both sides is cheap enough that special handling is over-engineering — replay everything and let the window refill.
Interview question
Q: An engineer proposes resetting one topic's offsets to reprocess three days of a two-hour interval join. Tell me what happens, what you would do instead, and what signal you would put in place before anyone tries this again.
What a strong answer covers: the watermark is the minimum across inputs, so one-sided reset moves the operator's clock backwards · state cleanup stops and memory grows until the job dies · output is wrong in both directions and not self-correcting · the alternative is a separate aligned backfill job or a replay of both sides from the same event time · and the pre-emptive signal is the spread between inputs' maximum event times, plus watermark alignment configured as a guard.
Quick check
Quiz: Why does replaying one side of an interval join stop state from being cleaned up? The operator's watermark is the minimum across inputs, so it moves backwards with the replayed source and no cleanup timer fires.
Flashcard: You reset offsets on one input of a two-hour interval join to three days ago. What does the output look like? — Near-nothing while the replay runs, because the live side's records are outside the replayed event times, then a burst when the watermark jumps forward, with state growing to failure in between.