skip to content

In Flink, what is the difference between a checkpoint and a savepoint?

level: middleimportance: should knowfreq 68%

answer

  1. one is Flink's, one is yours
  2. triggered on a timer versus on purpose
  3. cheap and incremental versus portable
  4. what survives a job cancellation
  5. the artefact you upgrade a job with

basics

~20 s

Checkpoints are Flink's automatic periodic snapshots for failure recovery, owned and cleaned up by the system. Savepoints are deliberately triggered, self-contained snapshots in a portable format, taken for planned work: upgrades, rescaling, migrations and cluster moves.

solid answer

~50 s

Mechanically they are both snapshots of the same state, and both are produced by the same barrier machinery. The difference is **ownership and intent**. A checkpoint is triggered by Flink on an interval, is optimised for cheap and frequent writing (it can be incremental and in a backend-specific format), is owned by the running job, and is deleted when the job is cancelled unless you configure retention. A savepoint is triggered by a human or an API call — `bin/flink savepoint <jobId>`, or `bin/flink stop` to take one and stop atomically — is always self-contained, is never deleted by Flink, and in Flink 2.3 defaults to a canonical format that can be restored into a different state backend. That makes it the artefact for planned work: new job code, a different parallelism, another cluster. Savepoints cost more to write, which is why you do not take one every thirty seconds.

code

bash · 6 lines
bash
# planned upgrade: stop with a savepoint, then restore the new build
bin/flink stop --savepointPath s3://state/savepoints 4f21c9a3b7e1
bin/flink run -s s3://state/savepoints/savepoint-4f21c9-1a2b3c target/job-v2.jar

# take a savepoint without stopping (e.g. before a risky config change)
bin/flink savepoint 4f21c9a3b7e1 s3://state/savepoints

go deeper

for a junior

Recall the intent split: checkpoints are automatic and for crash recovery, savepoints are manual and for planned changes like deploying a new version. Knowing which one you take before a release is enough at this level.

for a middle

Explain ownership, format and lifecycle — incremental and backend-native versus self-contained and canonical, system-deleted versus operator-deleted — and name the CLI commands for taking one and restoring from it.

for a senior

Show the operational discipline: stop-with-savepoint for clean shutdowns, stable operator UIDs, keeping the savepoint as a rollback until the new version proves itself, and the maximum-parallelism ceiling on rescaling.

for a principal

Own the upgrade and rollback policy for a platform of jobs: where snapshots are stored, retention and cost, how version upgrades are staged, and what standards teams must follow so any job can be restored years later.

## Same machinery, different contract Both a checkpoint and a savepoint are produced by the barrier-based snapshot algorithm, and both capture the same two things: every operator's state and every source's read position. If you only look at the algorithm, they are the same object. The difference is the contract around it — who triggers it, who owns it, what format it is written in, and what you are allowed to do with it afterwards. ## Checkpoints: the system's business A checkpoint exists so that the job can recover from its own failures without human involvement. - **Triggered by Flink** on the configured interval. - **Optimised for frequency.** With an incremental-capable state backend a checkpoint can upload only the state files that changed since the previous one, which makes large state affordable to snapshot often. The format is whatever that backend finds cheapest. - **Owned by the job.** Flink retains a bounded number of completed checkpoints and deletes older ones. When the job is cancelled the checkpoint data goes with it, unless you asked for externalized retention (`execution.checkpointing.externalized-checkpoint-retention: RETAIN_ON_CANCELLATION`; the default is `NO_EXTERNALIZED_CHECKPOINTS`), in which case it survives and *you* become responsible for deleting it. - **Restored automatically** on failover, with no operator action. ## Savepoints: your business A savepoint exists so that a human can do something deliberate to a job. - **Triggered explicitly**, by CLI, REST or the client API. - **Always self-contained.** It never references files belonging to an earlier snapshot, so you can copy or move the directory and it still restores. - **Canonical format by default** — a unified, backend-independent layout, which is what lets you restore a savepoint taken with one state backend into a job using a different one, and what gives the best chance of restoring across Flink versions — though Flink 2.0 does not guarantee state compatibility with 1.x snapshots at all. Writing it costs more than a checkpoint, because incremental tricks and backend-native layouts are given up for portability. A *native*-format savepoint is available when you want a faster write and do not need to switch backends. - **Always aligned.** A savepoint never uses unaligned mode, even on a job with unaligned checkpoints enabled. - **Never deleted by Flink.** The lifecycle is entirely yours. ## The operations savepoints exist for The classic four: 1. **Job upgrade.** Stop with a savepoint, deploy a new jar, restore. State survives the code change as long as the operators that own it can still be identified — which is why setting stable operator UIDs with `uid("...")` on every stateful operator is the single most important habit for an upgradeable Flink job. 2. **Rescaling.** Restore at a different parallelism; keyed state is redistributed by key group, which is why the maximum parallelism is fixed at first run and cannot be changed by a savepoint. 3. **Migration.** Move a job to another cluster, another region, or another Flink version. 4. **A/B or blue-green.** Start a second job from the same savepoint and compare, since a self-contained savepoint can be restored more than once. ## Doing it safely ```bash # take a savepoint and stop the job in one atomic operation bin/flink stop --savepointPath s3://state/savepoints <jobId> # resume the new build from it bin/flink run -s s3://state/savepoints/savepoint-abc123 target/job.jar ``` Stopping *with* a savepoint is better than cancel-then-savepoint: it stops the sources, takes the snapshot and terminates as one step, so no records are processed after the snapshot and any transactional sink commits cleanly instead of leaving an open transaction behind. The separate `--drain` flag additionally emits a final `MAX_WATERMARK` to flush event-time windows; use it only when the job will not be resumed. ## The overlap people get wrong Two things blur the line and cause confusion in interviews. First, a **retained checkpoint can be resumed from manually** with the same `-s` flag as a savepoint. So "you can only resume from savepoints" is wrong, and Flink's own compatibility table goes further: an aligned checkpoint supports changed job code, rescaling and a Flink minor-version upgrade. What you give up is ownership and portability. A checkpoint cannot switch the state backend, is not self-contained or relocatable (an incremental one references files from earlier checkpoints), and an *unaligned* checkpoint supports neither arbitrary job changes nor a minor-version upgrade. Second, because the job owns its checkpoints, the one you meant to upgrade from may already be gone — subsumed by a newer checkpoint, or deleted on cancellation because retention was never enabled. Treating a retained checkpoint as your upgrade path works right up until the day you change a backend, move the files, or find the checkpoint was unaligned, and then it does not. ## Choosing in practice Run checkpointing continuously with an interval matched to how much replay you can tolerate, and enable retention so a crashed-and-cancelled job leaves something behind. Take a savepoint before every deliberate change — deploy, rescale, migration — and keep it until the new version has proved itself, because it is your rollback.

  • Why must every stateful operator have an explicitly set UID before you rely on savepoints?
    State in a savepoint is filed under the operator's ID. Without an explicit uid(), Flink derives one from the topology, so inserting or reordering an operator changes the generated IDs and the restore either fails or silently drops that operator's state. An explicit, stable UID decouples the state's identity from the shape of the graph, which is the whole point of an upgradeable job.
  • Can you change a job's parallelism when restoring from a savepoint, and is there a limit?
    Yes — keyed state is stored in key groups that are redistributed across the new subtasks on restore. The limit is the maximum parallelism, which fixes the number of key groups and is baked in when the job first runs. You can scale freely up to it, but exceeding it requires reprocessing state through the State Processor API rather than a plain restore.
  • Why prefer stop-with-savepoint over cancelling the job and then taking a savepoint?
    Because cancel-then-savepoint is not a thing: cancelling ends the job, so the snapshot must come first, and doing them separately leaves a window where records are processed after the snapshot. Stop-with-savepoint stops the sources, snapshots and terminates as one operation, so nothing is processed past the cut and transactional sinks commit their final transaction instead of leaving one open.

A checkpoint is a game's autosave, taken constantly and overwritten; a savepoint is the save file you name yourself before trying something risky, and you decide when to delete it.

saying these in an interview costs you the question

  • Says savepoints are just checkpoints with a longer interval
  • Thinks retained checkpoints cannot be resumed from at all
  • Believes Flink deletes old savepoints automatically
  • Claims savepoints are always incremental to save space
  • Relies on generated operator IDs and expects upgrades to restore

context