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.

+2 s+1 s0−1 s−2 s 06121824 h +20 ppm, never synced: +1.73 s/day −10 ppm: −0.86 s/day +20 ppm, resync 6 h
Two unsynchronized clocks with opposite drift are 2.6 seconds apart after one day. Periodic resynchronization keeps the error bounded.

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.

clientserver RTT = T1 − T0 = 40 ms T0 = 10:00:00.000 T1 = 10:00:00.040 Ts = 10:00:05.100 client sets clock to Ts + RTT/2 = 10:00:05.120 error ≤ ±RTT/2 = ±20 ms

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.

coordinator3:00 node A2:50 node B3:25 offset −10 min offset +25 min average = (0 − 10 + 25) / 3 = +5 adjust +5 → 3:05 adjust +15 → 3:05 adjust −20 → 3:05
The coordinator sends adjustments, not absolute times, so the reply's network delay doesn't distort the result. Node B slews backward gradually instead of jumping.

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.

PaperGusella & Zatti, "The Accuracy of the Clock Synchronization Achieved by TEMPO in Berkeley UNIX 4.3BSD," IEEE Trans. Software Eng. 15(7), 1989

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.

stratum 0stratum 1stratum 2clients GPS atomic S1S1 S2S2S2S2

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.

P1P2P3 m1m2 a · 1 · (1,0,0) b · 2 · (2,0,0) c · 2 · (1,1,0) d · 3 · (1,2,0) e · 1 · (0,0,1) f · 4 · (1,2,2) label: event · Lamport · vector
b and c share Lamport time 2, and e (1) is smaller than b (2), yet none of these pairs is causally related. Their vectors show it: (2,0,0) and (1,1,0) don't dominate each other, so b ∥ c. f's vector (1,2,2) dominates a, c, d, and e, so all four happened before it.

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.

coordinatorqueue: [P2] P1in critical section P2waiting P3idle grant request
Simple, fair, and cheap, but the coordinator is a single point of failure and a throughput bottleneck. A crashed coordinator also can't be told apart from a long queue.

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.

P1P2P3 req(8)req(12) OK (8 < 12) P1 in CS deferred OK P2 in CS request OK 2(N − 1) messages per entry · any crashed process blocks everyone

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.

Token ring Grid quorum (Maekawa) T token moves clockwise ask row + column: 5 of 9 nodes any twoquorumsintersect
AlgorithmMessages per entryEntry delay (message times)Main weakness
Centralized32Coordinator failure and bottleneck
Lamport3(N − 1)2Any process failure blocks progress
Ricart–Agrawala2(N − 1)2Any process failure blocks progress
Token ring1 to ∞0 to N − 1Lost token; idle circulation
Maekawa≈ 3√N2Deadlock-prone without extra messages
Majority quorum≈ N2Needs 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.

client 1lock svcclient 2storage token 33 paused (GC stop) lease expires token 34 write(34) ✓ write(33) rejected: 33 < 34
Without the token check, both clients would write and the "exclusive" lock would have failed silently.

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.

PessimisticOptimistic
IdeaLock first, then read and writeRead and work freely, then validate at commit; retry on conflict
MechanismLock service, row locks, two-phase lockingVersion numbers, compare-and-swap, conditional writes
Wins whenContention is high and retries are expensiveContention is low and conflicts are rare
Fails asWaiting, lock convoys, deadlockRepeated 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.

site A site B P1 P2 P3 waits forwaits forwaits for neither site alone sees the cycle
ApproachHowCost
PreventionBreak a Coffman condition, such as always acquiring locks in a global orderLess concurrency; needs to know locks in advance
AvoidanceGrant only if the system stays in a safe state (Banker's algorithm)Needs future demand; rarely practical in distributed systems
Timestamp schemesWait-die or wound-wait: older transactions win, younger ones abortSome unnecessary aborts
DetectionBuild a global wait-for graph, or send probes along edges (Chandy–Misra–Haas)Message overhead; risk of phantom deadlocks from stale views
TimeoutsAbort any wait that takes too longSimple; 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.

true time → TT.now() at commit commit starts s = latest TT.now() after waiting commit wait ≈ 2ε earliest > s: release locks

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.

Burrows, OSDI 2006 · ZooKeeper programmer's guide

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