Google's Spanner, described in the 2012 OSDI paper 'Spanner: Google's Globally-Distributed Database', was the first system to offer externally consistent transactions across data centres on different continents. If transaction T1 commits before T2 starts, as observed by anyone in the real world, then T1's commit timestamp is smaller than T2's. A snapshot read at any timestamp therefore sees a state that could have existed at an instant of real time. The ingredient that made this practical was TrueTime, a clock API that admits how wrong it might be.

Google Cloud's managed Spanner and its operational features are covered in Cloud Spanner. This page concentrates on the protocol. It covers the rules Spanner follows when picking timestamps, the short proof that commit wait gives external consistency, the safe-time machinery that lets replicas serve reads without locks, and a simulation you can run to watch the guarantee hold and fail. It ends with what to borrow when you do not have GPS receivers and atomic clocks in every data centre. All figures for TrueTime itself come from the 2012 paper. Present-day uncertainty bounds are not given here, because the paper's numbers are the ones that are documented.

What external consistency asks for

Start with the problem. A database that assigns commit timestamps from local clocks can order two transactions backwards. Suppose T1 commits on a server whose clock runs 5 ms fast, and a client that saw T1 commit immediately starts T2 on a server whose clock is accurate. T2 can receive the smaller timestamp. A snapshot read between the two timestamps then sees T2's effects without T1's, even though T2 started after T1 finished. If T1 removed a user from an access list and T2 posted a private photo, a reader can see the photo with the old access list.

External consistency rules this out. It is linearizability applied to transactions: the timestamp order must agree with real-time order for every pair of non-overlapping transactions. Linearizability versus sequential consistency covers the single-object version. Lamport clocks and hybrid logical clocks, described in hybrid logical clocks, preserve causal order when a message passes between the two servers. External consistency is stronger, because it must also hold when the only link between T1 and T2 is a person, a phone call or a separate system.

TrueTime: an interval, not a timestamp

TrueTime replaces 'what time is it?' with 'what range must the time be in?'. Its API has three calls:

CallReturnsGuarantee
TT.now()an interval [earliest, latest]the true absolute time of the call lies inside the interval
TT.after(t)true if t has definitely passedequivalent to t < TT.now().earliest
TT.before(t)true if t has definitely not arrivedequivalent to t > TT.now().latest

The half-width of the interval is the uncertainty epsilon (ε). In the 2012 deployment each data centre had a set of time masters. Most had GPS receivers, and the others, called Armageddon masters, had atomic clocks, so the two groups failed for unrelated reasons. A daemon on every machine polled several masters, used a variant of Marzullo's algorithm to discard liars, and synchronised its local clock. Between polls it assumed a worst-case local drift of 200 microseconds per second. With a 30-second poll interval that adds up to 6 ms, plus about 1 ms for communication delay to the masters, so ε followed a sawtooth from roughly 1 to 7 ms over each poll interval, averaging about 4 ms. Machines whose clocks drifted beyond the bound were evicted.

The engineering point is that ε is measured and bounded, not assumed. Everything that follows only needs the guarantee that true time lies inside the interval. The size of ε decides how much latency the guarantee costs, not whether it holds.

Commit wait and the proof

Spanner's read-write transactions follow two rules, which the paper calls Start and Commit Wait. Start: the coordinator leader picks a commit timestamp s that is no smaller than TT.now().latest, computed after the commit request arrives. Commit wait: the coordinator does not let any client see the transaction's effects until TT.after(s) is true, meaning s has definitely passed.

Commit wait: why T2 always gets a larger timestamptrue timeT1 leader: TT.now() at commit[e1, l1], picks s1 = l1uncertainty interval, width 2εs1commit waituntil TT.after(s1)T1 visibleT2 starts laterpicks s2 = TT.now().latests2s1 is below true time when T1 becomes visible; T2 starts after that, and its timestamp is at least true timeso s1 is below s2 for every pair of transactions where one finishes before the other starts
The leader picks s1 at the top of its uncertainty interval and holds the result until s1 is certainly in the past. Any transaction that starts after T1 becomes visible picks a timestamp no smaller than true time at that moment, so it lands above s1.

The proof is a chain of four inequalities, where t(e) is the true time of event e:

  1. s1 < t(T1 commit), by commit wait: T1 is not visible until s1 has passed.
  2. t(T1 commit) < t(T2 start), by assumption: T2 started after T1 finished.
  3. t(T2 start) ≤ t(T2 commit request arrives), by causality.
  4. t(T2 commit request arrives) ≤ s2, by the start rule and TrueTime's guarantee.

So s1 < s2. Nothing in the chain depends on the servers talking to each other, which is why the guarantee survives a phone call between the two clients. Commit wait lasts about 2ε, because the leader must wait until the bottom of its interval passes what was the top. In the paper, much of that wait overlaps the Paxos round needed to replicate the commit anyway, so the visible latency added is often smaller than the full 2ε.

Worked example: simulating commit wait

The simulation below gives each node a hidden clock error within ±ε and has it report TrueTime intervals honestly. T1 commits on one node and T2 starts on another up to 2 ms after T1 becomes visible. It counts how often s2 ≤ s1, which is a violation, with and without commit wait.

import random

EPS = 0.007          # 7 ms: the top of the 2012 sawtooth

class Node:
    def __init__(self):
        self.offset = random.uniform(-EPS, EPS)    # true error, unknown to the node

    def now(self, t):                              # TT.now() at true time t
        local = t + self.offset
        return (local - EPS, local + EPS)          # [earliest, latest]

def run(trials, commit_wait):
    violations, waits = 0, []
    for _ in range(trials):
        a, b = Node(), Node()
        t = random.uniform(0, 1000)
        s1 = a.now(t)[1]                           # start rule: TT.now().latest
        if commit_wait:                            # wait until TT.after(s1)
            start = t
            while a.now(t)[0] <= s1:
                t += 0.0001
            waits.append(t - start)
        t += random.uniform(0, 0.002)              # client sees commit, starts T2
        s2 = b.now(t)[1]
        violations += s2 <= s1
    return violations, (sum(waits) / len(waits) if waits else 0.0)

random.seed(1)
for cw in (False, True):
    v, w = run(100_000, cw)
    print(f"commit_wait={cw}: {v} of 100000 orderings violated; mean wait {w*1000:.1f} ms")
commit_wait=False: 43238 of 100000 orderings violated; mean wait 0.0 ms
commit_wait=True: 0 of 100000 orderings violated; mean wait 14.1 ms

Without commit wait, about 43 percent of back-to-back transactions on different nodes receive timestamps in the wrong order, even though every node reports honest intervals. With it, there are no violations, and the price is a wait of 2ε (14 ms at ε of 7 ms). Change EPS to 0.001 and the wait drops to about 2 ms. That is the whole economic argument for investing in better clocks: uncertainty turns directly into commit latency.

Read-write transactions across Paxos groups

Real transactions touch data in many Paxos groups, each with its own leader. Spanner runs two-phase commit across the groups, with one group's leader acting as coordinator. Timestamps are woven into the protocol:

  • Writers take locks at each participant leader. Deadlocks are avoided with wound-wait: an older transaction wounds, or aborts, a younger lock holder, and a younger one waits for an older one.
  • Each participant leader picks a prepare timestamp larger than any timestamp it has already assigned, and logs the prepare through Paxos.
  • The coordinator picks a commit timestamp s that is at least every participant's prepare timestamp, greater than any timestamp it has assigned, and at least TT.now().latest when it received the commit message. It then performs commit wait before telling participants and clients.
  • Each Paxos leader holds a time-based lease, 10 seconds by default in the paper. Leaders assign timestamps only within their lease interval, and Spanner keeps lease intervals of successive leaders disjoint. A new leader therefore cannot assign a timestamp lower than one its predecessor already used, so timestamps within a group increase monotonically across leader changes.

Commit wait and two-phase commit stack: a cross-group write pays a prepare round, a commit round and a wait of about 2ε, partly overlapped. Distributed transactions in practice compares this with the commit paths of other systems.

Reads without locks: safe time

TrueTime pays off most on reads. Because every write has a timestamp that respects real time, a read at timestamp t is a consistent snapshot. It needs no locks and can be served by any replica that is far enough along. Each replica tracks a safe time, t_safe, the highest timestamp at which it is certain it has every write:

t_safe = min(t_paxos_safe, t_tm_safe)

t_paxos_safe = timestamp of the highest applied Paxos write
t_tm_safe    = infinity if no transaction is prepared but uncommitted here,
               else min(prepare_ts of those transactions) - 1

serve read at t   if   t <= t_safe   else wait (or go to a fresher replica)

The second term matters. A transaction that has prepared but not committed might commit at any timestamp at or above its prepare timestamp, so a replica cannot answer reads beyond that point until the outcome is known. A read-only transaction picks its timestamp as TT.now().latest, which is the simplest choice that is guaranteed to respect real-time order, and then reads at that timestamp anywhere t_safe allows. For reads confined to one Paxos group, the leader can instead use the timestamp of its last committed write. That is often lower, so it waits less. Stale reads choose an older timestamp on purpose and can be served immediately by a nearby replica.

Schema changes use the same idea. A change is assigned a timestamp in the future, and transactions on either side of it use the matching schema, without stopping the world. Snapshot semantics in general, and the anomalies that remain at snapshot isolation, are covered in snapshot isolation. Spanner's read-write transactions use locks and are serializable, so they do not have write skew.

Failure modes

  • ε spikes. If a time master is unreachable or a machine's oscillator misbehaves, ε grows, and commit wait grows with it. Alert on uncertainty, not only on clock offset, because ε is the term that turns into latency.
  • A clock lying beyond its bound. If true time leaves the interval, external consistency is silently lost. That is why the design uses two independent time sources, cross-checks masters, and evicts machines whose drift exceeds the assumed rate.
  • Long prepared transactions. A participant stuck between prepare and commit pins t_tm_safe, and replica reads at newer timestamps block. A slow coordinator shows up as read latency on other replicas.
  • Leader loss. Disjoint leases mean a new leader may have to wait for the old lease to expire before assigning timestamps, which is a short write-unavailability window for that group.
  • Hot timestamps or keys. TrueTime orders transactions but does not spread load. Monotonically increasing keys still concentrate writes on one split.

Copying the idea without atomic clocks

Most systems lack Spanner's clock infrastructure, but the idea transfers in three forms. First, use bounded uncertainty without waiting. CockroachDB uses hybrid logical clocks with a configured maximum clock offset (500 ms by default). Instead of commit wait, a read that encounters a value inside its uncertainty window restarts at a higher timestamp, which gives serializability and avoids stale reads under the offset assumption, without full external consistency for causally unrelated transactions. Second, measure your bound. Background on synchronisation, drift and leases is in clocks in distributed systems, and AWS publishes ClockBound, a daemon that reports a clock-error bound applications can read. Third, wait when it matters. Applying commit wait only to the operations that need real-time ordering, such as revocations and fencing changes, buys the guarantee where it counts.

def commit_externally_consistent(txn, clock):      # clock.now() -> (earliest, latest)
    s = max(clock.now()[1], txn.max_prepare_ts, last_assigned_ts + 1)
    replicate_commit(txn, s)                       # overlaps the wait below
    while clock.now()[0] <= s:
        sleep(0.0005)
    release_locks_and_reply(txn, s)

Trade-offs

ApproachGuaranteeCost
TrueTime plus commit waitexternal consistency; lock-free snapshot reads anywherededicated time infrastructure; about 2ε on every read-write commit
HLC plus uncertainty restartsserializable, no stale reads within the offset boundrestarts under contention; no protection if clocks exceed the bound
Single timestamp oraclea strict order from one servicea round trip to the oracle for every transaction; a central bottleneck
Causal or session consistency onlycheap and highly availablereal-time anomalies between independent clients

What to do next

  1. Run the simulation, then vary EPS and the client delay to see where violations appear and how the wait scales.
  2. Measure clock offset and the error bound on your own fleet. Know your ε before relying on any timestamp ordering.
  3. List the operations in your system where real-time order matters, such as permission revocation or fencing, and decide how each is protected.
  4. If you use CockroachDB or a similar database, find its maximum-offset setting and the monitoring that kills nodes exceeding it.
  5. If you run Spanner, use stale or bounded-staleness reads where a few seconds of age is acceptable, and keep strong reads for paths that need them.
  6. Watch for long-running prepared transactions: they show up as blocked replica reads elsewhere.
  7. Reread the Start and Commit Wait rules and the four-step proof until you can reproduce them without notes.
Key takeaway: TrueTime returns an interval guaranteed to contain true time, with uncertainty of roughly 1 to 7 ms in the 2012 deployment. Spanner picks each commit timestamp at the top of the interval and waits until it has certainly passed before revealing the commit, and that is enough to prove that real-time order matches timestamp order. Prepare timestamps, disjoint leader leases and per-replica safe time extend the guarantee to cross-group transactions and lock-free snapshot reads. Without atomic clocks, measure your clock bound and use uncertainty restarts or targeted waits where real-time order matters.