Distributed Systems
The Chandy-Lamport Snapshot: Photographing a Running Distributed System
The Chandy-Lamport Snapshot is a 1985 algorithm that records a consistent global state of a running distributed system — every process's local state plus the messages still in flight on every channel — without ever pausing the computation. There is no shared clock and no way to freeze every machine at the same instant, yet the algorithm captures a picture the system could genuinely have passed through, using nothing but small control messages called markers and the assumption that each channel delivers in order.
It is the theoretical bedrock of modern fault tolerance: Apache Flink's exactly-once stream checkpoints, distributed deadlock and termination detection, and rollback recovery all descend directly from it.
- TypeDistributed global-snapshot algorithm (non-blocking)
- Published1985, K. Mani Chandy & Leslie Lamport, ACM TOCS 3(1):63-75
- AssumesFIFO, reliable, directed channels; no process/link failures during snapshot
- Message costExactly one marker per channel: O(E); O(N^2) for N fully-connected processes
- GuaranteeRecords a consistent cut - a global state reachable on a valid execution
- Used inApache Flink exactly-once checkpoints, rollback recovery, deadlock/termination detection
Interactive visualization
Press play, or step through manually. The visualization is yours to drive — try it before reading on.
Watch the 60-second explainer
A condensed visual walkthrough — narrated, captioned, under a minute.
The Problem: A Photograph With No Shutter
Suppose you want to know the total amount of money in a banking system spread across many machines, or whether a distributed computation has terminated, or you simply want a checkpoint to roll back to after a crash. You need a global state: the local state of every process together with the state of every communication channel. The trouble is there is no global clock, and machines cannot be frozen simultaneously. If each process just reports its own state whenever it feels like it, the results are mutually inconsistent — one process may report a dollar as already spent while the recipient reports it as not yet received, so a dollar vanishes.
Two facts make a distributed snapshot genuinely hard. First, a message that has been sent but not yet received is real, live state sitting inside a channel; ignore it and your snapshot is wrong. Second, the property we actually want is a consistent cut: a partition of all events into 'before' and 'after' such that no message is received 'before' the cut while being sent 'after' it. A cut that violates this describes a state that no lawful execution could ever produce — an effect preceding its cause. Chandy and Lamport's insight was that you can construct such a cut on the fly, while the system keeps running, provided channels deliver messages in order.
The Algorithm: Markers and Two Rules
The entire protocol is a special control message — the marker — plus two rules. Any process may initiate.
- Initiation: the initiator records its own local state, then immediately sends a marker on every outgoing channel, before sending any further application messages.
- Marker rule (receiving a marker on channel c): If the process has not yet recorded its state, it (1) records its own state now, (2) records the state of channel
cas empty, and (3) begins logging every application message that arrives on each of its other incoming channels. If the process has already recorded its state, it records the state of channelcas exactly the sequence of messages it logged oncsince recording its state and before this marker arrived. - Propagation: the moment a process records its own state (whether as initiator or on first marker), it sends a marker on all its outgoing channels before any subsequent application message.
The algorithm terminates when every process has received a marker on all of its incoming channels; at that point each process has one recorded local state and one recorded set for each incoming channel. A separate collection step (any process can broadcast a request, or results flow back to the initiator) gathers the fragments into one global snapshot.
Why It Works: FIFO and the Four Message Cases
Correctness rests entirely on FIFO channels and the discipline of sending the marker before any post-snapshot message. Classify every application message m sent from p to q by whether its send and its receive fall before or after the respective process recorded its state:
- Sent pre-snapshot, received pre-snapshot: already reflected in both local states; not in the channel.
- Sent pre-snapshot, received post-snapshot: genuinely in flight - this is exactly what the logging phase captures as the channel state.
- Sent post-snapshot, received post-snapshot: belongs to the future; recorded in neither.
- Sent post-snapshot, received pre-snapshot: the inconsistent case - and it is impossible. Because p sends its marker before m, and the channel is FIFO, the marker reaches q before m; q records its state on that marker, so m can only arrive post-snapshot.
Eliminating the fourth case is precisely what guarantees a consistent cut. The formal payoff is a reachability theorem: if Si is the global state when the snapshot began and Sf the state when it ended, the recorded snapshot S* is reachable from Si, and Sf is reachable from S*. The photograph may never have existed at any wall-clock instant, but it is a state the system could have passed through on a valid execution — which is all any downstream analysis needs.
Cost, Complexity, and What It Assumes
Message overhead is minimal and exact: exactly one marker traverses each directed channel, so a run costs O(E) markers for E channels — N(N-1), i.e. O(N2), for N fully-connected processes, but only O(E) for sparse topologies. Each marker is a tiny fixed-size control packet. Time is bounded by the network diameter times the maximum message-propagation delay: the snapshot 'wavefront' spreads outward like markers flooding a graph. Space per process is one recorded local state plus the logged in-flight messages, which is proportional to how much traffic is genuinely in transit when the marker passes — usually small, but unbounded in the worst case of a very busy channel.
Crucially the algorithm is non-blocking: no process ever stops computing or delays application messages to take part; it merely records and logs alongside its normal work. The price is a set of assumptions that must hold: channels are FIFO, reliable (no loss, no duplication), and unidirectional, and no process or link fails during the snapshot. Drop FIFO and the four-case argument collapses — a post-snapshot message can overtake the marker and be mislogged. Failures during the run can leave some incoming channel without its marker forever, so the collection step must be paired with a failure detector or timeout in real deployments.
Variants and Look-alikes
Chandy-Lamport's FIFO requirement spawned a family of alternatives. Lai-Yang (1987) needs no control messages at all: it colors every application message white (pre-snapshot) or red (post-snapshot) by piggybacking a single bit, and a receiver infers the channel state from the colors it sees — trading marker traffic for a per-message tag and the burden of retaining message history. Mattern's (1993) approach uses per-process counters (vector-clock-like) of messages sent and received to deduce in-flight counts without FIFO. Spezialetti-Kearns optimizes the common case of several processes initiating snapshots concurrently, merging their marker regions so the work is shared rather than duplicated.
It is worth contrasting the snapshot with two cousins it is often confused with. A stop-the-world checkpoint also yields a consistent state, but by halting every process — trivially correct and completely impractical at scale. Lamport clocks and vector clocks supply the happened-before ordering that defines a consistent cut, but they timestamp events rather than materialize the channel contents; Chandy-Lamport is what you run when you need the actual bytes in flight, not just their causal order. And unlike a database's write-ahead log — which records a single node's history durably — the snapshot coordinates many nodes into one coherent instantaneous picture.
Real Systems: Flink, Recovery, and Stable-Property Detection
The algorithm's biggest modern deployment is Apache Flink. Flink's Asynchronous Barrier Snapshotting (Carbone et al., 2015) is a direct descendant: it injects checkpoint barriers — Chandy-Lamport markers — into the data streams flowing between operators. When an operator has received the barrier on every input, it snapshots its state to durable storage (RocksDB/HDFS/S3). For acyclic dataflow graphs Flink relies on barrier alignment and need not log in-flight records at all; for cyclic graphs it logs the backward-edge messages exactly as Chandy-Lamport prescribes. This is precisely how Flink delivers exactly-once processing: on failure it restores the last complete snapshot and replays. Kafka Streams and Spark Structured Streaming use conceptually similar barrier/offset checkpoints.
The classic application is stable-property detection. A stable property is one that, once true, stays true — deadlock, termination, token loss, or garbage that will never be referenced again. The snapshot's reachability guarantee makes it a valid oracle: if a stable property holds in the recorded state S*, then (since Sf is reachable from S*) it holds in the real current state; if it is false in S*, it was false when the snapshot began. So you can detect a distributed deadlock or decide that a computation has finished by evaluating the predicate on a snapshot taken without ever stopping the system — the same principle behind checkpoint/rollback recovery and distributed debuggers that reconstruct a consistent global view from a live, uncoordinated cluster.
| Method | Channel model | In-flight capture | Cost / tradeoff |
|---|---|---|---|
| Chandy-Lamport (1985) | FIFO required | Log messages between recording state and the marker | One marker per channel (O(E)); non-blocking; needs FIFO |
| Lai-Yang (1987) | Non-FIFO OK | Color/tag every app message (white/red); no control messages | Zero extra messages but piggybacks a bit + unbounded channel history |
| Mattern (1993) | Non-FIFO OK | Vector-counter of sent/received messages per process | Counting-based; needs message counters, tolerant of reordering |
| Flink ABS (2015) | FIFO dataflow edges | Align barriers; log only on cyclic edges (none for DAGs) | Barrier alignment adds latency; acyclic graphs skip channel logging |
| Stop-the-world checkpoint | Any | Trivially none (system halted) | Simple and exact but blocks all progress; unacceptable at scale |
Frequently asked questions
Does the Chandy-Lamport snapshot capture a state the system was actually in?
Not necessarily at any single wall-clock instant. It records a consistent cut, which is a global state the system could have passed through on some valid execution. Formally the snapshot S* is reachable from the state when recording began and can itself reach the state when recording ended, so any stable property you evaluate on it gives a sound answer about the real system.
Why does the algorithm require FIFO channels?
FIFO is what guarantees the marker acts as a clean dividing line. A process sends its marker before any post-snapshot application message, so on an in-order channel the marker always arrives before those later messages. This makes the 'sent-after but received-before' case impossible, which is exactly the case that would produce an inconsistent cut. Without FIFO you must switch to a variant like Lai-Yang or Mattern that tags or counts messages instead.
What exactly is recorded as a channel's state?
The set of messages that were in flight: sent before the sender took its snapshot but not yet received when the receiver took its. Operationally, a process records a channel as empty if it first learns of the snapshot from that channel's marker, and otherwise records the messages it logged on that channel between recording its own state and receiving the marker on it.
How many marker messages does a snapshot cost?
Exactly one marker per directed channel, so O(E) markers for E channels. For N fully connected processes that is N(N-1), i.e. O(N^2), but for realistic sparse topologies it is far cheaper. Markers are tiny fixed-size control packets, and the algorithm never blocks application traffic, so the runtime overhead is very low.
How does Apache Flink use Chandy-Lamport?
Flink's Asynchronous Barrier Snapshotting inserts checkpoint barriers (markers) into its data streams. Each operator snapshots its state once it has seen the barrier on all inputs, writing to durable storage. On failure Flink restores the latest complete snapshot and replays, which is how it achieves exactly-once semantics. For acyclic pipelines it uses barrier alignment and skips in-flight logging; for cyclic ones it logs loop-back records just as the original algorithm requires.
What happens if a process crashes during the snapshot?
The basic algorithm assumes no failures during the run. A crash can leave some incoming channel without its expected marker, so that process never finishes recording and the snapshot can stall. Real systems pair the protocol with timeouts or failure detectors, and combine snapshots with checkpoint/rollback so that a failed run is simply abandoned and retried rather than corrupting anything.