In Apache Flink, checkpoints are automated, lightweight snapshots managed internally by the engine to recover state automatically when task failures occur. Savepoints are manually triggered, version-portable state snapshots created by engineers for operational lifecycle tasks such as application upgrades, cluster rescaling, and infrastructure migrations. While checkpoints prioritize minimal processing overhead and fast recovery times, savepoints generate self-contained, canonical metadata designed for long-term administrative operations.
Automated checkpoints and barrier alignment
Checkpoints run continuously on a scheduled interval configured by the application. Flink coordinates distributed snapshots using stream alignment barriers based on the Chandy-Lamport algorithm:
Source Stream ---> [Record A] [Record B] [Checkpoint Barrier] [Record C] ---> Operator
When an operator with multiple input channels processes data, it performs barrier alignment:
- The operator pauses consumption on channels that deliver the checkpoint barrier early.
- It continues reading from slower channels until all input barriers arrive.
- Once all barriers align, the operator writes its state to persistent storage like Amazon S3 and forwards the barrier downstream.
Unaligned checkpoints under backpressure
When downstream operators experience severe processing bottlenecks, input channel buffers fill up, creating backpressure that stalls aligned checkpoints. Barriers get trapped behind backlogged records, causing checkpoints to time out and leaving the cluster without recent recovery points.
To resolve this, Flink provides unaligned checkpoints. When enabled, operators snapshot the contents of in-flight channel buffers alongside operator state as soon as the barrier arrives at the head of any channel, bypassing the alignment wait. This ensures checkpoints complete successfully even under heavy backpressure.
Manual savepoints for pipeline management
Savepoints are explicitly triggered through the Flink command-line interface or REST API before maintenance actions.
Unlike checkpoints, which may use internal binary formats optimized for a specific RocksDB version, savepoints serialize state into a standardized, canonical format. This portability allows engineers to stop a running job, refactor business code in Java or SQL, change operator parallelism from ten to fifty, and resume execution cleanly without state loss.