Distributed systems
Part 5 of 6 · Time, Clocks & OrderingHybrid Logical Clocks & TrueTime - Commit Wait, Uncertainty Intervals & External Consistency
HLC algorithm (l, c) with skewed nodes and a max-offset guard (runnable); CockroachDB max-offset self-termination, MongoDB cluster time, YugabyteDB; TrueTime intervals and commit wait (runnable) plus CockroachDB-style uncertainty restarts; external consistency vs serializability vs SI; TrueTime vs HLC vs timestamp oracle vs single leader.
- 1Gist
- 2Maps
- 3Q&A
- 4Sandbox
Voice readout needs Web Speech Synthesis in this browser.
Question ladder
L1
What does HLC add to a Lamport clock?
Answer
The same causal order, with the physical component kept within the skew bound of real time so stamps work for snapshots and TTLs.
L2
When does c reset?
Answer
When physical time, or a new l from the max, moves past the previous l. c goes back to 0.
L3
Why does Spanner wait after choosing s?
Answer
s is TT.now().latest, the latest possible now. Waiting until s is definitely past means any later transaction picks a greater timestamp.
L4
How long is commit wait?
Answer
About twice the uncertainty. In the simulation, epsilon 7 ms waits 14.5 ms because the loop sleeps 0.5 ms.
L5
What is an uncertainty restart?
Answer
A read at t that sees a value in (t, t + max offset] bumps its read timestamp to that value and retries.
L6
Is serializable the same as externally consistent?
Answer
No. Serializable matches some serial order. External consistency also matches real time.
L7
What did the example max-offset guard print?
Answer
It refused a remote timestamp 10000 ms ahead of the local clock.
Failure modes
A fast clock drags every HLC forward
Without the guard, one node 10 seconds ahead pulls later receivers with it. The example raises once the gap exceeds 500 ms.
Commit wait on NTP-grade epsilon
An example 50 ms uncertainty is about 100 ms added to every read-write commit.
Max offset set blindly
Too low and nodes shut down on ordinary noise. Too high and uncertainty windows and restarts grow.
Misconceptions
HLC gives external consistency.
It preserves causality and stays near wall time. Real-time order across machines needs commit wait or a different protocol. CockroachDB still allows a causal reverse inside max offset.
Commit wait is the same as an uncertainty restart.
Commit wait delays the writer until the timestamp is in the past. A restart moves the reader's timestamp forward when a value falls in the uncertainty window.
Serializable isolation already matches wall-clock order.
It matches some serial order. External consistency also requires that order to respect real time.
Interviewer traps
Promising Spanner's external consistency on NTP.
Price two times epsilon. If epsilon is tens of milliseconds, say you would use HLC and uncertainty restarts instead.
Ignoring a node that is far in the future.
Name the max-offset guard and CockroachDB's self-termination at 80 percent of max offset.
Design scenario
Same prompt for every reader.
Requirements
Commits respect causality. Reads of one key should not miss a write that might have happened before them. State whether you claim external consistency.
Traffic / scale
Example skew: one node 80 ms fast, one 40 ms slow. Example max offset: 500 ms.
Latency
Commit wait of two times a multi-millisecond epsilon is a product decision. Read restarts should be rare.
Consistency
HLC gives the causal MVCC stamp. Uncertainty restarts cover the single-key window. Full external consistency is a stronger claim.
Availability
A node past the offset bound should stop rather than serve stale reads. A timestamp oracle outage stops transactions.
Failure assumptions
- Physical clocks disagree by tens of milliseconds.
- One node can jump far ahead.
- A reader can observe a write stamped slightly in its future.
Constraints
- Do not stamp MVCC versions with the raw wall clock.
- Do not claim strict serializability if you only implemented HLC.
Prompt
A geo-distributed SQL database needs MVCC timestamps. You do not have GPS and atomic clocks in every region.
API
What does a commit return, and does the client wait for commit wait or for a restart?
Data
What pair or interval do you store as the version?
Architecture
Where is the max-offset check, and when would you add a timestamp oracle instead?
Who pays for real-time order
Prefer
A bound you can name
TrueTime waits about two times epsilon. CockroachDB refuses to serve if a node is too far from the cluster and restarts reads inside the offset.
- The example HLC preserves causal order and stays 80 ms from true time, the fast node's skew.
- A 10 second future stamp is refused by the max-offset guard.
- Commit wait in the example is just over two times epsilon because the simulation steps 0.5 ms.
Alternative
Plain wall-clock MVCC timestamps
A slow node can commit in the past, hide a write from readers, or overwrite newer data.
- There is no logical counter when physical time stalls.
- There is no wait to make the stamp a past time everywhere.
- A central oracle avoids that lie and adds a round trip.
From a physical reading to a safe commit
The sequence is the HLC trace. The flowchart is Spanner's commit wait, including the self-loop while earliest is still at or below s.
- 1
Advance (l, c)
Physical time ahead of l resets c to 0. Otherwise c increments. A message can pull l forward. - 2
Refuse an implausible future
If the sender is more than max offset ahead, do not drag the cluster. The example refuses a stamp 10000 ms ahead. - 3
Commit-wait or restart the read
Spanner waits until s is in the past. CockroachDB moves the read timestamp past a value inside the uncertainty window. - 4
Name the anomaly you still allow
CockroachDB does not claim full strict serializability across unrelated keys. A causal reverse is possible inside max offset.
Overview
Distributed databases want MVCC timestamps that are causally ordered (like Lamport) and close to real time (so "read as of 10:00" and TTLs make sense). A hybrid logical clock (HLC) gives both: a physical component that tracks the largest wall time seen, plus a small logical counter for ties. HLC alone still cannot promise that a transaction which started after another committed (in real time, possibly on a different machine) gets a later timestamp. Spanner closes that gap with TrueTime, a clock API that returns an uncertainty interval, and commit wait, which waits out the uncertainty before acknowledging. CockroachDB, without atomic clocks, uses HLC plus a configured maximum clock offset, and restarts reads that hit a value inside their uncertainty interval.
Hybrid logical clocks
An HLC timestamp is (l, c): l is the largest physical time this node has seen (its own or in a message), c counts events that happened while l stayed the same.
- Local event or send: if physical time
pt > l, set l = pt and c = 0; otherwise keep l and increment c. - Receive message (lm, cm): l' = max(l, lm, pt). If l' equals both the old l and lm, c = max(c, cm) + 1; if it equals only the old l, c + 1; if only lm, cm + 1; otherwise 0.
- Compare timestamps as (l, c) pairs: the clock condition (a -> b implies
ts(a) < ts(b)) holds like Lamport. - Bound: l never exceeds the maximum physical clock in the system, so with bounded skew it stays within that bound of true time, and c stays small.
Sequence
- 1
"Node C (80 ms fast)"
"Step 1: send at physical 1080, HLC (1080, 0)"
- 2
"Node C (80 ms fast)" → "Node A (accurate)"
"message carries (1080, 0)"
- 3
"Node A (accurate)"
"Step 2: physical 1005 is behind, keep l = 1080, HLC (1080, 1)"
- 4
"Node A (accurate)" → "Node B (40 ms slow)"
"Step 3: send (1080, 2)"
- 5
"Node B (40 ms slow)"
"Step 4: physical 970, still adopt l = 1080, HLC (1080, 3)"
- 6
"Node B (40 ms slow)"
"Step 5: later physical 1160 passes l, HLC resets to (1160, 0)"
- 7
"Node A (accurate)"
"Step 6: causality preserved, l stays within max skew of true time"
Lesson map
Hybrid Logical Clocks & TrueTime - Commit Wait, Uncertainty Intervals & External Consistency
HLC algorithm (l, c) with skewed nodes and a max-offset guard (runnable); CockroachDB max-offset self-termination, MongoDB cluster time, YugabyteDB; TrueTime intervals and commit wait (runnable) plus CockroachDB-style uncertainty restarts; external consistency vs serializability vs SI; TrueTime vs HLC vs timestamp oracle vs single leader.
Architecture. Architecture
Select a node to see why it exists, or an edge to see the protocol, direction, effect, and consequence.
Mermaid export
flowchart TB c["Node C (80 ms fast)"] a["Node A (accurate)"] b["Node B (40 ms slow)"] c -->|message carries (1080, 0)| a a -->|Step 3: send (1080, 2)| b
HLC in code (runnable)
# Hybrid Logical Clock (Kulkarni et al., 2014): timestamp = (l, c).
# l tracks the largest physical time seen; c is a small counter that breaks ties
# when physical time has not moved past l. Causality is preserved like Lamport,
# but l stays within the clock-skew bound of real time, so it is human readable.
from dataclasses import dataclass
MAX_OFFSET_MS = 500 # like CockroachDB's default --max-offset (example guard value)
@dataclass
class HLC:
name: str
skew_ms: int # this node's physical clock error (example values)
l: int = 0
c: int = 0
def pt(self, true_ms: int) -> int:
return true_ms + self.skew_ms
def now(self, true_ms: int):
"""Local event or send."""
p = self.pt(true_ms)
if p > self.l:
self.l, self.c = p, 0
else:
self.c += 1
return (self.l, self.c)
def recv(self, true_ms: int, m):
ml, mc = m
p = self.pt(true_ms)
if ml - p > MAX_OFFSET_MS: # sender is implausibly far in the future: refuse
raise ValueError(f"{self.name}: remote ts {ml} is {ml-p} ms ahead of local clock")
new_l = max(self.l, ml, p)
if new_l == self.l == ml: self.c = max(self.c, mc) + 1
elif new_l == self.l: self.c += 1
elif new_l == ml: self.c = mc + 1
else: self.c = 0
self.l = new_l
return (self.l, self.c)
A = HLC("A", skew_ms=0)
B = HLC("B", skew_ms=-40) # 40 ms slow
C = HLC("C", skew_ms=+80) # 80 ms fast
trace = []
m = C.now(1000); trace.append(("C send", 1000, m))
r = A.recv(1005, m); trace.append(("A recv", 1005, r))
m2 = A.now(1006); trace.append(("A send", 1006, m2))
r2 = B.recv(1010, m2); trace.append(("B recv", 1010, r2))
r3 = B.now(1011); trace.append(("B local", 1011, r3))
r4 = B.now(1200); trace.append(("B local", 1200, r4))
print(f"{'event':8} {'true':>5} {'phys':>5} hlc(l,c)")
for ev, t, ts in trace:
node = {"A": A, "B": B, "C": C}[ev[0]]
print(f"{ev:8} {t:>5} {node.pt(t):>5} {ts}")
print("causal order preserved:", trace[0][2] < trace[1][2] < trace[2][2] < trace[3][2] < trace[4][2])
print("max (l - true time) observed:", max(ts[0] - t for _, t, ts in trace), "ms (bounded by worst skew)")
# A node whose clock jumped 10 s ahead must not drag the whole cluster's HLC with it.
bad = HLC("D", skew_ms=10_000)
try:
A.recv(1300, bad.now(1300))
except ValueError as e:
print("guard:", e)Output:
event true phys hlc(l,c)
C send 1000 1080 (1080, 0)
A recv 1005 1005 (1080, 1)
A send 1006 1006 (1080, 2)
B recv 1010 970 (1080, 3)
B local 1011 971 (1080, 4)
B local 1200 1160 (1160, 0)
causal order preserved: True
max (l - true time) observed: 80 ms (bounded by worst skew)
guard: A: remote ts 11300 is 10000 ms ahead of local clockThe guard at the end matters in practice: one node with a wildly wrong clock would otherwise drag every HLC in the cluster into the future. CockroachDB's real protection is cluster-wide: nodes continuously measure their offset to peers, and a node whose offset to a majority of the cluster exceeds 80% of --max-offset (default 500 ms) shuts itself down rather than risk violating consistency.
Who uses HLC
| System | What the hybrid clock is used for |
|---|---|
| CockroachDB | MVCC transaction timestamps; uncertainty intervals bounded by max offset |
| MongoDB | Cluster time ($clusterTime, operationTime) for causally consistent sessions since 3.6, signed so clients can't advance it arbitrarily |
| YugabyteDB | HLC-based MVCC timestamps across tablets |
TrueTime and commit wait
TrueTime's API is TT.now() returning an interval [earliest, latest] guaranteed to contain true time, plus TT.after(t) and TT.before(t). Google bounds the uncertainty (epsilon) with GPS receivers and atomic clocks in every datacenter; the Spanner paper reports epsilon typically varying between about 1 and 7 ms in a sawtooth between syncs.
- A read-write transaction acquires locks and picks commit timestamp s = TT.now().latest.
- Commit wait: the leader waits until TT.after(s) is true, that is, until s is definitely in the past everywhere.
- Only then does it release locks and acknowledge the client.
- Any transaction that starts after the acknowledgment must see
TT.now().latest > s, so its timestamp is larger: external consistency. - Read-only transactions and snapshot reads pick a timestamp and read lock-free from any sufficiently up-to-date replica.
Decisions
- 1
Step 1: acquire locks
- nextStep 2: s = TT.now().latest
- 2
Step 2: s = TT.now().latest
- nextStep 3: is TT.now().earliest > s yet?
- ?
Step 3: is TT.now().earliest > s yet?
- no, keep waiting about 2 x epsilonStep 3: is TT.now().earliest > s yet?
- yes, s is in the past everywhereStep 4: release locks and acknowledge
- 4
Step 4: release locks and acknowledge
- nextStep 5: any later transaction gets a timestamp > s (external consistency)
- 5
Step 5: any later transaction gets a timestamp > s (external consistency)
Commit wait and uncertainty restarts (runnable)
// Part 1: Spanner-style commit wait with a TrueTime-like interval clock (simulated).
// TT.now() returns [earliest, latest]; true time is guaranteed to be inside.
// Part 2: CockroachDB-style read uncertainty restart with a max clock offset.
// All epsilon / offset values are EXAMPLE values.
class TrueTimeSim {
constructor(public trueMs: number, private epsMs: number) {}
now(): { earliest: number; latest: number } {
return { earliest: this.trueMs - this.epsMs, latest: this.trueMs + this.epsMs };
}
sleep(ms: number): void { this.trueMs += ms; }
}
function commit(tt: TrueTimeSim): { ts: number; waitedMs: number } {
const s = tt.now().latest; // pick commit ts no earlier than any possible "now"
let waited = 0;
while (tt.now().earliest <= s) { // commit wait: until s is definitely in the past
tt.sleep(0.5); waited += 0.5;
}
return { ts: s, waitedMs: waited }; // only now release locks / acknowledge the client
}
console.log("Part 1: commit wait ~ 2 x epsilon");
for (const eps of [1, 4, 7]) {
const tt = new TrueTimeSim(10_000, eps);
const t1 = commit(tt);
// T2 starts strictly after T1 was acknowledged, possibly on another machine.
const t2 = commit(tt);
console.log(` eps=${eps}ms T1.ts=${t1.ts} waited ${t1.waitedMs}ms T2.ts=${t2.ts} T1<T2: ${t1.ts < t2.ts}`);
}
console.log("\nPart 2: uncertainty interval restart (max offset 500 ms)");
const MAX_OFFSET = 500;
// Versions of key k written by other nodes, stamped with THEIR clocks.
const versions = [{ ts: 9_700, v: "old" }, { ts: 10_300, v: "new" }];
function read(readTs: number): string {
let ts = readTs;
for (;;) {
const uncertaintyLimit = ts + MAX_OFFSET;
const visible = versions.filter((x) => x.ts <= ts).sort((a, b) => b.ts - a.ts)[0];
const uncertain = versions.find((x) => x.ts > ts && x.ts <= uncertaintyLimit);
if (!uncertain) return `read@${ts} -> ${visible?.v ?? "nothing"}`;
// The 'future' write may really have happened before us (its writer's clock ran ahead),
// so we cannot ignore it: bump our read timestamp past it and retry.
console.log(` read@${ts}: value at ${uncertain.ts} is within uncertainty window (${ts}, ${uncertaintyLimit}] -> restart`);
ts = uncertain.ts;
}
}
console.log(" ", read(10_000));
console.log(" ", read(11_000));Output:
Part 1: commit wait ~ 2 x epsilon
eps=1ms T1.ts=10001 waited 2.5ms T2.ts=10003.5 T1<T2: true
eps=4ms T1.ts=10004 waited 8.5ms T2.ts=10012.5 T1<T2: true
eps=7ms T1.ts=10007 waited 14.5ms T2.ts=10021.5 T1<T2: true
Part 2: uncertainty interval restart (max offset 500 ms)
read@10000: value at 10300 is within uncertainty window (10000, 10500] -> restart
read@10300 -> new
read@11000 -> newExpectedPart 1: commit wait ~ 2 x epsilon eps=1ms T1.ts=10001 waited 2.5ms T2.ts=10003.5 T1<T2: true eps=4ms T1.ts=10004 waited 8.5ms T2.ts=10012.5 T1<T2: true eps=7ms T1.ts=10007 waited 14.5ms T2.ts=10021.5 T1<T2: true Part 2: uncertainty interval restart (max offset 500 ms) read@10000: value at 10300 is within uncertainty window (10000, 10500] -> restart read@10300 -> new read@11000 -> new
Press Run. Snippets must be self-contained — no network, files, or native modules.
Part 1 shows the cost: commit latency grows with uncertainty (just over 2 x epsilon here because the simulation polls in 0.5 ms steps). Part 2 shows CockroachDB's alternative: instead of every writer waiting, a reader that sees a value with a timestamp slightly in its future cannot know whether that write really happened first (the writer's clock may run ahead), so it moves its read timestamp up and retries. Most reads never hit this; contended keys written just before a read do.
External consistency vs serializability
| Guarantee | Meaning | Example system |
|---|---|---|
| Serializable | Result equals some serial order of transactions | PostgreSQL SERIALIZABLE (single node), CockroachDB |
| Strict serializable / external consistency | Serial order also respects real time: if T1 committed before T2 started, T1 comes first | Spanner; a single-node database trivially |
| Snapshot isolation | Reads from a consistent snapshot; write skew allowed | Many MVCC databases by default |
CockroachDB provides serializable isolation and, thanks to uncertainty restarts, does not show stale reads for single keys, but it does not claim full strict serializability across unrelated keys (a "causal reverse" anomaly is possible when clocks are skewed within max offset). Spanner pays commit wait to rule that out.
TrueTime vs HLC vs timestamp oracle
| Approach | How ordering is guaranteed | Latency cost | Dependency | Failure mode if clocks misbehave |
|---|---|---|---|---|
| TrueTime + commit wait | Interval bound + waiting | About 2 x epsilon per read-write commit | Tight, well-monitored clock infrastructure | If the bound is violated, external consistency is lost; Spanner's design treats bad clocks as failed machines |
| HLC + max offset (CockroachDB) | Causal HLC + uncertainty restarts | Occasional read restarts; no commit wait | NTP-grade sync within max offset | Exceeding max offset can produce stale reads, so nodes self-terminate |
| Central timestamp oracle (Percolator, TiDB PD) | One service hands out increasing timestamps | One round trip to the oracle per transaction (often batched) | Highly available oracle | Oracle outage stops transactions; cross-region latency |
| Single leader | Leader's log order | Every write goes through one node | Leader availability | Failover pause |
AWS Time Sync with ClockBound, and PTP-based fleets, make TrueTime-style bounded clocks more accessible outside Google, which is why "small epsilon + commit wait" is no longer purely a Google-only design point.
What happens if you choose otherwise
- Plain wall-clock MVCC timestamps: a transaction on a slow-clock node can commit "in the past" and be invisible to readers that should see it, or overwrite newer data.
- Commit wait with NTP-grade epsilon (say 50 ms, example value): about 100 ms on every write, which most OLTP workloads cannot afford.
- Max offset set too low: nodes self-terminate on normal clock noise; too high: more uncertainty restarts and larger stale windows. Monitor clock offset as a first-class SLO.
- Timestamp oracle in one region for a global database: every transaction elsewhere pays a cross-region round trip.
Interview Q&A
What problem does HLC solve that Lamport clocks don't?
Answer
Lamport timestamps are causally ordered but unrelated to real time. HLC keeps the causal guarantee while staying within the clock-skew bound of physical time, so timestamps work for MVCC snapshots, TTLs and "as of" queries.
Why does Spanner wait after choosing a commit timestamp?
Answer
The chosen timestamp is the latest possible current time. Waiting until it is definitely in the past ensures any transaction that starts after the commit is acknowledged, anywhere, picks a larger timestamp.
How long is commit wait and how do you make it shorter?
Answer
Roughly twice the clock uncertainty. You shrink it by improving clock infrastructure (local GPS/atomic references, PTP), and Spanner overlaps it with replication so it is often partly hidden.
What is a read uncertainty restart in CockroachDB?
Answer
When a read at timestamp t sees a value with a timestamp in (t, t + max offset], it cannot tell whether that write happened before the read in real time, so it bumps its read timestamp past the value and retries.
Is serializability the same as external consistency?
Answer
No. Serializable means equivalent to some serial order; external consistency (strict serializability) also requires that order to respect real time across transactions.
What do l and c mean?
Answer
l is the largest physical time this node has seen, from its own clock or a message. c counts events that happened while l stayed the same. Compare the pair lexicographically.
When does CockroachDB stop a node?
Answer
Nodes measure offset to peers. A node whose offset to a majority exceeds 80 percent of max-offset, 500 ms by default, shuts itself down rather than risk the consistency bound.
Why is a 50 ms epsilon a poor commit-wait budget?
Answer
The wait is about twice epsilon, so an example 50 ms NTP-grade uncertainty is about 100 ms on every read-write commit. The lesson says most OLTP workloads cannot afford that.
Check yourself
Pick epsilon of 1, 4, and 7 ms as in the snippet. State the commit timestamp and the wait before you look at the output. Then say whether your system would rather wait on the writer or restart the reader.
Elsewhere in the library
These pages stay as they are. This lesson only points at them: MVCC, Snapshot Isolation & Write Skew, Two-Phase Commit — Protocol, Coordinator & Participants, Replication Lag & Session Guarantees — Read-Your-Writes, Monotonic Reads & Consistent Prefix, Raft Consensus — Leader Election, Log Replication & Safety.
Go Deeper
- Spanner: Google's Globally-Distributed Database (OSDI 2012)
- Kulkarni et al., Logical Physical Clocks and Consistent Snapshots in Globally Distributed Databases (HLC paper)
- Cockroach Labs: Living Without Atomic Clocks
- MongoDB: Causal Consistency and Read and Write Concerns
- Google Cloud: Spanner TrueTime and external consistency
- Jepsen: CockroachDB beta-20160829 analysis (clock skew findings)