A Flink job checkpoints every 60 seconds. State has grown so that one checkpoint now takes 90 seconds to complete. No configuration changes. What happens over the next hour, and which metric will not show it?
Show the full answer Hide the answer
What happens, minute by minute
Checkpoint N starts at t=0 and is still running when N+1 falls due at t=60. With the default limit of one concurrent checkpoint, the scheduler cannot start N+1 and the trigger is simply skipped. The effective interval becomes the checkpoint duration, so it is now 90 seconds and climbing in step with state size.
Within the hour, the duration and the state size feed each other. Aligned checkpointing makes barriers queue behind in-flight records at every operator input, and the slowest input decides the alignment time. As the job backpressures, alignment takes longer, so the checkpoint takes longer, so more input accumulates behind the barrier. The loop runs one way.
If the sink is configured for exactly-once through two-phase commit, output becomes visible only when a checkpoint completes. A consumer of the job's output sees its end-to-end latency jump from about 60 seconds to 90 and rising, while the job's own per-record processing latency has not changed at all.
Eventually the checkpoint timeout fires (10 minutes by default), the checkpoint is aborted, and depending on the tolerable-failed-checkpoints setting the job either keeps limping or fails. On restart it restores from the last completed checkpoint, so recovery now means restoring a larger state plus re-reading everything since that checkpoint, which is the longest outage of the sequence and arrives after the signal has been visible for an hour.
What the metrics show, and do not show
Consumer lag on the input topic looks normal, because the job is still consuming. Throughput is flat. Error rate is zero. Pod CPU and memory look busy but not alarming. This is a failure with no error signal, which is why it is usually found by a downstream team asking why their table is late.
The signals that do show it: lastCheckpointDuration, lastCheckpointSize, checkpointAlignmentTime, numberOfFailedCheckpoints, and the gap between the sink's commit timestamp and the event time it covers. Alert on checkpoint duration as a fraction of the interval, at a threshold like 50%, rather than on failure, because by the time checkpoints fail the restore is already expensive.
The fix, in order
- Unaligned checkpoints. Barriers overtake in-flight records, so alignment time stops growing with backpressure. The price is that in-flight data goes into the checkpoint, making it larger and the restore slower.
- Incremental checkpointing with RocksDB, so only changed files are uploaded rather than the whole keyed state.
- State time-to-live or a bounded key space. The other two treat the symptom. State growth is monotonic unless something expires it, so nothing here self-heals: a job with unbounded keys and no TTL returns to this condition at a larger size.
- Raise the interval deliberately rather than letting it drift, and accept the matching increase in replayed work after a failure.
When not to checkpoint at all
If the job is a stateless map, filter or route, this entire failure class does not exist. A plain consumer loop committing offsets after its write, with an idempotent sink, has no state to snapshot, restarts in seconds and needs none of these settings. The checkpointing machinery is the price of keyed state and event-time windows. Teams that adopt a stateful framework for stateless work pay this bill and get nothing for it.