SysDesignPrep.com
Study guide 81 of 183

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 placementWrite latency (roughly)Survives
3 zones in one regiona few millisecondsa zone failure
3 regions on one continenttens of millisecondsa region failure
5 regions across continents100+ millisecondsa 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.

Open in your browser to sign in

Google does not allow sign-in inside this app's built-in browser. Open this page in Safari and sign in there. The link opens this same page.

Tap the ⋯ or share button at the top or bottom of the screen, then Open in browser. Or copy the link and paste it into Safari.