skip to content

A Flink job's checkpoints time out under heavy backpressure — how do unaligned checkpoints help?

level: seniorimportance: nice to knowfreq 45%

answer

  1. the barrier is stuck in a queue
  2. let the marker jump the line
  3. then you must save the queue too
  4. bigger snapshots, slower restore
  5. the job is still too slow

basics

~20 s

Under backpressure a barrier crawls through full network queues and alignment stalls the snapshot. Unaligned checkpoints let the barrier overtake queued records and persist the in-flight buffers as part of the snapshot, so checkpoint duration stops depending on throughput.

solid answer

~50 s

With aligned checkpointing, a barrier can only advance as fast as the records queued ahead of it, and a multi-input task then idles while the slowest channel catches up — so when the job is backpressured, checkpoint duration tracks the queue drain time and eventually exceeds the timeout. Unaligned checkpointing breaks that coupling: the barrier is moved to the front of the output queue immediately, and the records it jumped over are written into the checkpoint as **in-flight data**. Checkpoint duration then depends on state and buffer size, not on throughput. In Flink 2.3 you enable it with `enableUnalignedCheckpoints()` on the `CheckpointConfig` (key `execution.checkpointing.unaligned.enabled`), usually with `setAlignedCheckpointTimeout(...)` so checkpoints start aligned and switch only when a barrier is already running late. The costs are real: bigger snapshots, extra I/O, longer recovery because the buffers must be re-injected — and it treats the symptom, not the backpressure itself.

code

java · 8 lines
java
import org.apache.flink.core.execution.CheckpointingMode;

CheckpointConfig cfg = env.getCheckpointConfig();
cfg.setCheckpointingConsistencyMode(CheckpointingMode.EXACTLY_ONCE);
cfg.setMaxConcurrentCheckpoints(1);
cfg.enableUnalignedCheckpoints();
// start aligned; switch once a barrier is 30s late reaching a task
cfg.setAlignedCheckpointTimeout(Duration.ofSeconds(30));

go deeper

for a junior

Not expected at this level. It is enough to know that alignment means a task waits for barriers on all its inputs, and that waiting costs time.

for a middle

Be able to say why a barrier moves slowly through a backpressured job and what alignment buffers while it waits. Knowing that a mode exists which persists in-flight data instead is a good sign at this tier.

for a senior

Demonstrate the diagnosis first: split checkpoint duration into sync, async and alignment before choosing a remedy. Explain what unaligned checkpoints store, what they cost in size and recovery, and that they do not remove the backpressure.

for a principal

Frame it as a fault-tolerance safety valve, not a tuning knob: decide when a fleet should default to timeout-based switching, what it does to checkpoint storage cost and recovery objectives, and when the real answer is to fix capacity or skew.

## Why backpressure kills aligned checkpoints A barrier is an in-band element of a channel: it cannot overtake the records ahead of it. When a job is backpressured, every network queue between tasks is full, so a barrier injected at the source has to wait behind thousands of buffered records at every hop. Meanwhile any task with multiple inputs is *aligning*: it has blocked the channel that already delivered the barrier and is idling on it while the congested channel slowly drains. The result is that checkpoint duration becomes a function of queue depth and throughput. As backpressure worsens, checkpoints take longer; when they exceed the configured timeout (ten minutes by default) they are abandoned, the recovery point ages, and a failure now means replaying far more input than the interval suggests. With the default of zero tolerable checkpoint failures, each expiry also fails the whole job over, which under sustained backpressure can turn into a restart loop. Worse, aborted checkpoints mean a transactional sink never commits, so an exactly-once pipeline stops producing visible output entirely. ## What unaligned checkpointing changes Unaligned checkpointing lets the barrier jump the queue. When a barrier arrives at a task's input, the task immediately moves it to the front of its output buffers and forwards it downstream, then snapshots two things instead of one: - the operator's state, as usual; - the **in-flight records** — the contents of the input and output network buffers that the barrier overtook. Because the snapshot now contains the buffered data as well as the state, the cut is still consistent: on recovery Flink restores the state *and* re-injects those buffers into the channels, reproducing exactly the situation at the moment of the cut. This is much closer to the textbook Chandy-Lamport algorithm, which records channel contents rather than draining them. Checkpoint duration is then bounded by how fast state and buffers can be written to durable storage — a function of size, not of how congested the pipeline is. A job that could not complete a checkpoint in ten minutes under backpressure will often complete one in seconds. ## Turning it on sensibly ```java CheckpointConfig cfg = env.getCheckpointConfig(); cfg.setCheckpointingConsistencyMode(CheckpointingMode.EXACTLY_ONCE); cfg.setMaxConcurrentCheckpoints(1); cfg.enableUnalignedCheckpoints(); cfg.setAlignedCheckpointTimeout(Duration.ofSeconds(30)); ``` The timeout is the part worth understanding. With it set, a checkpoint *starts* aligned — cheap, small — and a task switches it to unaligned once the barrier's start delay, the time since the checkpoint was triggered, exceeds that budget. Steady-state checkpoints stay small, and only the ones that hit trouble pay the unaligned price. Zero, the default, means always unaligned. Two constraints follow from the design. Unaligned checkpointing only applies in exactly-once mode — under at-least-once Flink ignores the setting, since that mode already skips blocking and accepts the duplicates instead — and Flink documents it as requiring a single concurrent checkpoint, since overlapping snapshots of the same in-flight buffers have no coherent meaning. Two things stay aligned regardless: savepoints, and, since 2.0, the connections inside a Sink V2 sink's commit topology (writer to committer), because committables must reach the committer before the checkpoint-complete notification. ## What it costs - **Snapshot size.** Every persisted buffer is bytes written per checkpoint that an aligned checkpoint would not have written. On a wide, heavily buffered job this can be substantial, and it repeats every interval if you leave the timeout at zero. - **I/O and storage.** More writes to checkpoint storage, and larger retained snapshots. - **Recovery time.** Restoring means re-injecting the buffers, not just loading state. Flink 2.3 adds two experimental options, both off by default, to shorten that window: `execution.checkpointing.unaligned.recover-output-on-downstream.enabled`, and `execution.checkpointing.during-recovery.enabled`, which lets the job checkpoint while it is still restoring in-flight data. - **It does not fix the backpressure.** The job is still too slow for its input; you have only stopped that fact from destroying your fault tolerance. Treating unaligned checkpoints as a performance fix is the classic misreading. ## The diagnostic path Before reaching for it, read the checkpoint statistics and split duration into its synchronous, asynchronous and alignment parts: - **Alignment dominant** — a genuinely lagging channel. Look for skew across upstream subtasks, or one slow branch of the graph. Unaligned checkpointing helps here, and so does fixing the skew. - **Asynchronous dominant** — the state write is slow. That is state size or checkpoint storage throughput, and unaligned checkpointing will make it worse, not better. - **Synchronous dominant** — the state snapshot itself is expensive. Again, unaligned does not help. So unaligned checkpointing is a targeted answer to one specific shape of problem, which is exactly why it is a differentiator question rather than a core one: plenty of strong senior engineers have never needed it, and reaching for it reflexively is a worse answer than not knowing it at all.

  • Why does an unaligned checkpoint have to persist the in-flight network buffers?
    Because the barrier overtook those records, the snapshot's state does not reflect them, yet the sources have already moved past them. Storing the buffers restores the missing piece: on recovery Flink loads the state and re-injects the buffers into the channels, so the restored job is in exactly the situation the cut described. Without them, those records would be silently lost.
  • When would enabling unaligned checkpoints make things worse?
    When alignment was never the problem. If the duration breakdown shows the asynchronous phase dominating, the bottleneck is state size or checkpoint storage throughput, and adding buffer contents to every snapshot increases exactly the thing that is already slow. It also lengthens recovery and inflates storage cost, so it is wasted on a job whose checkpoints are slow for any other reason.
  • Why start aligned and switch on a timeout rather than running unaligned always?
    Because steady-state checkpoints are small and cheap when aligned, and most of them complete quickly. The timeout keeps that cheap path for the normal case and pays the unaligned price only for the attempts that are already in trouble, which is usually a transient backpressure spike. Always-unaligned means every checkpoint carries the buffer cost even when there was nothing to fix.

An aligned checkpoint is a photographer waiting for every queue at the counter to clear before taking the picture; an unaligned one photographs the queues themselves and puts the people back where they stood afterwards.

saying these in an interview costs you the question

  • Says unaligned checkpoints relieve backpressure itself
  • Thinks unaligned means the barrier is simply skipped
  • Believes unaligned checkpoints make snapshots smaller
  • Enables it without reading the alignment portion of the duration
  • Assumes recovery is unchanged when in-flight data is persisted

context