Imagine you're sitting in a café in Tokyo. You tap "buy" on a pair of sneakers. Three hundred milliseconds later, a server in Iowa updates a ledger, a cache in Belgium invalidates a key, and a fraud model in Oregon decides you're not a bot. To you, it's magic. To the distributed systems engineer who built that stack, it's a minor miracle—one that required solving a problem that philosophers and computer scientists have been wrestling with for decades.
That problem is simple to state and brutal to solve: How do you make a distributed system behave as if it were a single machine, when its components are scattered across an ocean and separated by the speed of light?
For most of computing history, the answer was "you don't." Then Google built Spanner, and for a decade, the industry got used to a certain answer: TrueTime. Atomic clocks. GPS receivers. Uncertainty intervals. A beautiful, brutalist approach to time that only Google could afford.
But here's the thing nobody talks about at conferences: Google doesn't actually use TrueTime for everything anymore. In recent years, Spanner's internals have undergone a quiet but profound evolution—a migration from a clock-centric model to a Timestamp Oracle (TSO) model that borrows concepts from systems like Percolator and TiDB. And the reasons why are fascinating, because they reveal the real constraints of planetary-scale computing: not physics, but economics; not correctness, but latency.
Let's dig in.
The Original Sin of Distributed Systems: You Can't Trust Clocks
Before we can appreciate what Google did, we need to understand why time is such a nightmare in distributed systems.
Every machine has a clock. Every clock drifts. Even with NTP (Network Time Protocol), your two servers in the same datacenter might disagree by a few milliseconds. Across continents? Tens or hundreds of milliseconds. In a distributed system, that's an eternity.
Why does this matter? Consider a simple transaction:
- Client A writes
x = 1to Server 1 at "time" T1. - Client B reads
xfrom Server 2 at "time" T2. - If T1 < T2, we'd like Server 2 to see the write from Server 1.
But what if Server 1's clock is 50ms ahead of Server 2's? Then a write that happened "later" in real time might look "earlier" to the second server. You get causal inconsistencies: reading stale data after a write has already been acknowledged.
The traditional fix is to use logical clocks (Lamport timestamps, vector clocks) that don't depend on wall-clock time. But logical clocks have their own problem: they can't tell you when something happened in real time. For databases—especially ones handling financial transactions—that's a dealbreaker. You need real timestamps.
So you're stuck choosing between:
- Correctness (use logical clocks, accept poor real-time semantics)
- Real-time timestamps (use wall clocks, accept inconsistency)
- Strong consistency (use coordination protocols like Paxos, accept horrific latency)
Google's Spanner was supposed to be the system that refused to choose.
TrueTime: The Ferrari of Clocks
Spanner's original trick was audacious: make the clocks so accurate that the uncertainty is negligible.
TrueTime isn't a single clock—it's an API. TT.now() returns an interval [earliest, latest] guaranteed to contain the actual current time. Inside Google's datacenters, this interval is typically just a few milliseconds wide, sometimes under 1ms. How?
- Atomic clocks (cesium or rubidium) in every datacenter, providing a stable reference.
- GPS receivers in every datacenter, providing an absolute time signal from satellites.
- A daemon (
TTdaemon) in every machine that compares local clocks against these references and computes an uncertainty bound.
The magic was in Spanner's commit-wait protocol. When a transaction commits with timestamp s, Spanner waits until TT.now().earliest > s before acknowledging. This guarantees that any subsequent transaction anywhere on Earth will see a timestamp strictly greater than s. External consistency—the gold standard of distributed consistency—achieved with a few milliseconds of wait time.
It was genius. It was also expensive. And it had a scaling problem.
The Hidden Cost of TrueTime
Here's what the conference talks gloss over: TrueTime requires a physical infrastructure that is staggeringly expensive to replicate.
Atomic clocks cost tens of thousands of dollars each. GPS antennas need rooftop access and clear sky views. The whole apparatus needs constant calibration, monitoring, and failover. Google can afford this because they run their own datacenters, own their own fiber, and have an army of SREs. But even for Google, TrueTime has a subtler issue: the commit-wait latency is bounded by clock uncertainty, not by your actual workload.
If your TrueTime uncertainty is 7ms (which was common in early deployments), then every write transaction pays a 7ms tax—minimum. That's fine for batch workloads. It's murder for low-latency systems where users expect sub-10ms response times globally.
And here's the kicker: as Spanner scaled out to more regions and more workloads, the commit-wait tax became a bottleneck. Not because clocks got worse, but because the coordination overhead of synchronizing state across regions started to dominate.
Around 2017–2019, Google engineers started asking: Do we really need TrueTime for every transaction?
Enter the Timestamp Oracle (TSO)
The answer, it turned out, was no. And the alternative was hiding in plain sight: a centralized timestamp oracle.
The idea is simple. Instead of every node having its own hyper-accurate clock and coordinating via commit-wait, you have a single logical service that doles out monotonic timestamps. Want to commit a transaction? Ask the TSO for a timestamp. It increments a counter, returns the value, and you're done. No waiting. No uncertainty intervals. Just a number.
This is, of course, the model used by Google Percolator (Google's incremental processing system for search indexes), TiDB, and CockroachDB (to some extent). The TSO is the source of truth for "now."
Why is this faster? Because the TSO can hand out timestamps in microseconds. There's no commit-wait. The TSO doesn't need to be atomic-clock-accurate—it just needs to be monotonically increasing and globally unique. That's a much weaker requirement.
But wait—doesn't a centralized TSO reintroduce the single point of failure that TrueTime was designed to eliminate?
Yes. And that's the tradeoff. But modern TSO implementations are surprisingly resilient:
- Batched timestamp allocation: The TSO can hand out timestamp ranges in bulk (e.g., "you get 1,000,000 to 1,001,000"), so nodes don't need to call it for every transaction.
- Leader election with Paxos/Raft: If the TSO dies, a new leader takes over in milliseconds, with a persisted monotonic counter.
- Regional TSOs: You can run a TSO per region, synchronized asynchronously, for even lower latency.
The interesting thing is that Spanner's newer "TSO" mode is closer in spirit to Percolator than to the original TrueTime design. Google hasn't deprecated TrueTime—it's still there for the workloads that need the absolute strictest guarantees. But for most transactions, the TSO path is now the default.
The Architecture Nobody Talks About
Let's get concrete. Here's a simplified view of what Spanner's internals look like today, based on public papers and engineer talks.
The Old Model (TrueTime-Centric)
[Client] --> [Spanner Server] --> [TrueTime API] --> [Commit-Wait] --> [Replica] --> [ACK]
Every commit pays:
- A TrueTime lookup (network hop to TTdaemon, ~100 microseconds)
- A commit-wait delay (up to 2× uncertainty, often 5–10ms)
- Replication via Paxos (network round-trip to quorum, 50–200ms cross-region)
The New Model (TSO-Centric)
[Client] --> [Spanner Server] --> [TSO RPC] --> [Timestamp] --> [Replica] --> [ACK]
Every commit pays:
- A TSO RPC (network hop to TSO leader, ~1ms in-region)
- Replication via Paxos (same as before)
- No commit-wait
The savings are enormous. For a transaction that writes to a single region, the old model might take 15ms just in clock overhead. The new model takes 1ms. That's a 15× improvement in tail latency, and it compounds across every transaction in a workload.
But here's where it gets really interesting: the TSO model is not just faster—it's more flexible.
Why TSO Wins for Cross-Region Commits
Cross-region commits are the hardest case in distributed systems. You're coordinating state between machines 10,000 km apart, where the speed of light itself adds 60–80ms of round-trip latency. Every millisecond of overhead is brutal.
The TrueTime approach handles this by ensuring that no matter how far apart two replicas are, their clocks are within a known bound. But that bound is global—if one datacenter's GPS antenna is flaky, the entire system's uncertainty grows.
The TSO approach handles it differently: the timestamp is assigned once, at the source. A transaction that commits in Tokyo gets a timestamp from the TSO. A transaction that commits in Iowa gets a timestamp from the same TSO. Since the TSO is monotonic, the ordering is guaranteed—no matter what the clocks say.
This has a subtle but powerful implication: the commit-wait is replaced by a TSO round-trip, which can be optimized independently of clock uncertainty.
Specifically:
- Regional TSOs with clock skew bounds: Google can run a TSO per region, with each TSO synchronized via a lightweight protocol (like a simplified TrueTime just for the TSO itself). The TSO hands out timestamps in batches, and the batch size can be tuned based on observed clock skew.
- Lease-based allocation: A region can "lease" a range of timestamps for a fixed duration (say, 1 second). During that window, it can allocate timestamps locally without contacting the global TSO. When the lease expires, it renews. This reduces cross-region RPCs to once per second per region, rather than once per transaction.
- Failover with bounded staleness: If the global TSO fails, a new leader is elected, but regions can continue serving reads (and in some cases, writes) using their leased ranges. The system degrades gracefully instead of halting.
This is the architecture that makes low-latency cross-region commits possible. Not by eliminating the speed of light, but by stopping paying for it on every transaction.
The Engineering Curiosity: Why Now?
If TSO is so obviously better, why didn't Google do this from the start?
Three reasons:
1. Hardware Wasn't Ready
In 2012, when the Spanner paper was published, the idea of running a reliable, low-latency TSO at Google's scale was untested. Atomic clocks were the known solution. The TSO model had been used in Percolator (2010), but that was for batch processing, not for a globally-distributed OLTP database. Nobody knew if it would scale.
2. The Workload Was Different
Early Spanner workloads were internal Google services—Ads, Play, etc.—that could tolerate a few milliseconds of commit-wait. As Spanner became a product (Cloud Spanner), external customers started demanding sub-10ms latencies for cross-region transactions. The economics changed.
3. The TSO Model Needed New Primitives
Modern TSO implementations rely on:
- Raft/Paxos-based leader election (well-understood by 2015)
- Lease-based timestamp allocation (refined by systems like CockroachDB)
- Batched RPCs (standard practice by 2018)
- Monotonic counters with crash recovery (solved by WAL + fsync)
By 2020, all the pieces existed. The migration was inevitable.
The Tradeoffs (Because There Are Always Tradeoffs)
I want to be careful here not to oversell TSO. It's not a free lunch. Here's what you give up:
- Single point of contention: The global TSO is a hot path. Even with batching, at extreme scale (millions of transactions per second), the TSO can become a bottleneck. Google mitigates this with regional TSOs and leases, but it's still a coordination point.
- Failure modes are different: TrueTime fails "safe" (uncertainty grows, commit-wait increases). TSO fails "hard" (no timestamps, no commits). The failover story is more complex.
- Clock skew between TSOs: If you run multiple TSOs, you need to bound their skew. This is where TrueTime still shows up—not for every transaction, but for TSO synchronization. The irony is delicious: TrueTime didn't go away; it moved up the stack.
- Correctness depends on TSO availability: In a TrueTime system, a node with a broken clock can still participate (it just has wider uncertainty). In a TSO system, a node that can't reach the TSO is dead in the water for writes.
Google's answer to all of these is the same: hybrid mode. Some transactions use TrueTime, some use TSO, and the system can switch based on workload characteristics. This is the real story: not a replacement, but a coexistence.
What This Means for the Rest of Us
If you're building a distributed system—whether it's a database, a ledger, or a globally-distributed cache—the Spanner story offers some hard-won lessons:
1. Don't Over-Engineer for Correctness
TrueTime is beautiful, but it's overkill for 90% of workloads. If your application can tolerate a few milliseconds of inconsistency (and most can), a TSO or even a simple Lamport clock will do. Save the atomic clocks for when you really need them.
2. Centralization Isn't Always Bad
The distributed systems community has spent decades demonizing central coordination. But a well-designed TSO can be faster than a fully decentralized approach, because it eliminates the need for consensus on every operation. Coordination is expensive; sometimes it's cheaper to just ask a central authority.
3. Latency Is a Feature
The original Spanner paper focused on correctness. The newer Spanner focuses on latency. That shift reflects a broader truth in systems engineering: once you have correctness, performance is the only thing that matters. Users don't care that your system is externally consistent if it takes 500ms to respond.
4. Hybrid Systems Win
The future isn't "TrueTime vs. TSO." It's both, with the system choosing the right tool for each transaction. This is the pattern we see everywhere in modern infrastructure: polyglot persistence, multi-modal databases, adaptive consistency. The best systems are the ones that can change their mind.
The Code (Because You Asked)
Here's a rough sketch of what a TSO client might look like in Go, inspired by CockroachDB's implementation:
type TSOClient struct {
conn *grpc.ClientConn
batchSize int
cached []int64
mu sync.Mutex
}
func (c *TSOClient) GetTimestamp(ctx context.Context) (int64, error) {
c.mu.Lock()
defer c.mu.Unlock()
if len(c.cached) > 0 {
ts := c.cached[0]
c.cached = c.cached[1:]
return ts, nil
}
// Request a batch of timestamps from the TSO
resp, err := c.conn.RequestTimestamps(ctx, &pb.TimestampRequest{
Count: int64(c.batchSize),
})
if err != nil {
return 0, err
}
c.cached = resp.Timestamps[1:]
return resp.Timestamps[0], nil
}
And here's the commit-wait equivalent in the TrueTime world:
func (s *SpannerServer) Commit(ctx context.Context, txn *Transaction) error {
ts := s.tt.Now()
// Write to Paxos group
if err := s.replicate(ctx, txn, ts); err != nil {
return err
}
// Wait until TrueTime guarantees the timestamp has passed
for {
now := s.tt.Now()
if now.Earliest.After(ts.Latest) {
break
}
time.Sleep(time.Millisecond)
}
return s.ack(ctx, txn)
}
The difference is stark: the TSO path has no loop, no sleep, no waiting. It's just a number and a write. That's the whole point.
The Bigger Picture: Time Is a Service
The most profound insight from Spanner's evolution is this: time is not a property of the physical world—it's a service you provide.
TrueTime provided time as a hardware-accelerated service, with atomic clocks and GPS. TSO provides time as a software service, with monotonic counters and leases. Both are valid. Both are correct. But they have different cost profiles.
And in the end, cost is what drives architecture. Google didn't move away from TrueTime because it was wrong. They moved toward TSO because it was cheaper—in latency, in hardware, in operational complexity.
This is the lesson for anyone building distributed systems: the theoretically optimal solution is rarely the one that ships. The one that ships is the one that balances correctness, latency, cost, and complexity in a way that matches the workload.
Spanner's journey from TrueTime to TSO is a masterclass in that balance. It's not a story about abandoning a beautiful idea. It's a story about layering abstractions—keeping TrueTime where it's needed, adding TSO where it's better, and letting transactions flow through the fastest path available.
And that, in the end, is what global consensus really looks like: not a single, perfect clock, but a carefully orchestrated dance of timestamps, leases, and acknowledgments—all tuned to make a tap in Tokyo feel instantaneous.
