Spanner, TrueTime and distributed SQL
How globally distributed SQL databases work: ranges replicated with Paxos or Raft, distributed transactions with two-phase commit, Spanner's TrueTime and commit wait, hybrid logical clocks in CockroachDB and YugabyteDB, latency costs, data placement, and when to use them.
Reading is half of it. See this used in a real interview: walk through Design a Payment System →
Traditional relational databases scale up on one machine; NoSQL databases scale out but give up transactions and SQL. Distributed SQL databases (Google Spanner, CockroachDB, YugabyteDB, TiDB) aim for both: SQL, ACID transactions and strong consistency across many machines and even continents. They come up when a design needs global scale and strict correctness (payments, inventory, ledgers). Knowing how they work explains their latency costs and when they are worth it.
The architecture
- Data is split into ranges (or tablets) by primary key, each a few hundred megabytes.
- Each range is replicated, typically 3 or 5 times across zones or regions, using Paxos or Raft. Each range has its own consensus group and leader. See how Raft works.
- Ranges split and move automatically as data grows and load shifts, like an automatically sharded database. See sharding and partitioning.
- A SQL layer on every node plans queries and routes reads and writes to the right ranges.
Transactions across ranges
A transaction touching one range commits through that range's consensus group. A transaction touching several ranges uses two-phase commit coordinated over the consensus groups: since every participant is itself replicated, a coordinator failure does not leave the transaction stuck the way classic 2PC can. See distributed transactions and idempotency.
Concurrency control uses multi-version storage (MVCC) with timestamps, so the ordering of transactions depends on timestamps being meaningful across machines.
Spanner and TrueTime
Spanner's key idea: TrueTime, a clock API backed by GPS receivers and atomic clocks in every data centre, which returns an interval [earliest, latest] guaranteed to contain the true time, with uncertainty usually a few milliseconds.
- A transaction gets a commit timestamp, then waits out the uncertainty (commit wait) before making the commit visible.
- After that wait, any transaction that starts later, anywhere in the world, is guaranteed to get a larger timestamp.
- Result: external consistency (strict serializability): transaction order matches real-time order globally, and consistent snapshot reads anywhere without locks.
The cost is a few milliseconds of commit wait per write, made small by investing in precise clocks.
Without atomic clocks
CockroachDB and YugabyteDB run on ordinary hardware, using hybrid logical clocks (physical time plus a logical counter) and a configured maximum clock offset. When a read encounters a value with a timestamp in its uncertainty window, it may need to restart at a higher timestamp. They provide serializable isolation, with linearizability per key, but not Spanner's full external consistency in every case. Clocks that drift beyond the configured limit are a real risk, so nodes shut themselves down when they detect excessive offset. See unique ids, ordering and time.
Latency is geography
Consensus needs a majority of replicas, so write latency is at least the round trip to the nearest majority:
| Replica placement | Write latency (roughly) | Survives |
|---|---|---|
| 3 zones in one region | a few milliseconds | a zone failure |
| 3 regions on one continent | tens of milliseconds | a region failure |
| 5 regions across continents | 100+ milliseconds | a region failure, with global reads |
Reads from the leader are fast if the leader is near the client; reads elsewhere either go to the leader, or use follower reads at a slightly stale timestamp.
Data placement
To keep latency low, place data near its users:
- Pin leaders for a table or partition to the region where most writes happen.
- Partition by region (geo-partitioning): European users' rows live in Europe, which also helps with data residency. See privacy and data deletion.
- Global tables for small, read-mostly reference data, replicated everywhere for fast local reads at the cost of slower writes.
See multi-region architecture.
When to use distributed SQL
Good fits:
- Strong consistency and transactions needed at a scale beyond one primary database. See scaling a relational database.
- Multi-region active-active writes with correctness (financial ledgers, inventory, user accounts).
- Avoiding the complexity of application-level sharding.
Less good:
- Small workloads where a single Postgres with replicas is simpler and faster.
- Latency-critical writes that cannot afford consensus round trips.
- Very high-volume append workloads (metrics, logs) better served by wide-column or time-series stores. See data modelling for Cassandra and DynamoDB.
Hot spots still matter: monotonically increasing keys send all inserts to one range. Use hashed or random keys. See hot keys and skew.
In the interview
For a global payment system or ledger: "Accounts and ledger entries live in a distributed SQL database, geo-partitioned so each account's leader is in its home region; transfers within a region commit in a few milliseconds; cross-region transfers pay a consensus round trip; serializable isolation prevents anomalies." For a ticketing system with global sales, explain why a single-region primary for each event may be simpler.
Checklist
- Ranges replicated by consensus, split and moved automatically.
- Cross-range transactions via 2PC over replicated participants.
- TrueTime commit wait (Spanner) or hybrid logical clocks with uncertainty restarts.
- Write latency set by the distance to a majority.
- Leader pinning, geo-partitioning and global tables for locality.
- Choose it for correctness at scale; avoid hot sequential keys.