Skip to content

How to fix Flink checkpoints that fail or time out

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

A Flink job with failing checkpoints usually shows it in one of two ways. Either the Failed counter on the Flink web UI’s Checkpoints tab keeps rising while the job status cycles between RESTARTING and RUNNING, or nothing fails outright but each checkpoint takes longer than the last until one expires with the message “Checkpoint expired before completing” (the CheckpointFailureReason source defines it). In that second case the job can also spend so long checkpointing that it falls behind its source.

A Flink checkpoint is a consistent snapshot of every operator’s state and source position, taken while the job runs, that the job restores from after a failure.

A checkpoint can be slow at three stages on each subtask: the barrier reaching it, the synchronous snapshot, or the asynchronous write to checkpoint storage. A timeout at any stage restarts the job by default, because the tolerable checkpoint failure number is 0. The checkpoint monitoring documentation lists the statistics that show each stage.

This page is about checkpoints inside a Flink job. If the symptom is a Kafka consumer group falling behind and the consumer is not a Flink job, the guide to Kafka consumer lag is the better starting point. For how checkpoints fit into Flink’s wider fault tolerance model, see the complete Flink guide.

The sections below follow this order:

  1. Check the job for backpressure.
  2. Find which stage is slow on the slowest subtask.
  3. Change one setting and watch the metric that showed the problem.
  4. Raise the tolerable failure number only as a stopgap while the cause is found.

Rule out backpressure first

Checkpoint barriers flow through the job with the data, so a task that is backpressured holds its barrier behind a queue of records. Flink’s guide to checkpointing under backpressure says checkpoint time in this case is dominated by barrier propagation, which shows up as high alignment time and start delay.

Check the job’s backpressure before touching any checkpoint setting. Each subtask splits every second into time spent backpressured, idle and busy, and the three add up to about 1000 ms.

backPressuredTimeMsPerSecond
idleTimeMsPerSecond
busyTimeMsPerSecond

The backpressure monitoring documentation rates a task OK from 0 to 10% backpressured, LOW above 10 to 50%, and HIGH above 50%. In the web UI’s job graph a backpressured task is black and a busy task is red. A black task is a victim of the slowdown. The Flink post on finding the source of backpressure says the busiest (red) task downstream of the backpressured tasks will most likely be the bottleneck. A subtask that sits near 1000 ms idle is starved of input and is not slow.

If backpressure is HIGH, fix that first and check the checkpoints again afterwards. How to find the bottleneck behind Flink backpressure walks through that diagnosis.

Find the stage that is slow

The fix depends on which stage is slow, so the comparison table below comes first. The Flink web UI’s checkpoint History tab lists each checkpoint, and expanding one shows per-subtask statistics. The REST API returns the same figures.

F1 Backpressure-driven and storage-driven slow checkpoints, side by side
Backpressure-driven Storage-driven
What you see first Duration climbs with input rate, and the job falls behind its source. Duration climbs with state size, even at a steady input rate.
Metric that proves it High Start Delay and Alignment Duration, normal Sync and Async Duration. High Async Duration, normal Start Delay and Alignment Duration.
Usual cause A slow operator or sink holds barriers behind queued records. Full snapshots of large RocksDB state, or slow checkpoint storage.
Fix Remove the backpressure source, then buffer debloating or unaligned checkpoints. Incremental checkpoints, a minimum pause, FileSystemCheckpointStorage.
Risk of the fix Unaligned checkpoints add state-storage I/O. Incremental checkpoints report the delta, so a smaller figure is not smaller state.
Drawn from the Apache Flink stable documentation on checkpoint monitoring, checkpointing under backpressure and large state tuning, which are linked in the sections below. Checkpoint monitoring documentation

Read the stage on the slowest subtask

End to End Duration is set by the last subtask to acknowledge, so one slow subtask fails the whole checkpoint while the job-level average looks fine. Sort the subtasks by duration, then read the stages for the slowest one.

Barrier delay and alignment

Start Delay is how long the first barrier took to reach the subtask, and Alignment Duration is the time between processing the first and the last barrier. High values in either point back to backpressure.

Start Delay
Alignment Duration

Synchronous snapshot

The synchronous part snapshots operator state and blocks all other activity on the subtask, including processing records and firing timers. The large state tuning guide notes that heap-based timers may increase checkpointing times, so check timers when this stage is the slow one.

Sync Duration

Asynchronous write

The asynchronous part includes the time to write the checkpoint to the selected filesystem while the task keeps processing. A long Async Duration points to large snapshots, slow storage, or both. Compare Checkpointed Data Size across checkpoints to see whether state is growing.

Async Duration
Checkpointed Data Size

Why a failed checkpoint restarts the job

By default Flink tolerates no checkpoint failures. The Flink checkpointing documentation sets the tolerable checkpoint failure number to 0, so the first checkpoint that fails in the asynchronous phase, hits an IOException on the JobManager, or expires on timeout fails the job over. A failure in the synchronous phase always fails over the task, whatever the setting. A checkpoint timeout therefore often looks like a restart loop to the operator, not like a timeout.

The defaults below are listed in the Flink configuration reference.

execution.checkpointing.tolerable-failed-checkpoints: 0   # default
execution.checkpointing.timeout: 10 min                   # default

Raising the tolerable number is the fourth step above and means fewer restarts caused by checkpoints, but it also widens the gap to the last completed checkpoint, so a later failure replays more records. It is a stopgap while the real cause is found.

The checkpoint counters can also mislead. The checkpoint monitoring documentation notes that the Triggered, In Progress, Completed, Failed and Restored counts do not survive a JobManager loss and are reset if the JobManager fails over, and that Restored also counts restarts.

Fixes, one change at a time

Pick the metric that showed the problem, make one of the changes below, and watch that metric across several checkpoints before making the next one.

Version note. The keys below are the names in the Flink 2.x stable configuration reference. The Flink 1.20 release notes describe the move of checkpointing options to the execution.checkpointing.* prefix. On Flink 1.19 and earlier, the 1.19 documentation lists incremental checkpoints as state.backend.incremental, and checkpoint storage as state.checkpoint-storage with state.checkpoints.dir.

Set a minimum pause between checkpoints

If checkpoints regularly take longer than the interval, the next one starts as soon as the last finishes, and the large state tuning guide describes the job as constantly taking checkpoints. A minimum pause guarantees processing time between them. Flink does not allow a minimum pause together with concurrent checkpoints.

execution.checkpointing.interval: 60 s
execution.checkpointing.min-pause: 30 s

Turn on incremental checkpoints for large RocksDB state

When Async Duration grows with state size on RocksDB, the large state tuning guide says incremental checkpoints should be one of the first considerations. Each checkpoint then uploads only what changed since the last one.

execution.checkpointing.incremental: true

After the change, the checkpointing documentation notes that the UI and REST API report the delta, not the full state. A smaller Checkpointed Data Size after switching on incremental checkpoints does not mean the state shrank.

Reduce in-flight data with buffer debloating

Where alignment and start delay dominate, the lasting fix is the slow operator or sink found above: more parallelism, a faster external call, or less skew. While that is in progress, buffer debloating automatically reduces the in-flight data buffered between tasks. Flink’s guide to checkpointing under backpressure says buffer debloating works with aligned and unaligned checkpoints and has the most visible effect on aligned ones.

taskmanager.network.memory.buffer-debloat.enabled: true

Try unaligned checkpoints for backpressure-driven timeouts

Unaligned checkpoints let barriers overtake queued records and store those records as part of the checkpoint. They work only with exactly-once mode and one concurrent checkpoint, and the same guide warns that they add state-storage I/O, so they make a storage-bound job slower. An aligned checkpoint timeout starts each checkpoint aligned and switches to unaligned only when start delay passes the threshold.

execution.checkpointing.unaligned.enabled: true
execution.checkpointing.aligned-checkpoint-timeout: 30 s

Change checkpoint storage or the state backend last

A job that checkpoints fine in development and fails in production is often still on JobManager checkpoint storage. The checkpoints documentation describes JobManagerCheckpointStorage as holding state in the JobManager heap with a default limit of 5 MB per state, meant for development and small state, and encourages FileSystemCheckpointStorage for high availability setups.

execution.checkpointing.storage: filesystem
execution.checkpointing.dir: s3://bucket/flink/checkpoints

For the state itself, the state backends documentation says the HashMapStateBackend holds state as objects on the Java heap, and encourages the EmbeddedRocksDBStateBackend for jobs with very large state, long windows and large key or value states. Switching backend changes checkpoint cost as well as memory use, so plan it as a migration through a savepoint.

When checkpoints keep failing anyway

Large jobs fail checkpoints more often

A checkpoint completes only when every subtask has acknowledged it, and the checkpoint monitoring documentation lists the acknowledged subtasks and the end to end duration for each one, so a large job with many tasks has many chances to fail each checkpoint. The StreamShield paper from ByteDance describes region checkpointing, an internal change to ByteDance’s own Flink platform, and reports checkpoint success rising from 53.9% to 93.5% in its production cluster (Fig. 8). The paper presents it as a ByteDance system, and the Apache Flink checkpointing documentation lists no such option, so the figure shows the scale of the problem and is not a tuning result to reproduce.

Why Kafka lag stalls while checkpoints fail

Checkpoint failures also leave a mark in Kafka. The Flink Kafka source commits the current consuming offset when checkpoints are completed, so while checkpoints fail, offsets committed to Kafka stall or move in steps even though the job is still reading. Tools that read committed offsets show that as lag. The source’s own pending records metric is the live signal.

pendingRecords

A consumer group outside Flink can also stall for another reason, covered in what Kafka rebalancing is.

Restoring from the right checkpoint

When the job has to be redeployed, restoring from the latest valid checkpoint is where operator error creeps in. The Flink Kubernetes Operator job management documentation describes initialSavepointPath in the job manifest for starting a job from a chosen savepoint or checkpoint, which is the field to set deliberately during an incident.

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. Flex is one of several ways to watch checkpoints, and the comparison of tools for self-managed Apache Flink sets it beside the Flink web UI, the REST API and metrics stacks. Teams on Ververica Platform can start from the Ververica Platform tool comparison.

Diagnosis. Flex documents these views for each job in its job inspection documentation.

  • The Checkpoints tab counts Triggered, Completed, In progress, Restored and Failed checkpoints, shows Avg checkpoint size and Avg duration, and lists a history table whose Status is COMPLETED, IN_PROGRESS or FAILED.
  • The Configuration view on that tab shows the job’s checkpoint interval, timeout, minimum pause and storage as Flink reports them.
  • The Events tab is a reverse-chronological lifecycle log with lines such as “Triggered CHECKPOINT 91 with status FAILED” and “Flink job entered state RESTARTING”, so the failure and the restart sit in one place.
  • The Overview tab shows backpressure Status (OK, LOW or HIGH) and a backpressure level percentage.

Flex reads the Flink REST API through the endpoint set in its Flink cluster configuration. Support for Ververica Platform is in early access, as the Ververica Platform provider documentation states.

What stays in Flink. The job documentation does not describe the following, so they stay in the Apache Flink web UI or REST API.

  • Per-subtask start delay, alignment duration, and sync and async durations.
  • The failure reason for a failed checkpoint.
  • Alerting on checkpoint failures. The Flink primer on monitoring Flink applications lists the number of restarts, failed checkpoints and source lag among the first things to monitor, and calls successful checkpointing a strong indicator of the general health of an application. Flink exports checkpoint metrics such as numberOfFailedCheckpoints through its metric reporters to an external system, where alerts can be set. Flex documents no checkpoint alerts.

Fixes. Flex documents these actions.

  • The Inspect view has quick actions to Stop, Cancel, take a Savepoint and trigger a Checkpoint.
  • Submitting a job lets the operator choose a savepoint path, a claim mode, a restore mode and the Allow Non Restored State option, which covers the redeploy when failed checkpoints force one. Flink’s savepoint documentation explains what those options mean.

To see checkpoint history, Events and backpressure for each job in one place, try Flex.

FAQ

What is the default checkpoint timeout in Flink?

Ten minutes, as the Flink configuration reference lists for execution.checkpointing.timeout. A checkpoint that has not completed by then is discarded. With the default tolerable checkpoint failure number of 0, that expiry also fails the job over, so a timeout often shows up as a restart.

Why do Flink checkpoints get slower over time?

Two documented causes are growing state and backpressure. Compare Checkpointed Data Size across checkpoints for growing state, and Start Delay and Alignment Duration on the slowest subtask for backpressure. Without incremental checkpoints, each RocksDB checkpoint uploads the full state, so its size tracks state size, which is why the large state tuning guide names incremental checkpoints as an early fix.

When are unaligned checkpoints the right fix?

Only when checkpoints are slow because of backpressure and the job is not already limited by checkpoint storage throughput. Unaligned checkpoints need exactly-once mode with one concurrent checkpoint and add state-storage I/O. An aligned checkpoint timeout limits that cost to the checkpoints that need it.

What is the difference between a checkpoint and a savepoint?

Checkpoints are taken and managed by Flink for recovery from failures. Savepoints are triggered by an operator for planned work such as upgrades, rescaling or a backend migration. The Flink documentation compares the difference to backups versus recovery logs.

Related reading