Learn Labs
10. Consistency and Consensus

10.2 ID Generators and Logical Clocks

Why a single-node autoincrementing counter is nice: compact (64 bits, or 32 if you're sure you'll never exceed 4 billion records — but that is risky), and the ORDER of the IDs tells you the order in which records were created.

This single-node ID generator is another example of A LINEARIZABLE SYSTEM. Each request is an atomic FETCH-AND-ADD; linearizability ensures that if Aaliyah's post completes before Bryce's begins, Bryce's ID must be greater. (Concurrent posts may be ordered either way, as long as the IDs are unique.)

Three problems:

  1. Not fault-tolerant — a single point of failure
  2. Slow for records created in another region — potentially a round trip TO THE OTHER SIDE OF THE PLANET just to get an ID
  3. Could become a bottleneck at high write throughput

2.1 The alternatives, and what each loses

SchemeHowWhat it loses
Sharded ID assignmentOne node generates only even numbers, another only odd; generally reserve some bits for a shard numberOrdering. IDs 16 and 17 don't tell you which message was sent first — one node might have been ahead of the other
Preallocated blocksNode A claims 1–1,000, node B claims 1,001–2,000; each hands out from its block and requests a new block when running lowOrdering. One message may get an ID in 1,001–2,000 and a LATER message an ID in 1–1,000
Random UUIDs (v4)Generated locally on any node WITHOUT COMMUNICATION128 bits, and the order is RANDOM — comparing two IDs tells you NOTHING about which is newer
Wall-clock timestamp made uniqueTimestamp in the most significant bits, remaining bits filled with a shard number + per-shard sequence, or a long random value. Used by UUIDv7, X's Snowflake, ULIDs, Hazelcast Flake IDs, MongoDB ObjectIDsAt best APPROXIMATE ordering. An earlier write from a slightly fast clock and a later write from a slightly slow clock get inverted; with clock jumps, EVEN A SINGLE NODE'S TIMESTAMPS MIGHT BE ORDERED INCORRECTLY. Unlikely to be linearizable

2.2 Logical clocks

A LOGICAL CLOCK is an ALGORITHM THAT COUNTS THE EVENTS that have occurred. A timestamp from a logical clock DOESN'T TELL YOU WHAT TIME IT IS, but you can compare two timestamps to tell which is earlier and which is later.

Three requirements:

  1. Timestamps are compact (a few bytes) and unique
  2. Any two can be compared to determine which is earlier — they are TOTALLY ORDERED
  3. The order is CONSISTENT WITH CAUSALITY: if A happened before B, then A's timestamp < B's timestamp

A single-node ID generator meets these. The distributed ID generators above DO NOT meet the causal ordering requirement.

Lamport clocks (1978, Leslie Lamport — one of the most-cited papers in distributed systems)

A Lamport timestamp is a pair: (counter, node ID), with two rules:

  1. Every time a node generates a timestamp, it increments its counter and uses the new value.
  2. Every time a node sees a timestamp from another node, if that timestamp's counter is greater than its local counter, it increases its local counter to match.
AaliyahCalebBrycemsg (1,"Aaliyah")Aaliyah 0→1msg (1,"Caleb")Caleb 0→1reply (2,"Bryce")Bryce 1→2

Bryce starts at counter 0. Seeing counter 1 arrive, he raises his local counter to 1 (rule ②), then increments it to 2 when he generates his reply (rule ①).

Comparison: compare the counter first; if equal, compare the node ID lexicographically — (1,"Aaliyah") < (1,"Caleb") < (2,"Bryce").

Figure 10.2.1Lamport clocks (1978, Leslie Lamport — one of the most-cited papers in distributed systems)

⚠️ Although Lamport clocks provide a TOTAL ORDERING, THEY DO NOT PROVIDE LINEARIZABILITY — they are not a way of ensuring a value is up to date. They are MERELY a way of assigning IDs such that if A happened before B, A's ID is less than B's.

Two limitations:

  1. No direct relation to physical time — you can't find all messages posted on a particular date; you'd need to store the physical time separately
  2. If two nodes NEVER COMMUNICATE, one node's increments are never reflected in the other's counter. So events generated around the same time on different nodes could have WILDLY DIFFERENT counter values
Hybrid logical clocks (HLC)

Combines the advantages of physical time-of-day clocks with the ordering guarantees of Lamport clocks.

  • Like a PHYSICAL clock: it counts seconds or microseconds
  • Like a LAMPORT clock: when one node sees a greater timestamp from another, it MOVES ITS OWN LOCAL VALUE FORWARD to match. So if one node's clock runs fast, THE OTHERS WILL SIMILARLY MOVE THEIR CLOCKS FORWARD when they communicate
  • Every generated timestamp is ALSO INCREMENTED, ensuring the clock moves forward MONOTONICALLY even if the underlying physical clock JUMPS BACKWARD (e.g. from NTP adjustments)

⇒ You can treat an HLC timestamp ALMOST LIKE a conventional time-of-day timestamp, with the added property that ITS ORDERING IS CONSISTENT WITH HAPPENS-BEFORE. It doesn't depend on special hardware and requires only ROUGHLY SYNCHRONIZED CLOCKS. (Used by CockroachDB.)

Where these fit with MVCC: Lamport clocks and HLCs are a GOOD WAY of generating the transaction IDs that snapshot isolation needs (Ch 8), because they ensure THE SNAPSHOT IS CONSISTENT WITH CAUSALITY.

vs vector clocks
Lamport / HLCVector clock
Concurrent timestampsOrdered ARBITRARILY. You generally CAN'T TELL whether two timestamps were generated concurrently or one happened before the otherCAN detect concurrency: if A has a higher counter for one node and B a higher counter for another, A and B MUST BE CONCURRENT
SizecompactMuch larger — potentially ONE INTEGER FOR EVERY NODE in the system

2.3 Why logical clocks aren't enough — the privacy leak

User A does this……which is written toTimestamp
On their laptop: set account to privateAccounts DBts = 100
On their phone: upload embarrassing photoPhotos DBts = 95

The photos DB never read from the accounts DB, so its local counter is behind, and the photo gets a lower timestamp than the settings change.

A viewer who is not a friend reads with an MVCC snapshot at ts = 98:

  • photo upload (95) ≤ 98 ⇒ visible
  • privacy change (100) > 98 ⇒ not yet applied — the account still looks public

⇒ The viewer sees the photo they were not supposed to see.

Linearizability requires that if request A COMPLETED before request B BEGAN, then B must have the higher ID — EVEN IF A AND B NEVER COMMUNICATED WITH EACH OTHER. Lamport clocks can ensure only that a node generates timestamps GREATER THAN ANY OTHER TIMESTAMP THAT NODE HAS SEEN; no such guarantees can be made about timestamps IT HASN'T SEEN.

Possible fixes, and why they're bad: the photos DB could read the account status before writing — "but it's EASY TO FORGET SUCH A CHECK." The app could track the user's latest write timestamp — "but if the user uses a LAPTOP AND A PHONE, THAT'S NOT SO EASY." The simplest solution is a LINEARIZABLE ID GENERATOR.

2.4 Implementing a linearizable ID generator

Option A — a single node doing three things:

  1. Atomically increment a counter and return its value
  2. Persist the counter (so it doesn't generate duplicates after a crash)
  3. Replicate it for fault tolerance (single-leader replication)

(Used in practice: TiDB/TiKV calls it a TIMESTAMP ORACLE, inspired by Google's Percolator.)

The batching optimization:

Avoid a disk write and replication on every request: write a record describing A BATCH of IDs; once persisted and replicated, hand out those IDs in sequence. Before running out, persist the record for the next batch. SOME IDs WILL BE SKIPPED if the node crashes or you fail over — BUT YOU WON'T ISSUE ANY DUPLICATE OR OUT-OF-ORDER IDs.

The limits:

  • You can't easily SHARD it — multiple shards independently handing out IDs breaks linearizable order
  • You can't easily distribute it across REGIONS — in a geo-distributed database, all ID requests must go to a node in a single region
  • On the upside, the job is very simple, so a single node can handle a large request throughput

Option B — Spanner's approach: rely on a physical clock returning a RANGE of timestamps and wait for the uncertainty interval to elapse (Ch 9 §3.6).

This guarantees linearizable ID assignment WITHOUT ANY COMMUNICATION; even requests in different regions are ordered correctly, WITHOUT WAITING FOR CROSS-REGION REQUESTS. The downside is that you need HARDWARE AND SOFTWARE SUPPORT for tightly synchronized clocks and computing the uncertainty interval.

2.5 Why even a linearizable ID generator isn't enough for locks

You could use a logical clock to assign timestamps to lock requests and pick the LOWEST as the winner. If the clock is linearizable, you know future requests will generate greater timestamps.

But part of the problem is STILL UNSOLVED: HOW DOES A NODE KNOW WHETHER ITS OWN TIMESTAMP IS THE LOWEST? To be sure, IT NEEDS TO HEAR FROM EVERY OTHER NODE that might have generated a timestamp. If one of them has failed or is unreachable, THE SYSTEM WOULD GRIND TO A HALT. This is not the kind of fault-tolerant system we need.

To implement locks, leases, and similar constructs in a fault-tolerant way, WE NEED SOMETHING STRONGER. WE NEED CONSENSUS.


On this page