A checkpoint is a consistent snapshot of job progress: operator state, offsets/positions, and sometimes sink transaction markers. On failure, the job restarts from the last successful checkpoint instead of from the beginning of time.
time ---> process... checkpoint#1 ... process... checkpoint#2 ... CRASH
\
restart from #2
(state + offsets)What gets saved (typical)
- Kafka consumer offsets / input positions
- Window aggregates and keyed state
- Timers / watermarks progress
- Coordination for transactional sinks (Flink two-phase commit style)
Ops notes
- Too frequent → overhead; too rare → long recovery and more reprocessing
- Checkpoint timeout / size growth = smell for state or sink issues
- Store checkpoints on durable storage (S3, HDFS, GCS)
Interview tip: "Checkpoint = recover state + input position together." That pairing is what makes exactly-once / consistent recovery possible.