Incremental Checkpoint
also called Delta Snapshot, Incremental Snapshot
Persisting only the state files created since the previous snapshot, which decouples checkpoint cost from total state size while leaving recovery bounded by the full working set.
A job holding 2 TB of keyed state across 48 slots takes 11 minutes to complete a full checkpoint. The interval cannot be shorter than the checkpoint, so every failure costs at least 11 minutes of reprocessing, and the snapshot competes with record processing for network and CPU the whole time.
An incremental checkpoint uploads only the state files that are new since the last one. With a log-structured state backend the on-disk state is a set of immutable sorted files, so a checkpoint can reference files already uploaded and ship only the recent ones. Checkpoint duration stops tracking total state and starts tracking the write rate into state, which for most jobs is a small fraction. The 11 minutes becomes seconds.
The trap is the inference everyone draws next. Checkpoint duration and recovery time are different numbers, and only one of them improved.
Why it matters
Checkpoint duration sets the floor on the interval, and the interval sets how much work is redone after a failure and - with a transactional sink - how long before output becomes visible. Collapsing that duration is the difference between a job that can checkpoint every 30 seconds and one that cannot checkpoint faster than it fails.
Recovery, though, is a download. A replacement worker must fetch the files its state references and open the local store before processing resumes. Flink's own state-backend documentation says restoring from an incremental checkpoint may take longer than from a full one when network bandwidth is the bottleneck, because the referenced set is a chain of deltas and there is more to fetch; it is faster when CPU or IOPS dominate, since the local files do not have to be rebuilt from the canonical snapshot format. Restoring 2 TB at an aggregate 1 GB/s is about 35 minutes whatever the interval says.
Implementation patterns
- Pair it with the embedded key-value backend. Incremental snapshots depend on immutable on-disk files; a heap backend has nothing to send incrementally.
- Publish both numbers. Checkpoint duration and a measured restore time, from a drill, next to each other. Teams quote the first and plan with it.
- Cut recovery separately. Task-local state so a restart on the same node reads local disk; standby workers so rescheduling is not serial; less state per slot; and a tested savepoint path for planned changes.
- Watch the retained set. Deltas reference older files, so the storage held is larger than the logical state and old checkpoints cannot be deleted while a newer one references them. Set the number of retained checkpoints deliberately.
- Compaction shows up in the graph. When the backend compacts, large files are rewritten and the next "incremental" checkpoint is big. A sawtooth in checkpoint size is expected, not a defect.
Industry example
Flink introduced incremental checkpointing for its RocksDB backend and described the design in a January 2018 post on managing large state, framing it as the way to keep checkpoint cost proportional to change rather than to size. The documentation that accompanies it is unusually honest about the limit, stating plainly that restore may be slower or faster depending on which resource is the bottleneck. Both halves of that pairing matter in production: the feature is what makes terabyte-scale keyed state operable, and the recovery caveat is what teams forget until a node dies at peak.
Failure scenarios
- The planning error. Recovery budgeted from the checkpoint duration. The real restore is 30 minutes and the SLO assumed 1.
- Storage growth. Retained deltas plus retained checkpoints hold far more than the state size, and the bucket bill surprises someone.
- The long chain. Many small deltas between compactions mean many files to fetch, so restore degrades as the chain grows - worst at the moment a long-running job fails.
- Parallelism change. Rescaling redistributes key groups, so the restore is not a plain file copy and takes longer than any measured same-parallelism drill.
Trade-offs
| Choose | Gains | Pays |
|---|---|---|
| Incremental | checkpoints in seconds at terabyte scale; short intervals affordable | restore fetches a delta chain; more storage retained; sawtooth sizes |
| Full snapshots | one self-contained artefact; predictable restore | duration scales with total state; the interval has a hard floor |
When not to use it
If state is small - a few gigabytes per slot - a full snapshot already completes in seconds and incremental adds file bookkeeping, retained storage and a less predictable restore for no gain. Use full checkpoints and keep the mental model simple.
And if the state is large because nothing expires, bound the state before tuning the snapshot. A missing time-to-live or an unbounded join window will outgrow any checkpointing strategy; the tell is state size rising while the active key count is flat.
Interview question
Q: Your job's checkpoint duration dropped from 90 seconds to 6 after a configuration change, and the on-call engineer wants to set the recovery-time objective to one minute. What do you tell them?
What a strong answer covers: that the change altered upload volume, not download volume; that restore is bounded by fetching the referenced file set and opening the local store, and can be slower with deltas when bandwidth is the constraint; the arithmetic of state size over achievable download throughput; that the RTO must come from a measured drill including a rescale case; and the levers that genuinely reduce it - local recovery, standbys, smaller per-slot state.
Quick check
Quiz: Incremental checkpointing cut checkpoint time by 15x. What happened to recovery time? Little, and it can get worse - restore still fetches the whole referenced state, now as a chain of deltas.
Flashcard: What does an incremental checkpoint make proportional to what? - Checkpoint cost to the state change rate instead of total state size; recovery stays proportional to total state.