Streaming · Flashcard

A stream job keeps counts in memory. How does it recover after a crash without double-counting?

hard Resilience

Answer

It takes periodic checkpoints that store its state together with the input offsets it had reached. After a crash it restores the last checkpoint and replays input from those offsets, so state reflects each event exactly once (Flink works this way). Output to other systems still needs idempotent or transactional writes.

Review this in your daily deck All cards in Streaming