advanced 3 min answer

A job interval-joins click events to impression events with a two-hour window. To fix a parser bug the team resets only the click consumer's offsets to three days ago and lets the running job catch up. Walk through what happens, what the output looks like, and what stops it.

interval joinwatermarksreplaybackfillstate
Show the full answer Hide the answer

Minute by minute

The moment the click source restarts three days back, its event-time watermark goes backwards by three days. A join's watermark is the minimum across its inputs, so the whole operator's notion of "now" drops with it.

  • Minute 1. Replayed clicks carry event times from three days ago. The impression side of the join holds only records inside the two-hour window around the current watermark. The three-day-old clicks have nothing to match, so they join to nothing and are emitted as unmatched or silently dropped, depending on the join type.
  • Minutes 1 to 40. Live impressions keep arriving with today's timestamps. Because the operator's watermark is three days behind, no timer fires and no state is cleared, so the impression buffer grows with everything that arrives instead of being trimmed to two hours. At 20k events/s and 300 bytes a record that is roughly 6 GB per hour of growth that nothing will ever evict. State grows at the full arrival rate until the task runs out of managed memory and the job restarts, whereupon it does the same thing again.
  • When the replay catches up. The watermark jumps forward three days in one step. Every held timer fires at once, a burst of output lands, and anything now outside its allowed lateness is dropped. Attribution collapses during the replay and then spikes.

Where it amplifies

The failure is a positive feedback loop: the misalignment prevents state cleanup, the state growth causes a restart, and the restart re-reads from the same offsets. Output is wrong in both directions — under-reported during the replay, over-reported in the catch-up burst — so a revenue metric derived from it cannot be corrected by averaging.

What the user sees

Click-through and attribution rates fall to near zero, then overshoot. Nobody is paged, because throughput is high, error rate is zero and lag on the replayed partition is falling as intended.

What stops it

  • Replay both sides from the same event-time point. A join is a function of event time; only one side moving is not a replay, it is a misconfiguration.
  • Run the backfill as a separate job with its own sources, its own state and its own output table, then swap the results. This keeps the live job's watermark untouched and is the reason "backfill in place" is a bad default for any stateful join.
  • Watermark alignment and idleness settings, which bound how far one source may drift ahead of another, turning a silent misalignment into visible back-pressure. Alignment costs throughput, because the fast source is deliberately held back to the slow one, and that is the price of making this failure impossible rather than merely unlikely. Flink has shipped watermark alignment for its sources since the 1.15 release in 2022.
  • An alert on the gap between the maximum event time of each input, which is the signal that goes wrong first and is not derivable from lag in records.

What would have to be true for it to self-heal

Nothing about the running job. The impressions the replayed clicks needed have already passed out of the window, and no amount of waiting brings them back — the data required for the correct answer is no longer in state. Recovery means a second, properly aligned replay of both sides.

When not to use a stream-stream join at all

If the enrichment side is reference data rather than an event stream, model it as a compacted table and use a stream-table join. A table has no watermark to misalign: replaying the stream side alone is then safe, and this whole class of failure disappears. Prefer stream-table over stream-stream whenever one side is state rather than events — which is more often than teams assume.