Spanner / CockroachDB (distributed SQL)
SQL tables and transactions spread across ranges, with Paxos or Raft quorums replacing the single primary bottleneck.
Also worth naming: Google Spanner · Cloud Spanner · CockroachDB · distributed SQL · NewSQL
Distributed SQL is for the moment when one primary database is the bottleneck or the regional failure domain is too small. The trade is simple to say and hard to operate: you keep SQL and serializable transactions, and you pay quorum latency plus placement complexity.
What it is
Distributed SQL databases keep the SQL interface and transactional model, but split storage across many machines. The key idea is the same in both systems: a table is ordered by key, cut into splits or ranges, and each piece is replicated by consensus rather than by one primary streaming to passive replicas.
Spanner stores contiguous key ranges called splits and replicates each split with synchronous, Paxos-based voting; Google's docs say every write needs a majority of voting replicas before it commits, and the OSDI 2012 paper introduced Spanner as a globally distributed, synchronously replicated database with externally consistent transactions. CockroachDB stores a sorted key-value map divided into ranges; its docs say ranges are replicated by Raft, and range splits create new Raft groups as tables grow.
Use this page together with consistency trade-offs, replication and durability, leader election, and sharding. Distributed SQL is not magic Postgres. It is sharding plus consensus plus transaction coordination, packaged behind SQL. If one Postgres primary comfortably handles the write load and your failure target is one region, Postgres is still the simpler and often better answer.
When to reach for it
Reach for this when…
- You need SQL transactions over data that no longer fits one primary write leader comfortably.
- You need synchronous replication and zero-data-loss failover across nodes or zones, and can pay quorum latency.
- You need multi-region OLTP with data placement by region, not just read replicas.
- You want developers to keep relational modelling and SQL while the database handles range movement and consensus.
Not really this pattern when…
- A single Postgres primary plus read replicas meets the workload and recovery target; start with /learn/technologies/postgres.
- The hot path needs very low same-region latency and every extra network round trip hurts.
- The workload is append-heavy analytics, search, or event streaming rather than OLTP.
- The team cannot operate retries, schema locality, hot ranges, and regional placement yet.
How it works
Think in five layers.
1. The shard is a range, not a hand-written database split. Spanner calls the unit a split: a contiguous row-key range that can move independently and is replicated with Paxos, according to the Spanner reads and writes whitepaper. CockroachDB calls it a range: a contiguous chunk of a sorted KV map, with metadata stored in meta ranges and automatic range splits and merges. This agrees with sharding: range layout preserves scans, but the primary key can create hot tails.
2. Each range has a consensus owner for writes. In Spanner, a Paxos leader handles writes for a split; Spanner docs describe read-write, read-only, and witness replicas, and say write requests go first to the leader and commit after a write quorum agrees (replication docs). In CockroachDB, a range has a leaseholder colocated with the Raft leader; the leaseholder serves strong reads and proposes writes, while a majority of voting replicas coordinates the write (replication layer). Leader and leaseholder placement is why data locality matters.
3. Serializable transactions are real, but they are implemented differently. Spanner's default serializable isolation provides external consistency: if transaction A finishes before transaction B starts to commit, every client observes A before B. The mechanism is TrueTime plus commit timestamps; Spanner's docs say TrueTime gives monotonically increasing timestamps, and the write path performs commit-wait until the chosen timestamp is safely in the past before replying (TrueTime docs, reads and writes whitepaper).
A leader picks a commit timestamp, reaches a Paxos quorum, then waits out clock uncertainty before replying. That commit-wait is what makes external consistency visible to clients.
CockroachDB does not use specialised clock hardware. Its transaction layer uses MVCC timestamps from hybrid logical clocks, timestamp caches, read refreshing, and uncertainty handling. The current docs say strongly consistent reads go through the leaseholder, HLC timestamps track MVCC versions, and a read inside an uncertainty window can push a timestamp or lead to a retry (transaction layer). Retry errors that reach the client use SQLSTATE 40001 and include restart transaction in the message (retry error reference).
Current CockroachDB stable docs also say the default isolation level is SERIALIZABLE, while READ COMMITTED is supported as an option, is enabled by default via the sql.txn.read_committed_isolation.enabled setting, and was introduced as preview support in CockroachDB 23.2 (transactions, read committed, 23.2 launch). It trades away some serializable guarantees to reduce client-side retries.
CockroachDB uses hybrid logical clocks and MVCC snapshots. If a read sees a value that might be in its clock uncertainty interval, the transaction is pushed or retried so serializable ordering is preserved.
4. Multi-region is placement policy
CockroachDB's multi-region table localities are explicit: regional tables optimise reads and writes from one home region, regional by row optimises individual rows for their home region, and global tables optimise low-latency reads from every region while writes incur higher latency (table localities). Spanner chooses a regional, dual-region, or multi-region instance configuration; multi-region configs place voting replicas across regions, designate read-write regions and a default leader region, and can add read-only replicas that do not join write quorums (instance configurations).
Regional rows keep writes local to a home region. Global or multi-region reference data gives broad reads, but writes pay cross-region quorum or timestamp costs.
5. The latency bill is a write quorum. A write waits for a quorum for its range; if the quorum crosses regions, the network path crosses regions. Spanner docs say multi-region reads can be faster in more places while writes pay extra network latency because voting replicas are spread across regions (instance configurations). CockroachDB docs make the same shape visible: writes need Raft quorum, and global tables improve reads from all regions while writes are slower from any one region (replication layer, table localities).
Performance envelope
Distributed SQL performance envelope. Treat these as design constraints, not benchmark promises; actual numbers depend on regions, schema, contention, replicas, workload, and hardware.
| Dimension | Spanner | CockroachDB | Design consequence |
|---|---|---|---|
| Replication unit | [Split](https://cloud.google.com/spanner/docs/schema-and-data-model#database-splits), a contiguous key range | [Range](https://docs.cockroachlabs.com/docs/stable/architecture/distribution-layer#range-splits), a contiguous KV span | Primary-key order decides locality and hotspot risk. |
| Consensus | [Paxos-based synchronous replication](https://cloud.google.com/spanner/docs/replication) | [Raft group per range](https://docs.cockroachlabs.com/docs/stable/architecture/replication-layer#raft) | A write is not acknowledged until a voting quorum accepts it. |
| Regional fault tolerance | Base regional configs use [three read-write replicas and a 2-of-3 write quorum](https://cloud.google.com/spanner/docs/instance-configurations#replication). | With [3x replication, one failure can be tolerated](https://docs.cockroachlabs.com/docs/stable/architecture/replication-layer#overview). | Good fit for zero-data-loss node or zone failures, with a quorum latency cost. |
| Multi-region writes | [Voting replicas across regions add write latency](https://cloud.google.com/spanner/docs/instance-configurations#multi-region-configurations). | [Regional rows keep writes local; global tables make writes slower](https://docs.cockroachlabs.com/docs/stable/table-localities#when-to-use-regional-vs-global-tables). | Place write-heavy data near its writers; do not make hot mutable rows global. |
| Low-latency remote reads | [Stale reads with enough staleness can avoid a leader round trip](https://cloud.google.com/spanner/docs/replication#read-only-replicas). | [Follower reads use closed timestamps](https://docs.cockroachlabs.com/docs/stable/architecture/transaction-layer#closed-timestamps). | Great for read-mostly data where bounded staleness is acceptable. |
| Hotspot avoidance | [Monotonic first key parts create hotspots](https://cloud.google.com/spanner/docs/schema-design#choose-a-primary-key-to-prevent-hotspots). | [Sequential keys create hotspots; hash-sharded indexes can help](https://docs.cockroachlabs.com/docs/stable/performance-best-practices-overview#unique-id-best-practices). | Randomise or prefix the key when write distribution matters. |
Capabilities in interviews
Spanner for externally consistent global SQL
Use Spanner when the requirement is SQL transactions whose commit order matches real-world observation.
Spanner's distinctive promise is external consistency. The current docs say serializable isolation is the default and provides external consistency, where transactions appear serial and their order matches the order clients observe them to commit (isolation levels). TrueTime exposes clock uncertainty; commit-wait makes the chosen commit timestamp safe before the client hears success (TrueTime docs).
That is the honest reason to name Spanner: you need globally ordered SQL transactions more than you need the lowest possible write latency. If you only need regional ACID with read replicas, Postgres stays simpler.
Choose this variant when
- Strict global ordering matters to the business.
- You can choose Spanner instance configurations up front.
- The team is comfortable with schema design around primary-key locality.
CockroachDB for Postgres-compatible distributed SQL
Use CockroachDB when you want distributed SQL with a PostgreSQL-like wire protocol and Raft-replicated ranges.
CockroachDB accepts SQL on any node, converts it to KV operations, routes those operations to the right ranges, and commits through Raft. Its architecture docs describe a PostgreSQL-compatible SQL API, ranges replicated to nodes by default, and symmetric reads and writes from each node while routing to the responsible leaseholders (architecture overview).
The interview phrasing is precise: CockroachDB is not a normal Postgres cluster with clever failover. It is a distributed KV store with a SQL layer, MVCC transactions, range leaseholders, and Raft quorums underneath.
Choose this variant when
- You want horizontal write scaling with SQL.
- PostgreSQL compatibility matters, but exact Postgres behaviour is not assumed.
- You can implement transaction retry logic and test contention.
Regional and row-level placement
Put each row near the users who write it, rather than making every write global.
The practical multi-region shape is often regional ownership. CockroachDB exposes this directly: regional tables have one home region, regional-by-row tables assign rows to home regions, and global tables are for read-mostly data (table localities). Spanner exposes the placement through regional, dual-region, and multi-region instance configurations, plus leader-region choices for write-heavy workloads (instance configurations).
Use this with multi-region active-active: active-active usually works when records have a single writer region by construction.
Choose this variant when
- Users or tenants have a natural home region.
- Residency and latency both matter.
- The write path can route to the owner region.
Global read-mostly reference tables
Replicate small, rarely changed data broadly so every region can read it cheaply.
Global placement is best for reference data: feature flags, product catalog slices, country lists, pricing rules, or promotion codes. CockroachDB docs describe global tables as optimised for low-latency reads from every region, with higher write latency because writes must support that global read pattern (global tables). Spanner read-only replicas can scale reads and support low-latency stale reads without joining write quorums (replication docs).
Do not put a hot counter, stock row, or seat inventory row here unless you are prepared to pay for every global update.
Choose this variant when
- The table is read often and changed rarely.
- All regions need local reads.
- Stale reads are acceptable for some callers.
Operating knobs
Spanner vs CockroachDB
Pick Spanner when managed Google Cloud placement, TrueTime-backed external consistency, and Spanner-specific schema choices are acceptable. Pick CockroachDB when PostgreSQL compatibility, self-managed or multi-cloud deployment, and Raft-based ranges fit the organisation better. Both need careful keys, retries, and locality design.
Key and locality design
The primary key is now a physical design decision. Spanner docs warn that a monotonically increasing first key part sends inserts to one hot split; CockroachDB docs warn that sequential IDs create hotspots in a distributed database. Use UUIDv4, hash prefixes, tenant or region prefixes, or composite keys where the first component spreads load.
Isolation choice and retry policy
Use serializable by default for correctness-critical flows. In CockroachDB, READ COMMITTED arrived as preview support in 23.2 and current stable docs enable it as an optional isolation level, but SERIALIZABLE remains the default. It can reduce client-side retries by permitting anomalies that serializable prevents. If you enable it, name which transactions can tolerate those anomalies and keep serializable for inventory, balances, uniqueness, and bookings.
Region topology
Start with where writes originate. If most writes are in one region, a regional setup or single Postgres primary can be better. If each tenant has a home region, use row-level or table-level locality. If every write must be visible globally immediately, say the cross-region quorum cost out loud.
Operational ownership
Distributed SQL moves complexity from application sharding into the database, but the complexity still exists. You operate hot ranges, schema changes, retry storms, regional failover, lease or leader placement, and backup restore. If your team cannot debug those yet, the simpler Postgres answer may win.
Versus the alternatives
Distributed SQL against the real alternatives.
| Dimension | Single Postgres primary | Spanner | CockroachDB |
|---|---|---|---|
| Write scaling | One primary write leader; scale reads with replicas and partition tables. | Ranges split automatically; each split writes through Paxos. | Ranges split automatically; each range writes through Raft. |
| Consistency | ACID on one primary; sync replica optional for durability. | Default serializable isolation gives external consistency via TrueTime. | Default SERIALIZABLE isolation with MVCC, HLCs, leaseholders, and retries. |
| Multi-region | Usually active-passive or async read replicas unless you shard deliberately. | Regional, dual-region, and multi-region instance configs with leader regions. | Regional, regional-by-row, and global table localities. |
| Latency profile | Lowest for local writes on one primary. | Writes wait for quorum and commit-wait; cross-region quorums cost more. | Writes wait for Raft quorum; cross-region and global-table writes cost more. |
| Best interview answer | Default until scale or resilience disproves it: see /learn/technologies/postgres. | When external consistency and managed global SQL are the point. | When distributed Postgres-compatible SQL and explicit regional locality are the point. |
Failure modes & gotchas
Auto-increment IDs, timestamp-first primary keys, or ordered event IDs send new writes to the end of the keyspace. Spanner docs call this a common hotspot cause, and CockroachDB docs discourage sequential indexed keys for the same reason. Put a high-cardinality prefix first, use UUIDv4, hash-prefix the key, or use CockroachDB hash-sharded indexes when a sequential index is unavoidable.
A SQL transaction can look local in code while touching ranges whose leaders or leaseholders live in different regions. The result is extra network round trips at commit time. In a design, name the rows that participate in the transaction and where their leaders live.
Global or broadly replicated placement is attractive for reads, but each update must preserve the global read guarantee. Use it for read-mostly reference data. Do not put counters, inventory, or a popular seat row there unless the write rate is low enough for the extra latency.
Serializable distributed transactions can abort under contention. CockroachDB current docs say retry errors returned to the client use SQLSTATE 40001 and include restart transaction; SERIALIZABLE transactions that cannot be retried automatically require client-side retry handling. If the application treats a retry error as a hard failure, correctness survives but availability and user experience suffer.
Consensus keeps replicas of a range in agreement. It does not choose your tenant boundaries, cool a hot key, make scatter-gather joins cheap, or fix a transaction that spans too many hot ranges. You still design the primary key, indexes, locality, and transaction boundaries.
A single Postgres primary is often easier to reason about, cheaper to run, and faster for local writes. If the product needs ACID, joins, modest write scale, and one-region HA, saying “Postgres first, distributed SQL later if the write leader becomes the limit” is a stronger answer than over-building.
In production
Spanner-style global ledger (illustrative)
Externally consistent account transfers
Imagine a regulated wallet product where account balances are read in several regions, but the ledger must behave as one serial database. A Spanner-style design keeps account rows in key ranges, commits transfers through Paxos, and uses TrueTime-style commit timestamps so an audit read cannot observe a later debit without an earlier deposit that the customer already saw complete.
CockroachDB regional SaaS tenancy (illustrative)
Tenant rows homed by region
Imagine a collaboration product with customers in several jurisdictions. Tenant-owned rows use regional-by-row placement so edits route to the tenant home region. Shared reference tables are global because they are read far more than written. Cross-tenant analytics runs outside the OLTP path. The design avoids pretending every write is global while still giving a single SQL database abstraction to product engineers.
Good vs bad answer
Numbers in this section are illustrative.
Interviewer probe
“The product is global and needs SQL transactions. Would you use Spanner, CockroachDB, or Postgres?”
Weak answer
"Use distributed SQL because it is globally consistent and scales automatically. Put all tables in a global database so every region can read and write."
Strong answer
"I start by checking whether one Postgres primary is enough. If the write workload and RPO fit, Postgres is simpler and lower latency. I move to distributed SQL only if the single write leader or regional failure domain is the real constraint.
If I need managed global serializable transactions and I am on Google Cloud, Spanner is a candidate: splits are Paxos-replicated, and TrueTime plus commit-wait gives external consistency. If I need PostgreSQL-compatible distributed SQL or multi-cloud deployment, CockroachDB is a candidate: ranges are Raft-replicated, a leaseholder serves strong reads, and SERIALIZABLE is the default.
I would not make every table global. User-owned rows get a home region, reference tables can be global or read-only, and hot inventory rows stay near their writers. Then I design the primary key to avoid monotonic hotspots and build retry handling for serializable conflicts."
Why it wins: It starts with the simpler Postgres baseline, distinguishes Spanner from CockroachDB by mechanism, and treats region placement, key design, and retries as the real design work.
Interview playbook
When it comes up
- The prompt asks for SQL transactions plus horizontal write scale.
- The system must survive a node, zone, or regional failure without acknowledged-write loss.
- The interviewer asks how to keep multi-region writes consistent.
- You are tempted to say “just shard Postgres” but also need cross-shard transactions.
Order of reveal
- 11. Start from Postgres. I would keep a single Postgres primary if it meets write scale and recovery goals. Distributed SQL earns its cost only when one write leader or one region is the constraint.
- 22. Name the replication unit. The table is split into key ranges: Spanner splits or CockroachDB ranges. Each range has its own consensus group.
- 33. Name the consensus path. Spanner writes through Paxos; CockroachDB writes through Raft. A write waits for a voting quorum, so cross-region writes cost cross-region latency.
- 44. Explain transaction ordering. Spanner uses TrueTime and commit-wait for external consistency. CockroachDB uses HLC timestamps, MVCC, uncertainty intervals, read refresh, and retries to provide SERIALIZABLE isolation by default.
- 55. Place data deliberately. Regional rows near their writers, read-mostly reference data global or read-only, and hot mutable rows kept out of global placement unless the write rate is low.
Signature phrases
- ““Distributed SQL is sharding plus consensus behind SQL.”” — Prevents hand-wavy “it scales automatically” claims.
- ““A write waits for quorum; geography is in the latency.”” — Makes the PACELC cost concrete.
- ““Spanner uses TrueTime commit-wait; CockroachDB uses HLC uncertainty and retries.”” — Shows you can compare the two honestly.
- ““Postgres first unless the write leader or region is the real limit.”” — Signals judgement rather than vendor reflex.
Likely follow-ups
?“How does CockroachDB differ from Spanner on clocks?”Reveal
Spanner has TrueTime, an API with bounded clock uncertainty backed by Google infrastructure. It assigns commit timestamps and waits until the timestamp is definitely in the past before replying, which gives external consistency. CockroachDB uses hybrid logical clocks on ordinary machines. It preserves serializable isolation with MVCC timestamps, timestamp caches, uncertainty intervals, and retries rather than TrueTime commit-wait.
?“Why are cross-region writes slower?”Reveal
Because the write must be accepted by a quorum of voting replicas for the range. If the voters are in multiple regions, the leader or leaseholder waits on wide-area network latency. Read-only replicas or stale follower reads can make reads local, but they do not remove the quorum cost for writes.
?“When would you still pick Postgres?”Reveal
When one write primary handles the workload, the product can tolerate regional active-passive DR, and the team benefits from the simpler operational model. Postgres gives excellent local write latency, mature tooling, and fewer distributed failure modes. I would spend distributed SQL only when write scale, zero-data-loss failover, or multi-region placement demands it.
Worked example
Numbers in this section are illustrative.
Scenario. You are designing a global ticketing platform. Users browse worldwide, but each event has a fixed venue and a seat can be sold once.
Baseline. If sales are regional and the write peak fits one primary, use Postgres: one writer for seat state, row-level locks or conditional updates, read replicas for browsing, and a queue for payment side effects. That is simpler and has lower local write latency.
Distributed SQL shape. Move to distributed SQL when the seat inventory write leader or failure target no longer fits. Partition by event_id first, because the hot transaction is “hold or sell this seat for this event”. Place each event's inventory in the event's home region. That keeps the seat-hold transaction local even if buyers are global. User profiles can be regional by user, and read-mostly reference data such as venue maps or tax rules can be global or served by stale reads.
Spanner option. In Spanner, event inventory lives in splits whose leaders are near the event's write region. The hold transaction commits through Paxos and TrueTime external consistency, so two buyers cannot both observe success for the same seat. The price is quorum latency and commit-wait.
CockroachDB option. In CockroachDB, event inventory rows live in ranges with leaseholders in the home region. The hold transaction runs at SERIALIZABLE by default. If another transaction races it, CockroachDB may return a retryable SQLSTATE 40001 / restart transaction error. Application code treats retry errors as normal and re-runs the hold attempt.
What you say in the interview. “I am not making every table global. Seat inventory is regional to the event and serializable. Browsing data can be replicated broadly. A write waits for consensus, so a buyer far from the event may see higher checkout latency, but the business rule is one seat, one sale.”
Cheat sheet
- •Distributed SQL = SQL + automatic ranges/splits + consensus replication.
- •Spanner: splits, Paxos, TrueTime, commit-wait, external consistency.
- •CockroachDB: ranges, Raft, leaseholders, HLCs, uncertainty intervals, SERIALIZABLE by default.
- •A write waits for quorum; cross-region quorum means cross-region latency.
- •Regional rows for hot writes; global/read-only placement for read-mostly data.
- •Monotonic first key parts create hot ranges.
- •READ COMMITTED was introduced as a 23.2 preview option; current CockroachDB enables it as an option, but SERIALIZABLE remains the default.
- •Postgres is still better when one primary is enough.
Drills
Numbers in this section are illustrative.
Why is “just make the table global” a bad default for a hot inventory row?Reveal
Global placement optimises broad reads, but writes still need to preserve the global guarantee. A hot mutable row then pays quorum or timestamp costs on every update. Put the inventory row near its writer or owner region, and replicate read-only views separately when needed.
What is the difference between a CockroachDB leaseholder and a Raft leader?Reveal
The leaseholder is the replica that serves strong reads and proposes writes for a range. CockroachDB colocates the leaseholder with the Raft leader to reduce round trips. Raft leadership is the consensus role; the lease is the read/write serving authority clients route through.
What does Spanner commit-wait buy you?Reveal
It waits until the commit timestamp chosen from TrueTime is definitely in the past before the client hears success. That makes commit timestamp order line up with real-time observation, which is the external-consistency guarantee.
When should you answer Postgres instead of distributed SQL?Reveal
When one primary handles sustained writes, local write latency matters, and the recovery target fits same-region synchronous replication plus async DR. You can still partition tables, add read replicas, and shard later. Do not pay distributed consensus complexity before the single-primary limit is real.
What it is