advanced 3 min answer

A windowed job writes one row per hour into a reporting table. Overnight the autoscaler kills a task manager twice - at 03:14 and at 05:40 - and each time the job restores from a checkpoint taken about two minutes earlier and resumes with no errors. Next morning 2 of the 24 hourly rows are roughly double. Every checkpoint in the window completed successfully. What failed?

checkpointingtimersidempotent-sinkduplicatesrecovery
Show the full answer Hide the answer

The trigger

The checkpoint was correct. The restore was correct. What was not restored is the part of the system that lives outside the checkpoint: the rows already written to the reporting table.

Between the 03:12 checkpoint and the 03:14 kill, the job's event-time timer for the 02:00 window fired, the window was flushed and the sink wrote the 02:00 row. Restoring to 03:12 rewinds operator state and source offsets to a point before that flush. Registered timers are checkpointed state, so the 02:00 timer comes back pending, the replayed records rebuild the same window, and it fires a second time. The sink is a plain insert with no key constraint, so the second flush appends a second row. The dashboard sums the rows for the hour and shows double.

Two kills, two windows caught inside their post-checkpoint flush gap, two doubled rows. The other 22 hours flushed safely inside a checkpoint interval and are untouched.

Why nothing reported an error

Every signal a team watches was healthy, and correctly so. Checkpoint success means the state was durably written, not that the outside world agrees with it. Lag was zero. Throughput was normal. Restart count went up by two and nothing is wired to it.

At-least-once processing plus a non-idempotent sink is a duplicate generator with a trigger nobody treats as an event. The generator ran twice, in the quiet hours, at a rate of about one wrong row per restart.

The structural fix versus the tempting local fix

The tempting fix is to deduplicate in the reporting query. It hides the generator, and it fails the first time a replay writes two rows for the same hour with different values - then there is no basis to choose.

Two structural options, and you must pick one explicitly:

  1. Make the output part of the same atomic unit as the state. A transactional or two-phase-commit sink commits the write with the checkpoint. The price is that output becomes visible only at checkpoint boundaries, so a 60-second interval gives 60-second visibility - a latency floor you now own.
  2. Make the write idempotent on a deterministic key. (window_end, group_key, logic_version) with an upsert, so a replayed flush overwrites instead of appending. Never include the attempt, the job id or a timestamp in that key. Note the new hazard this creates: an upsert replayed before the source has caught up can overwrite a complete value with a partial one, which is worse because it looks right. Guard it by writing only on window completion and keeping the watermark in the row.

Detection worth having: a sink-side ratio of rows written to distinct keys written, alerted when it leaves 1.0, and a restart counter on the same dashboard as the business numbers.

Common weak answers

  • "The checkpoint was corrupt." It was not, and chasing it wastes the incident. The state was exactly what it claimed to be.
  • "Enable exactly-once mode." That makes state-and-offset commits atomic inside the engine. It says nothing about an external sink that cannot participate in the commit, which is where these rows came from.
  • "Shorten the checkpoint interval." It narrows the window in which a flush can be orphaned and does not close it, and it adds steady-state overhead for a partial mitigation.