Distributed systems · synchronization
Keeping time, taking turns
On one machine, a single clock orders every event and the kernel's locks decide who goes next. Spread the work across machines and both disappear. Each node has its own drifting clock, and there is no shared memory to hold a lock. This page covers how distributed systems agree on time, order events without it, and grant exclusive access to shared resources.
Why coordination is hard
Suppose two ATMs read the same balance of $100, each approves a $80 withdrawal, and each writes back $20. Both reads and both writes succeed, yet $60 has vanished from the books. That is a race condition: the result depends on the timing of operations that should have been serialized. Locally, a mutex prevents it. Across machines, three things get in the way:
- No common clock. Timestamps from different machines can't be compared reliably, so "who was first" is ambiguous.
- Unbounded message delay. A request can be slow or lost, and a slow node is indistinguishable from a dead one.
- Independent failure. The node holding a lock can crash, pause, or lose its network while everyone else waits.
The two generals problem captures the core difficulty. Two armies must attack at the same time and can only communicate by messenger through hostile ground. Every confirmation needs its own confirmation, and no finite number of messages over an unreliable channel makes both sides certain. Practical systems therefore aim for agreement that is probable, bounded, or eventually reached, never certain in every case.
Clock drift
Computer clocks count oscillations of a quartz crystal. No two crystals run at exactly the rated frequency, and temperature and age shift them further. The error is the drift rate, measured in parts per million (ppm). A typical server crystal drifts 10 to 50 ppm. At 20 ppm a clock gains or loses 72 ms per hour, about 1.7 seconds per day.
Two clocks with maximum drift ρ can diverge by up to 2ρ per unit time. To keep any two clocks within δ of each other, resynchronize at least every δ / 2ρ seconds. Synchronization comes in two flavors: external, matching an authoritative source such as UTC, and internal, keeping nodes close to each other whatever the true time.
One rule applies to every algorithm below: never step a clock backward. Software that measures durations or orders log entries breaks when time repeats. Instead, clocks are slewed, run slightly fast or slow until the error closes. For durations, programs should use a monotonic clock, which only moves forward, instead of the wall clock.
ReadPython time.monotonic · Google public NTP: leap second smearing
Cristian's algorithm
A client asks a time server for the time and measures the round trip on its own clock. Assuming the request and reply took equal time, the server's answer is half a round trip old when it arrives.
The accuracy bound comes from asymmetry: if the request took the full 40 ms and the reply took none, the estimate is 20 ms off. Clients improve accuracy by sending several probes and keeping the one with the shortest round trip, since it has the least room for asymmetric delay. The weakness is the single server: if it is wrong or unreachable, every client is wrong or unsynchronized.
PaperCristian, "Probabilistic Clock Synchronization," Distributed Computing 3(3), 1989
Berkeley algorithm
When no node has an authoritative clock, the goal shifts to internal agreement. A coordinator polls every node, estimates each one's offset using Cristian-style round-trip correction, averages the offsets, and sends each node the adjustment it needs.
Two refinements make it robust. A fault-tolerant average ignores clocks that differ wildly from the rest, so one broken crystal can't drag the group. If the coordinator fails, an election picks a new one. The group agrees with itself, but the shared time can drift away from UTC.
NTP and PTP
The Network Time Protocol synchronizes clocks across the internet through a hierarchy of servers. Stratum 0 devices are reference clocks such as GPS receivers and atomic clocks. Stratum 1 servers attach directly to them, stratum 2 servers sync from stratum 1, and so on. Each hop adds a little error but spreads the load.
Each NTP exchange records four timestamps: client send (T1), server receive (T2), server send (T3), client receive (T4). Server processing time cancels out:
offset θ = ((T2 − T1) + (T3 − T4)) / 2 delay δ = (T4 − T1) − (T3 − T2)
NTP queries several servers, discards outliers, and steers the local clock with a feedback loop. It typically holds a few milliseconds over the internet and well under a millisecond on a LAN. The Precision Time Protocol (IEEE 1588) reaches sub-microsecond accuracy by timestamping packets in network hardware. Finance, telecom, and some databases depend on it.
Papers and specsMills, "Internet Time Synchronization: the Network Time Protocol," IEEE Trans. Comm. 39(10), 1991 · RFC 5905 (NTPv4)
Toolchrony documentation (run chronyc tracking to see your own offset)
Logical clocks
Often the question is not "what time is it" but "which happened first". Lamport observed that order can be defined from causality alone. Event a happened before b (a → b) if they occur in the same process with a first, if a is the send of a message and b its receipt, or through a chain of such steps. Events with no such chain are concurrent.
Lamport timestamps
Each process keeps one counter. It increments before each event, attaches the counter to outgoing messages, and on receipt sets the counter to max(local, received) + 1. If a → b then L(a) < L(b). The converse doesn't hold: a smaller timestamp does not prove causality. Breaking ties by process ID gives a total order that every node agrees on.
Vector clocks
Each process keeps a vector with one counter per process. It increments its own entry on every event and takes the element-wise maximum on receipt. Now a → b exactly when V(a) ≤ V(b) in every entry and differs in at least one. If neither vector dominates, the events are concurrent. That is how replicated stores detect conflicting writes.
Vector clocks cost one entry per process, which grows with cluster size. Hybrid logical clocks combine physical time with a small logical counter. They stay close to wall-clock time, so they're readable, while still respecting causality.
PapersLamport, "Time, Clocks, and the Ordering of Events in a Distributed System," CACM 21(7), 1978 · Mattern, "Virtual Time and Global States of Distributed Systems," 1989 · Kulkarni et al., "Logical Physical Clocks," OPODIS 2014
Mutual exclusion
A critical section is code that touches a shared resource and must not run in two places at once. Any distributed mutual exclusion algorithm is judged on three properties, plus cost:
- Safety
- At most one process is in the critical section at any time.
- Liveness
- Every request is eventually granted. No deadlock, no starvation.
- Fairness
- Requests are granted in the order they were made, usually happened-before order.
- Cost
- Messages per entry, delay before entry, and tolerance of node or message failure.
The simplest solution simulates a single-machine lock. One coordinator holds a queue. A process sends request, waits for grant, enters, and sends release, for three messages per entry.
Mutex algorithms
Permission-based: Ricart–Agrawala
A process that wants the critical section sends a timestamped request to every other process and enters once all have replied OK. A process receiving a request replies at once unless it is inside the critical section or has its own earlier request pending. In those cases it defers the reply until it exits. Lamport timestamps decide "earlier", with process ID breaking ties.
Token-based
A single token circulates, and only its holder may enter. In a token ring it passes from neighbor to neighbor. Requests are served in ring order, so there's no starvation, but the token travels even when nobody needs it, and a lost token must be detected and regenerated without creating two. Suzuki–Kasami broadcasts requests so the token goes straight to a waiting process, and Raymond's algorithm routes it along a tree in O(log N) messages.
Quorum-based
Asking everyone is unnecessary. It is enough to ask a quorum, as long as any two quorums overlap. The shared member can't grant to both, so safety holds. With majority quorums, the system keeps working while a minority of nodes are down. Maekawa arranged N nodes in a √N × √N grid, with a node's quorum being its row plus its column. That cuts the cost to about 3√N messages. The same overlap idea underlies consensus protocols such as Paxos and Raft.
| Algorithm | Messages per entry | Entry delay (message times) | Main weakness |
|---|---|---|---|
| Centralized | 3 | 2 | Coordinator failure and bottleneck |
| Lamport | 3(N − 1) | 2 | Any process failure blocks progress |
| Ricart–Agrawala | 2(N − 1) | 2 | Any process failure blocks progress |
| Token ring | 1 to ∞ | 0 to N − 1 | Lost token; idle circulation |
| Maekawa | ≈ 3√N | 2 | Deadlock-prone without extra messages |
| Majority quorum | ≈ N | 2 | Needs a majority reachable |
PapersRicart & Agrawala, "An Optimal Algorithm for Mutual Exclusion in Computer Networks," CACM 24(1), 1981 · Maekawa, "A √N Algorithm for Mutual Exclusion in Decentralized Systems," ACM TOCS 3(2), 1985 · Suzuki & Kasami, "A Distributed Mutual Exclusion Algorithm," ACM TOCS 3(4), 1985 · Raymond, "A Tree-Based Algorithm for Distributed Mutual Exclusion," ACM TOCS 7(1), 1989
Distributed locks
In production, mutual exclusion is usually a service rather than an algorithm each application runs itself. A small, replicated lock service such as ZooKeeper, etcd, or Consul uses consensus internally and exposes simple acquire and release calls. Two ideas make these locks safe when clients fail.
Leases. A lock is granted for a limited time and must be renewed. If the holder crashes, the lease expires and the lock frees itself, so nobody waits forever. Choosing the duration is a trade-off: short leases recover quickly but risk expiring under a slow but healthy client.
Fencing tokens. A lease alone is not enough. A client can pause, for a garbage-collection stop, a VM migration, or a page fault, and wake up still believing it holds an expired lock. The fix is a number that increases with every grant, which the protected resource checks.
ZooKeeper's lock recipe uses ephemeral sequential nodes: each client creates a numbered node that disappears if its session dies, and the lowest number holds the lock. Each waiter watches only the node just ahead of it, which avoids a herd of clients waking at once. The node's sequence number doubles as a fencing token. Locks built on clock-dependent algorithms without fencing, such as Redis Redlock, have been criticized for exactly the pause scenario above.
PapersGray & Cheriton, "Leases: An Efficient Fault-Tolerant Mechanism for Distributed File Cache Consistency," SOSP 1989 · Burrows, "The Chubby Lock Service for Loosely-Coupled Distributed Systems," OSDI 2006 · Hunt et al., "ZooKeeper: Wait-free Coordination for Internet-scale Systems," USENIX ATC 2010
ReadKleppmann, "How to do distributed locking" · ZooKeeper recipes · etcd concurrency API
Locking strategy
Every distributed lock costs at least one network round trip. That cost is spent in one of two ways.
| Pessimistic | Optimistic | |
|---|---|---|
| Idea | Lock first, then read and write | Read and work freely, then validate at commit; retry on conflict |
| Mechanism | Lock service, row locks, two-phase locking | Version numbers, compare-and-swap, conditional writes |
| Wins when | Contention is high and retries are expensive | Contention is low and conflicts are rare |
| Fails as | Waiting, lock convoys, deadlock | Repeated aborts under heavy contention |
Granularity is the other dial. A coarse lock over a whole table is simple but serializes unrelated work. Fine locks on single rows or keys allow parallelism but multiply lock traffic and the chance of deadlock. Many systems use hierarchical or intention locks to get both.
-- optimistic update: succeeds only if nobody changed the row since it was read UPDATE accounts SET balance = 20, version = 8 WHERE id = 42 AND version = 7;
PaperKung & Robinson, "On Optimistic Methods for Concurrency Control," ACM TODS 6(2), 1981
Deadlock
Deadlock is a set of processes each waiting for a resource another holds, so none can proceed. It requires all four Coffman conditions at once:
- Mutual exclusion
- A resource can be held by only one process.
- Hold and wait
- A process holds resources while requesting more.
- No preemption
- Resources can't be taken away, only released.
- Circular wait
- A cycle of processes each waits on the next.
Locks spread across machines make deadlocks harder to see. Each site knows only its own waits, so a cycle that crosses sites is invisible locally.
| Approach | How | Cost |
|---|---|---|
| Prevention | Break a Coffman condition, such as always acquiring locks in a global order | Less concurrency; needs to know locks in advance |
| Avoidance | Grant only if the system stays in a safe state (Banker's algorithm) | Needs future demand; rarely practical in distributed systems |
| Timestamp schemes | Wait-die or wound-wait: older transactions win, younger ones abort | Some unnecessary aborts |
| Detection | Build a global wait-for graph, or send probes along edges (Chandy–Misra–Haas) | Message overhead; risk of phantom deadlocks from stale views |
| Timeouts | Abort any wait that takes too long | Simple; may abort healthy work |
Preventing deadlock does not remove blocking. A global lock order still makes processes wait, and timeouts or aborts move the cost to retries. The choice is about where to pay.
PapersCoffman, Elphick & Shoshani, "System Deadlocks," ACM Computing Surveys 3(2), 1971 · Chandy, Misra & Haas, "Distributed Deadlock Detection," ACM TOCS 1(2), 1983 · Dijkstra, "Hierarchical Ordering of Sequential Processes" (dining philosophers), 1971
In practice
Google Spanner and TrueTime
physical clocksbounded uncertainty
Spanner's TrueTime API returns an interval, not a single time: the true time is guaranteed to fall between earliest and latest. GPS receivers and atomic clocks in each datacenter keep that uncertainty ε to a few milliseconds. To commit, a transaction takes timestamp s = latest, then waits until earliest > s before releasing its locks. That commit wait guarantees any later transaction anywhere gets a larger timestamp.
Corbett et al., "Spanner: Google's Globally-Distributed Database," OSDI 2012
Amazon Dynamo and vector clocks
logical clocksconflict detection
Dynamo attaches a vector clock to each object version. When replicas return versions whose vectors don't dominate each other, the writes were concurrent, and both are returned for the application to merge. A shopping cart, for example, is merged by taking the union of items.
DeCandia et al., "Dynamo: Amazon's Highly Available Key-value Store," SOSP 2007
Chubby and ZooKeeper
lock servicesleases
Google built Chubby as a coarse-grained lock and naming service for leader election in systems such as GFS and Bigtable. ZooKeeper, its open-source counterpart, exposes lower-level primitives instead of locks directly. Applications compose those primitives into locks, barriers, queues, and leader election.
Tools and reading
Try it
- chronyInspect live offset, drift (ppm), and sources on any Linux machine.
- DistAlgoWrite Lamport clocks and mutex algorithms in Python-like code and run them.
- ZooKeeper getting startedBuild a lock from ephemeral sequential nodes; kill a client and watch the lock pass.
- etcd quick startUse etcdctl lock and leases from the command line.
- Raft visualizationSee quorums, terms, and leader election in action.
- TLA+Model-check a mutex algorithm for safety and liveness.
Read it
- van Steen & Tanenbaum, Distributed Systems, ch. 5Free textbook chapter on coordination.
- Kleppmann, Designing Data-Intensive Applications, ch. 8Unreliable clocks, process pauses, and fencing.
- Lamport, 1978Happened-before and logical clocks.
- Ricart & Agrawala, 1981Permission-based mutual exclusion.
- Mills, 1991NTP design.
- Corbett et al., 2012Spanner and TrueTime.