Turning on Flink checkpointing takes one line. Designing it takes numbers: how long a recovery may take, how stale your sink output may be, how big the state is, and how much storage bandwidth you can spend. Most checkpoint incidents come from defaults chosen by nobody, such as a tolerance of zero failures that turns one slow upload into a job restart.
This article treats checkpointing as a design exercise. It starts with what happens during a checkpoint and where the time goes, derives the interval from a recovery budget and a freshness requirement, then covers the timeout, the minimum pause, failure tolerance, unaligned checkpoints and the metrics that tell you whether the design is holding. Barrier theory and exactly-once proofs are covered in exactly-once stream processing; this page is about choosing the settings. Configuration keys use the Flink 2.x names; older 1.x releases used names such as state.checkpoints.dir and state.backend.incremental for the same ideas.
What one checkpoint does
The JobManager's checkpoint coordinator triggers a checkpoint by asking every source to inject a barrier into its output. Sources record their position, such as Kafka offsets, as part of the checkpoint. The barrier flows downstream with the records. An operator with several inputs waits until the barrier has arrived on every input, the alignment, so its state reflects exactly the records before the barrier. It then takes a short synchronous snapshot and forwards the barrier, and its state files are uploaded to durable storage asynchronously while processing continues. When every task has acknowledged, the checkpoint is complete, and the coordinator notifies operators, which is when transactional sinks commit the output they pre-committed for that checkpoint.
That sequence has four phases you can measure separately: start delay (how long the barrier takes to reach a task, which grows with backpressure because the barrier queues behind buffered records), alignment, the synchronous snapshot, and the asynchronous upload. A good design knows which phase dominates for its job, because each has a different fix.
The interval comes from two budgets
The interval trades steady-state cost (storage bandwidth, CPU, barrier traffic) against failure cost (replay after a failure and, with transactional sinks, staler output). Derive it from two requirements.
Recovery-time objective. After a failure the job restarts, restores state from the last completed checkpoint and re-reads all input since that checkpoint started. In the worst case that is one interval plus one checkpoint duration of input. Replay does not run against a frozen source: new data keeps arriving, so the backlog drains only at the difference between the job's maximum throughput and the ingest rate. A job running near its capacity cannot catch up quickly no matter how often it checkpoints.
Sink freshness. Exactly-once sinks such as Kafka transactions or two-phase file sinks make output visible only when a checkpoint completes. Readers using read-committed isolation therefore see data roughly one interval plus one checkpoint duration late. If downstream consumers need output within two minutes, the interval cannot be five minutes, whatever the recovery budget says.
def checkpoint_budget(rto_s, restart_s, restore_s, replay_rate, ingest_rate,
ckpt_duration_s, freshness_s=None):
"""Largest checkpoint interval that meets a recovery-time objective.
After a failure the job restarts, restores the last completed checkpoint and
re-reads everything since that checkpoint started. In the worst case that is
one interval plus one checkpoint duration of input.
Replay drains at (replay_rate - ingest_rate) because new data keeps arriving.
"""
catch_up_speed = replay_rate - ingest_rate
if catch_up_speed <= 0:
raise ValueError("job cannot catch up: replay must exceed ingest")
budget = rto_s - restart_s - restore_s
# backlog = ingest_rate * (interval + duration); time = backlog / catch_up_speed
interval = budget * catch_up_speed / ingest_rate - ckpt_duration_s
if freshness_s is not None: # transactional sink visibility
interval = min(interval, freshness_s - ckpt_duration_s)
return max(interval, 0)
# Worked example from the article
print(checkpoint_budget(rto_s=300, restart_s=30, restore_s=90,
replay_rate=200_000, ingest_rate=50_000,
ckpt_duration_s=25, freshness_s=120))
Worked example: an order enrichment job
A job reads 50,000 events per second from Kafka, joins them with customer state and writes to a transactional Kafka sink. Keyed state is 400 GB in RocksDB across 40 TaskManagers. At full speed the job can process 200,000 events per second. The business wants recovery within 5 minutes and enriched events visible to consumers within 2 minutes.
Measure first. A full checkpoint of 400 GB would upload 10 GB per TaskManager every time; with incremental checkpoints enabled the job uploads about 6 GB of changed files per checkpoint in total, and checkpoints complete in about 25 seconds. A restart takes about 30 seconds to reschedule, and restoring state takes about 90 seconds, because the TaskManagers download their share of the 400 GB.
Recovery budget: 300 seconds minus 30 for restart and 90 for restore leaves 180 seconds to replay. The job catches up at 200,000 minus 50,000, or 150,000 events per second, so 180 seconds absorbs a backlog of 27 million events, which is 540 seconds of input. Subtract the 25-second checkpoint duration and the recovery budget allows an interval of up to 515 seconds. Freshness is stricter: 120 seconds minus 25 gives 95 seconds. The design is therefore an interval of 60 seconds, which leaves margin for slower checkpoints, with a minimum pause of 20 seconds so at least a third of each cycle is guaranteed to be free of checkpoint work even when a checkpoint runs long.
Two conclusions follow. Recovery here is dominated by restore, not replay, so local recovery (a copy of state on the TaskManager's disk for restarts on the same machine) buys more than a shorter interval. And freshness, not failure handling, sets the interval, as is common with transactional sinks.
Timeout, minimum pause and failure tolerance
The timeout (execution.checkpointing.timeout, default 10 minutes) aborts a checkpoint that has not completed. Set it to several times the normal end-to-end duration. For a job whose checkpoints take 25 seconds, the default lets a stuck checkpoint occupy the only slot for 10 minutes while sink output stops committing.
The minimum pause (execution.checkpointing.min-pause, default 0) is the gap between the end of one checkpoint and the start of the next. Without it, a checkpoint that takes longer than the interval is followed immediately by another, and a struggling job spends all its time checkpointing. With max-concurrent-checkpoints left at its default of 1, the pause is what guarantees processing time.
Tolerable failures (execution.checkpointing.tolerable-failed-checkpoints, default 0) is how many consecutive checkpoint failures the job accepts before it fails and restarts. Zero means one slow upload or one storage throttling event restarts the job, which then has to restore and replay, putting more load on the same storage. A small number such as 2 or 3 rides out transient problems; the cost is that sinks commit less often during the failures, and you must alert on failed checkpoints so tolerance does not hide a persistent problem.
Retention. By default checkpoints are deleted when the job is cancelled (NO_EXTERNALIZED_CHECKPOINTS), and only the latest one is kept. Set retention on cancellation and keep two checkpoints, so an operator who cancels a job by mistake, or finds the latest checkpoint corrupt, still has something to restore from. For planned upgrades use savepoints, described in Flink savepoints.
# Flink 2.x configuration keys (config.yaml or per-job configuration)
execution.checkpointing.interval: 60 s
execution.checkpointing.min-pause: 20 s # guaranteed processing gap between checkpoints
execution.checkpointing.timeout: 5 min # default is 10 min
execution.checkpointing.max-concurrent-checkpoints: 1
execution.checkpointing.tolerable-failed-checkpoints: 3 # default 0: first failure restarts the job
execution.checkpointing.mode: EXACTLY_ONCE
execution.checkpointing.dir: s3://flink-ckpt/orders-enricher
execution.checkpointing.incremental: true # default false
execution.checkpointing.num-retained: 2 # default 1
execution.checkpointing.externalized-checkpoint-retention: RETAIN_ON_CANCELLATION
# Backpressure handling: start aligned, switch to unaligned if alignment drags on
execution.checkpointing.unaligned.enabled: true
execution.checkpointing.aligned-checkpoint-timeout: 30 s
taskmanager.network.memory.buffer-debloat.enabled: true
# Faster restarts on the same TaskManager (keeps a local copy of state)
execution.checkpointing.local-backup.enabled: true
Backpressure, alignment and unaligned checkpoints
Under backpressure, buffers between operators fill up and barriers queue behind the buffered records, so start delay and alignment grow until checkpoints time out exactly when the job is under the most stress. Two features address this, and they work best together.
Unaligned checkpoints let a barrier overtake buffered records. Instead of waiting, the operator snapshots its state together with the in-flight records in its buffers, and those records become part of the checkpoint. Checkpoint duration stops depending on how much data is queued. The costs are documented: more data written to checkpoint storage, so they are a poor fit when storage I/O is the bottleneck; no concurrent unaligned checkpoints, and they cannot run at the same time as a savepoint; and on recovery Flink generates watermarks after restoring in-flight data, unlike aligned checkpoints, which matters for operators that assume watermark ordering.
The aligned checkpoint timeout makes this a hybrid: every checkpoint starts aligned, and if it has not completed within the configured duration it switches to unaligned. Normal checkpoints stay small, and only checkpoints taken under backpressure pay the extra I/O. Buffer debloating (taskmanager.network.memory.buffer-debloat.enabled) attacks the cause by sizing network buffers to the throughput so less data is queued in the first place, which shortens aligned checkpoints and keeps unaligned ones smaller. Backpressure itself is explained in backpressure architecture.
Making each checkpoint cheap
The synchronous and asynchronous phases depend on the state backend. With the RocksDB backend, incremental checkpoints (execution.checkpointing.incremental, default false) upload only files created since the previous checkpoint, which turns a 400 GB upload into a few gigabytes for most jobs. The price is that a checkpoint references files from earlier ones, so restores read more files and storage cleanup is shared between checkpoints. Full, incremental and changelog-based checkpoints are compared in detail in Flink state backends compared.
Beyond the backend, expire state you no longer need with state time-to-live, because every byte of state must be uploaded and restored, and keep checkpoint storage in the same region as the cluster.
Reading the checkpoint metrics
The checkpoint tab in the web UI and the REST API break each checkpoint into the phases above, per task. Read them as a diagnosis rather than a single duration.
- High start delay on downstream tasks means barriers are stuck behind buffered data: backpressure. Look at the slowest operator, then consider debloating or unaligned checkpoints.
- High alignment duration on one multi-input operator means one input is much slower than the others, often from data skew on one key or one partition.
- High synchronous duration points at the operator itself, for example a heap state backend copying large structures.
- High asynchronous duration means uploads are slow: storage throttling, too much changed state, or full rather than incremental checkpoints.
- Checkpointed size close to full state size on every checkpoint means incremental checkpointing is off or ineffective.
import requests
def checkpoint_health(rest, job_id, interval_s, timeout_s):
"""Read the checkpoint summary from the Flink REST API and flag risks.
Field names follow recent Flink releases; check your version's REST docs."""
s = requests.get(f"{rest}/jobs/{job_id}/checkpoints", timeout=5).json()
counts, latest = s["counts"], s["latest"]["completed"] or {}
dur_s = latest.get("end_to_end_duration", 0) / 1000
issues = []
if counts["failed"] > 0:
issues.append(f"{counts['failed']} failed checkpoints since job start")
if dur_s > 0.5 * timeout_s:
issues.append(f"duration {dur_s:.0f}s is over half the {timeout_s}s timeout")
if dur_s > interval_s:
issues.append("checkpoints take longer than the interval: back to back, no idle gap")
return issues
Failure modes
- Timeouts only under load. Checkpoints are fine at night and expire at peak, because start delay grows with backpressure. Use the aligned checkpoint timeout and debloating, and fix the slow operator.
- Restart loops from zero tolerance. A storage hiccup fails a checkpoint, the job restarts, the restore hammers the same storage, and the next checkpoint fails too. Allow a few failures and alert on them.
- Kafka transactions expiring. A transactional sink's open transaction must outlive the longest checkpoint plus any restart. If the producer's transaction timeout is shorter, the broker aborts it and output is lost or the job fails; raise it within the broker's maximum, as the Kafka connector documentation advises.
- Lost recovery point. A cancelled job with default retention deletes its checkpoints. Retain on cancellation.
- Silent state growth. State that never expires slowly pushes restore time past the recovery objective. Track state size over time.
Trade-offs
| Setting | Lower value | Higher value |
|---|---|---|
| Interval | Less replay, fresher sink output, more storage and CPU | Cheaper steady state, longer recovery, staler output |
| Timeout | Stuck checkpoints noticed early, more spurious failures under load | Fewer failures, problems hidden longer |
| Min pause | More checkpoints when they run long | Guaranteed processing time, staler output |
| Tolerable failures | Fails fast, restart storms possible | Rides out blips, longer gaps between commits |
| Unaligned | Off: smaller checkpoints, slow under backpressure | On: steady duration, more checkpoint I/O |
What to do next
- Write down the recovery-time objective and sink freshness requirement for each job, then compute the largest allowed interval with the budget function above.
- Measure restart time, restore time, checkpoint duration and maximum throughput in a staging failover test instead of guessing them.
- Set the timeout to a few times normal duration, a minimum pause, two or three tolerable failures and retention on cancellation.
- Enable incremental checkpoints for large RocksDB state and local recovery if restore dominates recovery time.
- Turn on buffer debloating and an aligned checkpoint timeout for jobs that see backpressure, after checking the watermark caveat for your operators.
- Alert on failed checkpoints, duration over half the timeout and steady growth in checkpointed size. For state design, read Flink state in depth.