Skip to content

How to upgrade a Flink job from a savepoint without losing state

Flink
Chad Harris·October 3, 2026·12 min read

A Flink job needs an upgrade: new application code, a new Flink version or a different parallelism. The job is stopped with a savepoint and resubmitted, and then one of two things happens. The submission fails with an IllegalStateException or a StateMigrationException, or the job starts RUNNING but its counters, windows and deduplication state are empty, and the Kafka source reads from somewhere unexpected.

A savepoint is a consistent image of a streaming job’s execution state, created through Flink’s checkpointing mechanism and used to stop and resume, fork or update a job, as the savepoints documentation defines it. The upgrade only keeps state if the new job can map every piece of that image back onto its own operators. In practice a restore keeps its state when every stateful operator had an explicit uid before the savepoint was taken, the job was stopped without draining, the new parallelism does not exceed the stored max parallelism, and the restored block in the checkpoint statistics is filled in after the deploy.

This page covers planned upgrades and restores. If a restore has already failed or come up empty, the diagnosis section below starts with the cheapest rule-out. If checkpoints fail while the job is running, the guide to Flink checkpoints that fail or time out is the better starting point. If the job restores and then fails over and over, see a Flink job stuck in a restart loop. For how upgrades fit into running Flink in production more widely, see the complete Flink guide.

The safe upgrade procedure

The upgrading documentation bases every application upgrade and every Flink version upgrade on savepoints. The procedure below follows it in four steps, and each step has a check before the next one starts.

1. Check the job before taking the savepoint

Flink matches state in a savepoint to operators by operator ID. Each operator has a default ID derived from its position in the topology, so the upgrading documentation says a modified application can only be started from a savepoint if the operator IDs were set explicitly. The uid has to be in the running job before the savepoint is taken. Adding it only to the new version does not help, because the savepoint already holds the generated IDs.

.uid("orders-dedup")

The production readiness checklist asks for an explicit max parallelism as well, for the same reason. It cannot be changed after a job has started without discarding that operator’s state. The configuration key is listed in the Flink configuration reference.

setMaxParallelism(int maxParallelism)
pipeline.max-parallelism

To make a missing uid fail at submission instead of at the next upgrade, the upgrading documentation points to one execution config call.

ExecutionConfig#disableAutoGeneratedUIDs

2. Stop the job with a savepoint

The command. Stop with savepoint takes the savepoint and stops the job as one action. The CLI documentation describes it as the graceful way to stop a streaming job. The sources send a last checkpoint barrier that triggers the savepoint, and only after it completes do they finish.

bin/flink stop --savepointPath s3://bucket/savepoints :jobId

The REST equivalent is one call, and it is the call any management tool makes as well.

POST /jobs/:jobid/stop

What to avoid. The same documentation marks cancel with savepoint as deprecated in favour of stop. The savepoints documentation adds a related reason. Since Flink 1.15, a savepoint taken while the job keeps running is an intermediate savepoint that is not used for recovery and does not commit side effects, so a job that later fails falls back to an earlier checkpoint.

Do not drain the job for an upgrade. The drain flag sends a final maximum watermark, fires all event time timers and flushes windows, and the CLI documentation says it is meant for terminating a job permanently. Resuming a drained job can lead to incorrect results.

--drain

Large state. For a job with large state the client can time out while waiting. The savepoints documentation suggests detached mode, after which the savepoint status is read through the REST API.

-detached

Timing the stop. A savepoint is taken through the checkpoint mechanism, so the job’s backpressure and checkpoint duration before the stop are a fair guide to how long the stop will take. The guide to checkpointing under backpressure says barrier propagation can dominate checkpoint time under heavy backpressure, and the page on finding the bottleneck behind Flink backpressure covers that diagnosis.

3. Verify the savepoint before deploying

A completed stop prints the savepoint path. Confirm what the savepoint contains before the new version goes anywhere. The State Processor API has a SQL table function that reads savepoint metadata. It returns one row per operator, including the uid, the uid hash, the parallelism and the max parallelism.

LOAD MODULE state
SELECT * FROM savepoint_metadata('s3://bucket/sp-abc123')

Compare the operator-uid column with the uid() calls in the new code. An operator in the savepoint with no matching uid in the new job is state that the restore will refuse, or drop if told to. An operator-max-parallelism value lower than the parallelism planned for the new job is a restore that will fail.

Then check where the savepoint lives. The savepoints documentation requires the target directory to be accessible by the JobManager and every TaskManager, and the upgrading documentation requires that all savepoint data is accessible from a new Flink installation under the same absolute path.

4. Deploy and restore

Submit the new version with the savepoint path. The path can point to the savepoint directory or to its _metadata file.

bin/flink run -s s3://bucket/sp-abc123 :runArgs

The same restore can be set in configuration, which is how deployments that do not use the CLI usually pass it.

execution.state-recovery.path

Then confirm the job really restored from that savepoint. The REST API returns checkpoint statistics for a job, including a restored block with the external path and whether it was a savepoint.

GET /jobs/:jobid/checkpoints
latest.restored.external_path
latest.restored.is_savepoint

An empty restored block means the job started without state, whatever the submission command said.

A restore that does carry state can still look busy at first. Timers are checkpointed with state, so they are restored as well, and the Flink process function documentation says checkpointed processing-time timers that were due before the restore fire immediately. A window or timeout driven by processing time can therefore fire as soon as the job starts.

Claim mode decides whether Flink may delete the savepoint files after the restore. The default, NO_CLAIM, never deletes them, so keep the savepoint until the first checkpoint after the restore has completed. Until then a failure recovers from it. The mode has its own configuration key, and the claim mode section below compares the modes.

execution.state-recovery.claim-mode

Diagnose a failed or empty restore

The first question is whether the restore failed loudly or succeeded quietly with less state than expected. A loud failure points to the uid, schema, parallelism or path branches below, matched by the exception text in the table. A quiet empty restore points to allowing non-restored state on an operator whose uid changed, or to an empty restored block in the checkpoint statistics.

Rule out the cheapest cause first. A restore with nothing to restore from looks the same as lost state, so confirm that the savepoint path exists and that the JobManager and every TaskManager can reach it before reading any exception.

F1 Why a savepoint restore fails or comes up without state
What shows first What proves it, and the fix
Nothing to restore from The savepoint path does not exist, or the checkpoint was deleted when the job was cancelled Proof: execution.checkpointing.externalized-checkpoint-retention and the path itself. Fix: retain checkpoints on cancellation, or restore from the savepoint taken by stop.
Missing or changed operator uid Cannot map checkpoint/savepoint state for operator, or empty state after a restore that allowed non-restored state Proof: operator-uid in savepoint_metadata compared with the uid() calls in the new job. Fix: put the old uid() values back, or use setUidHash.
Incompatible state schema StateMigrationException naming the new and old state serializer Proof: the state type changed in a way POJO or Avro rules do not allow, or it is serialized with Kryo. Fix: revert the type change, or rewrite the savepoint with the State Processor API.
Parallelism above max parallelism The maximum parallelism of the restored state is lower than the configured parallelism Proof: operator-max-parallelism in savepoint_metadata. Fix: restore at or below the stored value.
Max parallelism changed The maximum parallelism of the latest checkpoint and the current maximum parallelism changed Proof: operator-max-parallelism in savepoint_metadata compared with setMaxParallelism or pipeline.max-parallelism. Fix: set the stored value, or rewrite the savepoint.
Wrong claim mode A savepoint that restored once is missing on the second restore Proof: execution.state-recovery.claim-mode on the first restore. Fix: restore with NO_CLAIM.
Drawn from the Apache Flink savepoints, upgrading, production readiness, schema evolution and checkpoints documentation, and the error messages in the Apache Flink source code, all linked on this page. Savepoints documentation

Causes and fixes by error

Nothing to restore from

Some failed restores have no state to restore at all. The checkpoints documentation says checkpoints are not retained by default and are deleted when a program is cancelled. A team that cancels a job and plans to restart from its last checkpoint finds nothing there unless retention was configured first.

execution.checkpointing.externalized-checkpoint-retention: RETAIN_ON_CANCELLATION

A retained checkpoint can be restored like a savepoint, by passing the path to its metadata file to run -s. The checkpoints versus savepoints documentation sets the limits. Unaligned checkpoints do not support arbitrary job upgrades or Flink minor version upgrades. Only canonical savepoints support a change of state backend or writing with the State Processor API. A canonical savepoint is the backend independent format that the savepoint command writes by default, as opposed to the native format. For a planned upgrade, a savepoint taken with stop is the snapshot to use.

bin/flink savepoint --type canonical :jobId

The savepoint path itself can also disappear. The savepoints documentation says the savepoint fails to trigger if neither a default directory nor a target directory is set, and the default comes from one key.

execution.checkpointing.savepoint-dir

Missing or changed operator uids

When the savepoint holds state for an operator that the new job does not have, the restore fails. The message comes from the Flink checkpoint restore code.

Cannot map checkpoint/savepoint state for operator <id> to the new program,
because the operator is not available in the new program.

This happens when a stateful operator was removed, and it also happens when its ID changed. The savepoints documentation says generated IDs depend on the structure of the program, so reordering stateful operators, or adding or removing stateless ones, can change the IDs of stateful operators that were never touched.

The error suggests a flag that makes the restore succeed. It has a short form and a configuration key.

--allowNonRestoredState
-n
execution.state-recovery.ignore-unclaimed-state

The flag skips state that cannot be mapped. If the operator was deleted on purpose, that is the right call. If its uid changed, the flag discards the state of an operator that still exists, and the operator starts as new. The savepoints documentation says a new operator is initialized without any state, and it warns that improper use of the flag can cause significant correctness issues. This is the usual path to a job that is RUNNING with empty state.

A Kafka source in that position is a specific case. The Flink Kafka connector documentation says the source does not rely on committed offsets for fault tolerance and commits them only to expose progress for monitoring. Its read position lives in Flink state, so a source that loses its state starts again from the starting offsets configured in its OffsetsInitializer. How offsets work on the Kafka side is covered in Kafka offsets.

The fix is to restore the old IDs, not to drop the state. Put the same uid() values back on the operators. If the original job never set uids, the upgrading documentation describes a low-level workaround that assigns the generated hash shown in the web UI or logs to the new operator.

setUidHash(String hash)

Table API and SQL jobs work differently, because the planner decides the operator topology. The upgrading documentation warns that any change to the query or to the Flink minor version can lead to state incompatibility. Its compatibility table also notes a configuration option for jobs affected by non-deterministic UIDs in Flink 1.15.0 and 1.15.1.

table.exec.uid.generation

Incompatible state schema

When a state type changed in a way Flink cannot migrate, the state backend throws a StateMigrationException. The wording below is from the heap keyed state backend, and the RocksDB keyed state backend uses the same phrase.

the new state serializer (...) must not be incompatible with the old state serializer (...)

The state schema evolution documentation limits evolution to POJO and Avro types.

  • POJO fields can be added or removed, but declared field types and the class name, including its namespace, cannot change.
  • Avro state can evolve as long as Avro’s own schema resolution rules accept the change, and the generated class cannot move to another namespace.
  • Keys cannot evolve at all.
  • Anything serialized with Kryo cannot evolve, including a POJO nested in a List inside a POJO.

Internal operator state counts too. The upgrading documentation lists window, join, reduce and built-in aggregation operators whose internal state type depends on their input or output type, so changing those types breaks the restore even when no user state changed.

Turning on state TTL is a version-dependent case. The Flink 1.20 state documentation says restoring state configured without TTL using a TTL-enabled descriptor, or the reverse, fails with a StateMigrationException. The Flink 2.3 state documentation says that from Flink 2.2.0 this migration is supported without restore-time errors.

When the schema change cannot be avoided, the State Processor API can read the old savepoint and write a new one. Its documentation lists changing state data types, adjusting max parallelism and reassigning operator uids among the changes it makes possible.

If the failure is a Kafka record that cannot be read with its registry schema, not a state type, the guide to Kafka deserialization errors covers that case.

Parallelism above max parallelism

Changing parallelism on restore is supported. The savepoints documentation says a program can be restored from a savepoint with a new parallelism. The limit is the max parallelism stored with the state, and going above it fails in the state assignment code.

The maximum parallelism (...) of the restored state is lower than the configured
parallelism (...). Please reduce the parallelism of the task to be lower or equal
to the maximum parallelism.

If no max parallelism was set, Flink chose one when the job first started. The production readiness checklist describes the default as 128 for a parallelism up to 128 and a rounded power of two above that. The rule in the Flink 2.3.0 source is more exact, with a floor of 128 and a ceiling of 2^15.

max parallelism = MIN(MAX(nextPowerOfTwo(parallelism + parallelism / 2), 128), 2^15)

A job first deployed at parallelism 4 therefore has a max parallelism of 128 and cannot be restored at 200. A job first deployed at parallelism 100 has 256. A second error in the same code fires when the max parallelism is set explicitly to a value different from the one stored in the snapshot.

The maximum parallelism (...) with which the latest checkpoint of the execution job
vertex ... has been taken and the current maximum parallelism (...) changed.
This is currently not supported.

Both have the same answer: restore with the stored max parallelism, which savepoint_metadata shows, and keep parallelism at or below it. Raising the ceiling for keyed state means rewriting the savepoint with the State Processor API.

The wrong claim mode

The claim mode decides who owns the snapshot files after a restore, and it explains most cases of a savepoint that worked once and then vanished. For an upgrade that might need a rollback, restore with NO_CLAIM so the savepoint stays intact as the fallback. The savepoints documentation describes the modes.

bin/flink run -s :savepointPath -claimMode :mode :runArgs
execution.state-recovery.claim-mode
  • NO_CLAIM is the default. Flink never deletes the snapshot files and several jobs can start from the same snapshot, and the original can be deleted manually once the first full checkpoint after the restore completes.
  • CLAIM hands the snapshot to Flink, which treats it like a checkpoint and may delete it once it is no longer needed for recovery. The documentation says it is not safe to delete it manually or to start two jobs from it, so a second restore from a claimed savepoint can find the files gone.
  • LEGACY is how Flink worked before 1.15, and the documentation calls it deprecated with ownership that is not well defined. The Flink 2.3 configuration reference lists only CLAIM and NO_CLAIM for the claim mode key.

A Flink version upgrade uses the same savepoint, with two extra checks. The upgrading documentation describes two ways to do it. An in-place upgrade stops all jobs with savepoints, shuts down the old cluster, upgrades it and restarts it. A shadow copy upgrade starts a new installation next to the old one, restores the jobs there and only then shuts the old cluster down, which leaves the old cluster as the rollback.

The upgrading documentation holds a compatibility table of which Flink version can resume a savepoint taken with which older version. In the Flink 2.3 documentation the table covers savepoints resumed with 1.17 to 1.20. It also notes that compatible refers to the internal savepoint format, not to SQL operators. For the move to the 2.x line, the Flink 2.0 release notes state that state compatibility is not guaranteed between 1.x and 2.x. A 1.x to 2.x upgrade is therefore something to test against a copy of the production savepoint before the production cutover, not something to assume.

The upgrading documentation lists two hard preconditions that make a migration fail. The first is RocksDB state that was checkpointed in semi-asynchronous mode, which is not supported for migration. The documentation says a job that used that mode can be switched to fully-asynchronous mode before the savepoint for the migration is taken. The second is that every savepoint file, including any written by the State Processor API, must be reachable from the new installation under the same absolute path.

Where Flex fits

This page is published by Factor House, which makes Flex, its Flink job management and monitoring product, so weigh this section with that in mind. The comparison of tools for self-managed Apache Flink sets Flex beside the Flink web UI and the REST API, including how each handles savepoints.

Upgrade actions. Flex documents these in its job documentation.

  • The Inspect view has quick actions to Stop, Cancel, take a Savepoint and trigger a Checkpoint.
  • The job overview table includes a Stoppable? field that indicates whether the job can be stopped gracefully through a savepoint, and a Last checkpoint entry.
  • Submitting a JAR sets the Parallelism, a Savepoint Path to resume from, a Claim mode and an Allow Non Restored State checkbox.
  • The Checkpoints tab lists every checkpoint with a Savepoint? column, so savepoints taken from Flex show beside periodic checkpoints, and it shows the average checkpoint duration.
  • The Topology tab shows a backpressure status of OK, LOW or HIGH and a backpressure level percentage for a selected task, which are the signals step 2 uses to judge how long a stop will take.

What stays in Flink. The job documentation does not say where the Stop action writes its savepoint or whether it drains, and it does not describe reading savepoint metadata, comparing operator uids or rewriting a savepoint. Those stay with the Flink CLI, the REST API and the State Processor API. For the meaning of each claim mode, the Flink savepoints documentation above is the reference.

Flink management tools use the REST API that Flink’s own dashboard uses, so a stop or savepoint from any tool is the same REST operation as the CLI’s. That endpoint accepts connections from any client by default and does not authenticate them, according to the Flink SSL setup documentation, which recommends an authenticating proxy in front of it. Anyone who can reach it can stop a job, cancel it without a savepoint or submit a new one, unless a proxy or a tool with its own access control sits in between.

The Flex role-based access control documentation describes granting or denying roles actions on Flink resources, including FLINK_SUBMIT, FLINK_JOB_EDIT and FLINK_JOB_TERMINATE, so those operations can be limited by role.

Flex ships as a single Docker container or JAR file, according to its GitHub repository, and its Flink provider documentation shows it connecting to a cluster through the Flink REST URL. To stop, savepoint and resubmit Flink jobs from one UI, Flex is pointed at that endpoint.

FAQ

Can Flink restore a savepoint with a different parallelism?

Yes, as long as the new parallelism is at or below the max parallelism stored with the state. The savepoints documentation says a program can be restored with a new parallelism. Above the stored max parallelism the restore fails, and the default when none was set is at least 128, so the stored value in savepoint_metadata is the one to read.

Is it safe to allow non-restored state on a restore?

Only when the unmapped state belongs to an operator that was removed on purpose. If an operator’s uid changed, the flag discards state for an operator that still exists, and the job starts that operator empty. Check the operator-uid column of savepoint_metadata against the new code before using it.

Where does the Kafka source resume after a savepoint restore?

From the position stored in the savepoint. The Flink Kafka connector documentation says the source does not rely on committed offsets for fault tolerance and commits them only to expose progress. A source whose state was not restored starts from the starting offsets configured in its OffsetsInitializer.

Can a Flink 1.x savepoint be restored on Flink 2.x?

Not with any guarantee. The Flink 2.0 release notes state that state compatibility is not guaranteed between 1.x and 2.x. Test the restore against a copy of the production savepoint on a separate 2.x cluster before cutting over.

Related reading