In Flink, how do checkpoint barriers produce a consistent snapshot across parallel operators?
answer
- a marker rides with the records
- never overtakes the data it follows
- one input channel arrives before the other
- the task waits, buffering the fast side
- every subtask must acknowledge
basics
~20 sBarriers are markers the source tasks inject into the record stream. Each task snapshots its state at the moment the barrier passes it, so every task's snapshot reflects the same logical cut of the stream even though no task ever stops the world.
solid answer
~40 sOn each interval the JobManager's checkpoint coordinator assigns a checkpoint ID and tells the sources to emit a **barrier** carrying that ID. Barriers travel in-band with the records and never overtake them, so the barrier cleanly separates "records before this checkpoint" from "records after it". A task with several input channels *aligns*: once the barrier for checkpoint N has arrived on one channel, it stops consuming that channel and buffers its records until the barrier for N has arrived on every channel. At that point the task's state contains exactly the pre-barrier records, it snapshots asynchronously, forwards the barrier downstream, and acknowledges to the coordinator. When every task has acknowledged, the checkpoint is complete and the coordinator notifies the tasks. This is Flink's asynchronous variant of Chandy-Lamport: no global pause, just a moving cut.
code
text · 9 lines-- aligned checkpoint at an operator with two inputs
channel A: r7 r8 [BARRIER-5] r9 r10 <- blocked, r9/r10 buffered
channel B: r3 r4 r5 r6 [BARRIER-5] ... <- still consumed
on BARRIER-5 from B:
snapshot state (contains r3..r8, nothing after)
emit BARRIER-5 downstream
ack checkpoint 5 to the coordinator
unblock A, process r9 r10go deeper
Recall that a barrier is a marker injected by the sources and travelling with the records, and that a task snapshots its state when the barrier passes. You are not expected to explain alignment in detail yet.
This is your tier. Walk the full path — coordinator triggers, sources inject, tasks align, snapshot asynchronously, acknowledge, coordinator completes — and explain precisely what alignment buffers and why dropping it yields at-least-once.
Read the three components of checkpoint duration (sync, async, alignment) and say what each one implicates. Be ready to trace a stuck checkpoint back to a single lagging subtask and explain why the recovery point ages while attempts are abandoned.
Be able to place the design against the Chandy-Lamport snapshot it derives from, and articulate why in-band markers with a central completion decision were the right trade for a dataflow engine with a known topology.
## The problem barriers solve A Flink job is a graph of parallel subtasks spread over many machines, each processing records at its own speed. To recover from failure you need a snapshot in which every subtask's state corresponds to the *same* prefix of the input. The naive approach — pause every task, snapshot, resume — would stall the pipeline for as long as the slowest state write takes. Barriers get the same consistency without a global pause. ## What a barrier is A checkpoint barrier is a lightweight control record injected into the data stream and carrying a checkpoint ID. It flows through exactly the same channels as ordinary records and, crucially, is never reordered around them. That single property is what makes it a *cut*: for any channel, everything before the barrier belongs to checkpoint N, everything after it does not. The sequence is: 1. The **checkpoint coordinator** on the JobManager decides it is time for checkpoint N and instructs every source subtask to trigger it. 2. Each source subtask snapshots its own state — its current read positions — and emits a barrier N into its output. 3. Barriers propagate downstream with the records. 4. Every task that has seen barrier N on all of its inputs snapshots its state, forwards barrier N, and acknowledges N to the coordinator with a handle to the persisted state. 5. When all subtasks have acknowledged, the coordinator writes the checkpoint metadata, marks N complete, and calls back into the tasks to tell them so. ## Alignment A task with more than one input channel — the parallel instances feeding a keyed operator, or the two sides of a join — will receive barrier N on its channels at different times, because upstream tasks run at different speeds. In `EXACTLY_ONCE` mode — the default in Flink 2.3 — the task **aligns**: the channel that already delivered barrier N is blocked and its subsequent records are buffered; channels that have not yet delivered it keep being consumed. When the last barrier arrives, the task's state reflects precisely the pre-barrier records from every channel — no more, no less — so it snapshots, forwards the barrier, and unblocks the buffered channels. Alignment is where the cost lives. The task idles on the fast channels while it waits for the slow one, and the buffered records occupy network memory. Under backpressure, when a barrier crawls through queues of already-buffered data, alignment can dominate checkpoint duration. ## The at-least-once alternative Set `CheckpointingMode.AT_LEAST_ONCE` and tasks stop blocking channels. A multi-input task still waits until barrier N has arrived on every input before it snapshots and forwards the barrier, but in the meantime it keeps consuming the channels that already delivered it, so records from *after* the barrier on those channels are processed into state. The snapshot is therefore *ahead* of the cut on those channels. Recovery still rewinds the sources to the checkpoint's positions, so those records are processed a second time — hence at-least-once. In exchange, the fast channels never stall. A task with a single input has nothing to align, so a job made only of one-to-one operators (`map`, `filter`, `flatMap`) behaves exactly-once even in this mode. This is a correctness-versus-latency choice, not a bug. ## Asynchronous state writing Snapshotting does not mean the task blocks while gigabytes are uploaded. The task takes a fast, in-memory copy-on-write or immutable view of its state, forwards the barrier, and continues processing while the write to checkpoint storage happens on a background thread. Only when that write finishes does the acknowledgement go out. This is why the reported *duration* of a checkpoint splits into synchronous, asynchronous and alignment components in the web UI, and why reading those three numbers separately is the standard first diagnostic step. ## Completion and the callback A checkpoint is complete only when *every* subtask has acknowledged it — one slow or stuck subtask means no completed checkpoint at all, and the recovery point silently ages. On completion the coordinator invokes the completion callback on the tasks. That callback is not decoration: it is exactly where a two-phase-commit sink commits its staged transaction, because completion is the first moment the job knows the snapshot it belongs to will never be rolled back. ## Relationship to Chandy-Lamport The algorithm is the classic Chandy-Lamport distributed snapshot adapted for acyclic dataflow with a known topology, published as *asynchronous barrier snapshotting*. The differences from the textbook version matter in interviews: Flink's sources are the only injection points, alignment replaces recording in-flight channel contents, and the coordinator holds the completion decision centrally rather than deriving it from marker exchange. Enabling unaligned checkpoints moves Flink *back* toward the textbook algorithm by persisting in-flight buffers instead of draining them.
- Why can a checkpoint barrier never overtake the records ahead of it in a channel?Because ordering is what makes the barrier a cut. If a barrier jumped ahead of records, some pre-barrier records would land in the post-barrier region and be replayed on recovery even though the snapshot already reflected them, or vice versa. Flink therefore treats the barrier as an ordinary in-band element of the channel, delivered in sequence with the data it separates.
- What does a large alignment time in the checkpoint statistics tell you?That at least one input channel is delivering its barrier much later than the others, so the task idles on the fast channels while it waits. The usual causes are backpressure making the barrier crawl through full network queues, skew putting far more work on one upstream subtask, or a slow single operator in one branch. Fix the imbalance, or enable unaligned checkpoints to stop paying for it.
- Which operator callback does a transactional sink use, and why that one?The completion notification the coordinator sends after every subtask has acknowledged. That is the earliest moment the job knows the snapshot is durable and will not be rolled back, so it is the only safe point to commit a staged transaction. Committing at snapshot time instead would publish output that a later abandoned checkpoint could force the job to reproduce.
Imagine several checkout queues merging into one counter, and a coloured card handed to each queue at the same moment. The counter waits until it has seen the card from every queue before taking a photo of the till.
saying these in an interview costs you the question
- Says Flink pauses the whole job to take a snapshot
- Thinks each operator triggers its own barrier independently
- Claims a barrier can overtake buffered records to speed things up
- Believes at-least-once mode still aligns, just without acknowledgements
- Thinks the checkpoint is complete once the sink sees the barrier