Why This Matters#
Distributed SQL is not "Postgres that scales". It is a latency contract with physics. You get SQL, serializable transactions and automatic sharding. In return, every write goes through a consensus round trip, and every transaction that touches more than one range or region pays for the distance between replicas. Spanner and CockroachDB hide the router, the resharding and the two-phase commit. They do not hide the speed of light. A row in Virginia that must survive the loss of Virginia cannot be committed until a replica in another region has acknowledged it, and that takes 10–70 ms however good the engineering is.
That is why "we'll just use Spanner" is a sentence interviewers lean into. The L5 candidate says it to make sharding and consistency go away. The L6 candidate says "Spanner in a regional configuration, 3 read-write replicas across zones, so a commit costs one in-region Paxos round of 5–10 ms. Primary keys are UUIDv4 so inserts don't all land on the last range. The order and its line items are interleaved under the customer so a checkout is a single-split transaction, and the analytics queries go to a separate store." The L7 candidate asks what the company is buying: one global database for 40 teams, or the right to stop hand-sharding Postgres. Those are different bills, different lock-in and different failure postures.
The L5 → L6 gap is not knowing that Spanner uses TrueTime. It is knowing that the database chooses where every byte lives and how many network hops every commit takes, and your schema and locality settings are the only levers you have over both. Choose the primary key and locality well and distributed SQL feels like a single-node database. Choose badly and you have built a slow, expensive single-node database with a hot range.
The L5 → L6 → L7 Contrast#
| Behavior | Senior (L5) | Staff (L6) | Principal (L7) |
|---|---|---|---|
| First move | "Use Spanner/CockroachDB so we don't have to shard" | "What is the write latency budget, and where are the writers? That decides regional vs multi-region and which tables are regional, per-row or global." | "Is this one team's database or the company's default OLTP platform? That's a vendor, cost and lock-in decision measured over five years." |
| Keys | Auto-increment IDs, like Postgres | UUIDv4 or hash-prefixed keys to spread inserts; interleave or colocate child rows under the parent | Publishes a key-design standard and a lint, because one sequential key in a shared cluster becomes everyone's hot range |
| Consistency | "It's strongly consistent, so we're done" | Knows the cost: serializable means retries (40001), contention means aborts; uses stale or follower reads where 5–15 s staleness is fine | Decides which data classes deserve global consistency and which belong in cheaper stores; doesn't make the whole company pay commit latency for logs |
| Multi-region | "Spread it across 3 regions" | Prices each table: home-region rows for user data, global tables for read-mostly config, region survival only where RPO 0 across regions is required | Writes the data-residency and survival policy per data class, signed off by legal and the business |
| Failure | "Consensus handles node failures" | Names what happens: a leaseholder dies, its ranges are unavailable for seconds until a lease moves; losing a region with zone survival makes its home rows unavailable | Designs the org's failure posture: which outage the business accepts, and what the 5-replica bill buys against it |
| Ownership | "The DBA team runs it" | Product teams own schemas, keys and query plans; the platform owns the cluster, upgrades, backups and capacity | Decides managed vs self-hosted, the exit path, and whether a single shared cluster is a blast radius the company will accept |
Why "Keys" separates levels
In a single-node database, an auto-increment primary key is the best choice: inserts append to the right edge of one B-tree that already sits in memory. In a range-partitioned database, the same key sends every insert to the last range in the key space, and so to one leaseholder on one node. The cluster can have 30 nodes and the insert path still runs at the speed of one. Range splits don't help: the database splits the hot range, and the new right-hand range is immediately the hot one. Spanner's schema documentation names this exact anti-pattern and recommends UUIDv4, a hashed key prefix, or bit-reversed sequences. CockroachDB offers hash-sharded indexes for the same reason. The Senior answer carries a single-node habit into a distributed system. The Staff answer chooses the key for write distribution and locality at the same time. The Principal answer makes that a reviewed standard, because a hot range in a shared cluster pages the platform team, not the team that caused it.
Why "Multi-region" separates levels
"Spread it across 3 regions" sounds like resilience. It is actually a choice about write latency, and the bill can arrive as a 60 ms regression on every checkout. In CockroachDB, a database with the default ZONE survival goal keeps each row's voting replicas in its home region: writes are fast locally, but if that region goes down, its rows are unavailable until it comes back. Switching to REGION survival moves the default to 5 replicas and, by the documentation's own description, raises write latency by at least the round trip to the nearest other region. Spanner's multi-region configurations make the same trade at the instance level: 99.999% availability in exchange for cross-region quorum on every write. The Staff answer chooses survival per database and locality per table, and says what each costs in milliseconds. The Principal answer recognises that "survive a region loss with zero data loss" is a business requirement with a price, and gets someone to sign it.
The 60-Second Pitch#
"The core is orders, payments and inventory: about 8K writes per second at peak, with multi-row invariants across customers and merchants, so hand-sharding Postgres would turn every cross-shard transfer into a saga. I'd put it on distributed SQL, CockroachDB or Spanner, in one primary region with 3 replicas across zones: commits cost one in-region consensus round, about 5–10 ms, and transactions are serializable. Primary keys are UUIDs so inserts spread across ranges, and child rows are colocated with their parent so most transactions touch one range. Reporting reads use follower or stale reads a few seconds behind, so they never contend with writers, and heavy analytics goes to a column store through CDC. I'm not making it multi-region on day one. When EU users need local writes, I'd make user-owned tables REGIONAL BY ROW so each row lives in its user's region, and keep global tables for read-mostly config. What I'm accepting: higher per-write latency and cost than a single Postgres primary, and client-side retry handling for serialization errors."
The Three Intents#
| Intent | Constraint | Strategy | Failure Mode | Correctness Bar |
|---|---|---|---|---|
| Outgrown the single primary | One region, writes beyond one Postgres primary, cross-entity transactions | Regional cluster (3 replicas across zones), UUID keys, colocation, retries | Hot range from a sequential key; retry storms on contended rows | Serializable; p99 commit < 20 ms in-region |
| Global users, local latency | Users in 3+ regions, residency rules, writes from everywhere | Multi-region with per-row home regions; global tables for reference data | Cross-region transactions sneak into the hot path; uniqueness checks fan out | Local-region p99 for home-region traffic; residency enforced by placement |
| Survive a region with RPO 0 | Financial or ledger data, no acknowledged write may be lost | Region survival (5 replicas, quorum across regions) or Spanner multi-region | Every commit pays a 20–70 ms cross-region round; costs roughly 2× the replicas | Zero acknowledged-write loss on full region failure |
🎯 Staff Move: "I'll design for the first intent: we've outgrown one Postgres primary, but all our writers are in one region. That means a regional configuration with zone survival, and I'll keep the schema shaped so that going multi-region later is a locality change per table, not a migration. Surviving a whole region with zero data loss is a separate requirement with a separate latency bill, and I'd want the business to ask for it explicitly."
The Staff Positions#
| Position | Rationale |
|---|---|
| Choose the primary key for distribution first | Sequential keys put every insert on one range and one leaseholder. UUIDv4, a hash prefix or a hash-sharded index spreads writes. |
| Design so most transactions touch one range | Single-range commits skip the distributed commit path. Interleave or prefix child rows with the parent key. |
| Regional by default; multi-region per table, by need | Every replica in another region is milliseconds on every write. Pay it only for data that needs it. |
| Serializable plus a retry loop in the client | Both systems abort conflicting transactions rather than return anomalies. Code that does not retry 40001 errors fails at peak. |
| Stale reads for anything that tolerates seconds | Follower and stale reads are served from the nearest replica without contending with writers. |
| Not for analytics, queues or blobs | Consensus-replicated row storage is the most expensive byte you own. Scans belong in a column store; blobs in object storage. |
| Plan the exit before the entry | Spanner is a proprietary managed service; CockroachDB's licensing has changed before. Schema and access code should stay as portable as the features you rely on allow. |
Architecture & Internals#
Both systems share one shape: a SQL layer on top of a transactional, sorted key-value store that is split into ranges, with each range replicated by its own consensus group. Spanner calls the unit a split and replicates it with Paxos. CockroachDB calls it a range and replicates it with Raft. The details differ; the design consequences are the same.
Ranges, Splits and Placement#
The table's primary key is the sort order of the key-value store. Rows are stored in key order, and contiguous key spans are cut into ranges. CockroachDB's current documentation lists a default maximum range size of 512 MiB and a minimum of 128 MiB, with 3 replicas per range by default (5 for system ranges) (CockroachDB replication zones). A range that grows past the maximum splits in two. Both systems also split on load: Spanner's documentation describes adding split boundaries between frequently accessed rows so each lands on a different server (Spanner schema and data model). A rebalancer then moves replicas and leases to even out disk and load across nodes.
What this changes in the design:
- Keys decide placement. Rows with adjacent keys live together; rows with distant keys are on different ranges, often different nodes. Colocation is something you build into the key.
- One range, one writer. Each range has one leaseholder (Spanner: the Paxos leader) that orders its writes. That is a per-range throughput ceiling of a few thousand writes per second in practice, whatever the cluster size.
- Splits are cheap, merges are lazy. The database reshards on its own. You never write a resharding runbook, but you also don't decide when it happens.
Per-Range Consensus and Leaseholders#
A write is durable once a majority of its range's replicas have it: 2 of 3, or 3 of 5. Reads go to the leaseholder, which can answer from its own copy without a quorum round because the lease guarantees no other replica is accepting writes. In current CockroachDB versions the Raft leader and the leaseholder are the same replica (leader leases); the documentation says detecting a node failure and moving the lease to a new node "should complete within a few seconds" (CockroachDB replication layer).
The two facts that matter in an interview:
- A write costs one round trip from the leaseholder to the nearest majority. With 3 replicas in 3 zones of one region, that's 1–2 ms of network plus disk sync, so 5–10 ms end to end in practice. With replicas in 3 regions, it is the round trip to the second-closest region.
- Losing a leaseholder costs seconds, not data. The range is unavailable for writes, and for leaseholder reads, until a new lease is established. That is a p99 spike for the ranges on that node, not an outage. Size client timeouts and retries for it.
Transactions: TrueTime vs Hybrid Logical Clocks#
Both systems give every transaction an MVCC timestamp, and both need clocks to make timestamps from different machines comparable. They diverge on how much they trust those clocks.
Spanner: bounded uncertainty, wait it out. TrueTime returns an interval [earliest, latest] guaranteed to contain true time. The OSDI 2012 paper reports that in production the uncertainty ε was a sawtooth between about 1 and 7 ms, so about 4 ms most of the time, driven by time masters with GPS and atomic clocks and a 30-second polling interval (Spanner paper, OSDI 2012). A read-write transaction picks a commit timestamp, then commit-waits until TT.now().earliest is past that timestamp before releasing locks and acknowledging the client. After the wait, any transaction that starts later is guaranteed a larger timestamp, which is external consistency (strict serializability). The wait overlaps with Paxos replication, so in practice it adds little on top of the consensus round.
CockroachDB: commodity clocks, an uncertainty window, and restarts. CockroachDB uses hybrid logical clocks (physical time plus a logical counter) on ordinary NTP-synchronised machines. It cannot bound uncertainty to milliseconds, so it assumes a configured maximum clock offset, --max-offset, with a default of 500 ms; a node whose clock drifts beyond the limit crashes rather than risk serving inconsistent data (cockroach start). Instead of making every write wait, it makes reads deal with uncertainty: if a read at timestamp t finds a value written in (t, t + max_offset], it can't tell whether that write happened before it, so it moves its timestamp forward and may have to restart. Transactions are SERIALIZABLE by default, and conflicts that cannot be resolved surface as retryable errors with SQLSTATE 40001 (CockroachDB transaction retry errors).
| Spanner (TrueTime) | CockroachDB (HLC) | |
|---|---|---|
| Clock assumption | Hardware-backed interval, ε about 1–7 ms (2012 paper) | NTP-class clocks, configured max offset (500 ms default) |
| Who pays for uncertainty | Writers: commit wait of about 2ε | Readers: uncertainty restarts on recently written keys |
| Guarantee | External consistency (strict serializability) | Serializable; weaker than Spanner's global real-time ordering between unrelated transactions |
| Failure if clocks misbehave | Widened ε, slower commits | Node self-terminates past the offset limit |
| Interview sentence | "Spanner spends milliseconds of commit wait to buy global ordering." | "Cockroach trades occasional read restarts for not needing atomic clocks." |
🎯 Staff Insight: You don't need to explain TrueTime's algorithm to pass. You need to say who pays for clock uncertainty: in Spanner it's every write, a few milliseconds; in CockroachDB it's reads that land on recently written keys, as occasional restarts. That tells the interviewer you know where the p99 comes from.
Distributed Transactions Across Ranges#
A transaction whose writes all fall in one range commits with one consensus round. A transaction across ranges needs an atomic commit protocol: Spanner runs two-phase commit across Paxos groups, with one group's leader as coordinator; CockroachDB writes provisional intents plus a transaction record. CockroachDB's parallel commits, on by default since v19.2, cut the cross-range commit from two sequential consensus rounds to one by marking the transaction record STAGING with the list of in-flight writes (CockroachDB parallel commits).
Commit latency, back-of-envelope (design numbers, not vendor SLAs):
single-range write, 3 replicas in one region ~ 1 consensus round -> 5-10 ms
multi-range write, same region ~ 1-2 rounds + intents -> 10-20 ms
write with quorum spanning regions ~ RTT to 2nd-closest region + disk
us-east + us-central + us-west (quorum = 2 of 3) -> 25-40 ms
client far from leaseholder/leader + 1 RTT client -> leader region
EU client writing a US-homed row +80-90 ms
contended row (hot counter) aborts and retries dominate
The formula worth saying out loud: commit latency ≈ RTT(client → leaseholder) + RTT(leaseholder → nearest quorum) + disk sync (+ commit wait in Spanner). Every multi-region design decision moves one of those terms.
Data Modeling — "The Entire Game"#
In distributed SQL the schema is the physical layout. The primary key decides which range a row lives on, which leaseholder orders its writes, and which other rows share its transaction path. Most performance problems on Spanner and CockroachDB are key-design problems, and they are expensive to fix after the table holds a few billion rows. See Schema Design for the general method and Partitioning for the partitioning theory behind it.
Step 1: Spread the Inserts#
-- Anti-pattern: every insert goes to the last range
CREATE TABLE orders (
order_id INT8 DEFAULT unique_rowid() PRIMARY KEY, -- roughly time-ordered
...
);
-- Option A: random UUID primary key (spreads inserts across all ranges)
CREATE TABLE orders (
order_id UUID DEFAULT gen_random_uuid() PRIMARY KEY,
customer_id UUID NOT NULL,
created_at TIMESTAMPTZ NOT NULL DEFAULT now(),
total_cents INT8 NOT NULL
);
-- Option B: keep a time-ordered column, but hash-shard the index (CockroachDB)
CREATE INDEX orders_by_time ON orders (created_at) USING HASH;
Spanner's documentation recommends four fixes for monotonically increasing keys: hash the key and lead with the hash, swap column order, use version 4 UUIDs (random high-order bits), or bit-reverse sequential values (Spanner schema and data model). CockroachDB's hash-sharded indexes prefix the key with a computed shard column; the documentation's examples use 16 buckets, and it notes the cost: range scans over sequential values must now read every bucket (CockroachDB hash-sharded indexes).
The rule: the leading column of every primary key and every high-write secondary index must not be monotonic. That includes timestamps, sequences, and ID schemes that embed time in the high bits (see ID Generation). A Snowflake-style ID that sorts well in Postgres is exactly the key that creates a hot tail here.
Step 2: Colocate What Commits Together#
A transaction that touches one range commits in one round; across ranges it pays the atomic commit protocol. So put rows that are written together under a common key prefix.
-- Spanner: physical interleaving of child rows inside the parent's key range
CREATE TABLE Customers (
CustomerId STRING(36) NOT NULL,
Name STRING(MAX)
) PRIMARY KEY (CustomerId);
CREATE TABLE Orders (
CustomerId STRING(36) NOT NULL,
OrderId STRING(36) NOT NULL,
Total INT64
) PRIMARY KEY (CustomerId, OrderId),
INTERLEAVE IN PARENT Customers ON DELETE CASCADE;
-- CockroachDB: the same effect by leading the child's key with the parent's key
CREATE TABLE orders (
customer_id UUID NOT NULL,
order_id UUID NOT NULL DEFAULT gen_random_uuid(),
total_cents INT8 NOT NULL,
PRIMARY KEY (customer_id, order_id)
);
Spanner's documentation notes that interleaving keeps a parent and its children together across splits as long as the parent row plus children stay under the split size limit. The design tension: a big parent (a merchant with 50 million orders) is itself a hot, oversized key range. Colocate by the entity whose transactions are small and frequent (customer, cart, account), not by the one that is merely the biggest.
Step 3: Secondary Indexes Are Distributed Too#
A secondary index is another sorted table with its own ranges. Every write to the base row is also a write to each index's range, which is usually on a different node. One table with 4 secondary indexes turns a single-range insert into a 5-range transaction.
| Index choice | Write cost | Read benefit | Staff default |
|---|---|---|---|
| No secondary index | 1 range | Only key lookups | For append-mostly tables read by key |
| Global secondary index | +1 range per index per write | Lookups by any column | 2–3 indexes max on hot-write tables |
STORING / covering columns | Larger index rows | Avoids a second lookup to the base table | For the 1–2 queries that dominate traffic |
| Index prefixed with the parent key | Usually same range as parent | Locality for per-customer queries | Default for child entities |
| Hash-sharded index on a time column | 1 range per write, spread | Scans touch every bucket | Only when you must query by time |
Step 4: Decide Locality Per Table#
Locality is a schema property. CockroachDB exposes three table localities in a multi-region database (CockroachDB table localities):
ALTER DATABASE shop SET PRIMARY REGION "us-east1";
ALTER DATABASE shop ADD REGION "europe-west1";
ALTER DATABASE shop ADD REGION "asia-southeast1";
-- User-owned rows live in the user's region (hidden crdb_region column,
-- defaulting to the gateway region that inserted the row)
ALTER TABLE accounts SET LOCALITY REGIONAL BY ROW;
-- Read-mostly reference data: fast consistent reads everywhere, slow writes
ALTER TABLE currencies SET LOCALITY GLOBAL;
-- Data accessed mostly from one place stays in one region
ALTER TABLE warehouse_jobs SET LOCALITY REGIONAL BY TABLE IN "us-east1";
Spanner makes the equivalent choice mostly at the instance level: regional, dual-region or multi-region configurations, with a default leader region where writes are fastest (Spanner instance configurations). For per-row placement inside one database, Spanner's geo-partitioning lets you create instance partitions with their own regional configuration and place rows in them by a placement key (Spanner geo-partitioning). The simpler and more common pattern is still to choose one instance configuration for the whole database and keep latency-sensitive writers close to the default leader region.
🎯 Staff Move: "I'll classify every table before I pick a region count. Accounts and orders are owned by one user, so they go
REGIONAL BY ROWand write locally. Currencies and feature config are read everywhere and written a few times a day, so they'reGLOBAL. The inventory ledger is written from one warehouse system, so it's regional by table. The only cross-region transaction left is a currency update, and I'm fine with that being slow."
The Tunable Tradeoff — Survival vs Write Latency vs Read Freshness#
Distributed SQL has three dials, and each one moves a term in the commit-latency formula.
| Dial | Options | What You Gain | What You Pay | Who Pays |
|---|---|---|---|---|
| Survival goal | Zone (default) / Region | Region: no acknowledged write lost and still writable if a region dies | Writes +RTT to the nearest other region; 5 replicas instead of 3 | Every writer, every request; finance for the replicas |
| Locality | Regional by table / by row / global | Local reads and writes for home-region traffic | Remote-homed rows cost cross-region RTT; global writes are slow | Users far from the home region; writers of global tables |
| Read freshness | Strong / bounded-staleness / exact-staleness (follower) reads | Stale reads served by the nearest replica, no contention with writers | Data seconds old | The feature that reads it; product must sign off |
| Isolation | Serializable (default) / Read Committed (CockroachDB) | Read Committed: no client-side 40001 retries | Application must handle write skew itself | Every engineer writing an invariant |
The documented numbers that anchor the dials:
- Region survival in CockroachDB requires at least 3 database regions, raises the replica count to 5, and the documentation states write latency increases by at least the round-trip time to the nearest region (CockroachDB survival goals).
- Global tables use a non-blocking transaction protocol: writes are given a timestamp in the future and commit-wait until it arrives, so reads in every region are consistent and local while writes are slower (CockroachDB global tables). The lead time scales with the maximum clock offset, which is why CockroachDB's multi-region guidance suggests lowering
--max-offsetto 250 ms where clocks are well synchronised. - Follower reads with exact staleness should be at least 4.2 seconds in the past to avoid blocking on conflicting writes (CockroachDB follower reads). Spanner's stale reads (exact or bounded staleness) play the same role.
Write latency by configuration (design numbers; US-East/US-Central/US-West, EU ~85 ms away):
config local writer remote writer
regional, zone survival 5-10 ms +RTT to region (30-90 ms)
multi-region, REGIONAL BY ROW, zone survival 5-10 ms (home) +RTT to home region
multi-region, REGIONAL BY ROW, region surv. 25-40 ms +RTT to home region
GLOBAL table write hundreds of ms (future timestamp + commit wait)
Spanner multi-region (writer near leader) ~RTT to nearest other read-write region + commit wait
🎯 Staff Insight: "Region survival" is not a database setting; it is a promise to the business that costs a cross-region round trip on every write, forever. Before you turn it on, find the person who will sign off on adding 20–40 ms to checkout and roughly 1.7× the replica bill. If nobody will, zone survival plus asynchronous backups to another region is the honest answer.
Anti-Patterns — What Kills Distributed SQL Deployments#
1. Sequential Primary Keys#
Auto-increment, timestamp-leading or time-ordered IDs funnel every insert to the last range. Symptom: one node at 100% CPU, the rest idle; write throughput flat no matter how many nodes you add. Fix: UUIDv4, hash prefix, bit-reversed sequence or hash-sharded index. See Hot Keys.
2. Treating It as Postgres With Infinite Scale#
Wire compatibility is not behaviour compatibility. A query plan that did a fast index scan on one node can fan out to 200 ranges on 30 nodes. Long-running transactions that held a few row locks on Postgres now hold intents that block other transactions across the cluster. Port the workload with its query plans reviewed, not just its schema. See PostgreSQL.
3. No Retry Loop#
Under serializable isolation, contention produces aborts by design. ORMs that treat 40001 as a fatal error turn a normal contention spike into user-facing 500s. Every transaction needs a bounded retry with jittered backoff, and the transaction body must be safe to re-run (no side effects like emails or payment calls inside it).
4. Hot Rows Disguised as Business Logic#
A global counter, a single "inventory remaining" row for a flash sale, or one account that every transfer debits. Consensus doesn't make one row faster; it makes it slower, because each write is a quorum round and conflicting writers abort. Shard the counter, batch the updates, or move the hot path out of the transactional store. See Concurrency Control and Ticket Drops.
5. Multi-Region Everything#
Making every table span regions with region survival "for resilience" puts a cross-region round trip on writes that never needed it: sessions, audit logs, idempotency keys. Classify data first; most of it is regional or disposable.
6. Analytics on the OLTP Cluster#
Full-table scans for dashboards compete with checkout for the same leaseholders, and on a consensus-replicated store every byte scanned is the most expensive byte you own. Use follower reads for light reporting and CDC into a column store for anything heavier. See OLAP Databases.
7. Giant Transactions and Bulk Deletes#
Deleting 50 million rows in one transaction writes 50 million intents, holds them until commit, and can stall every reader of those ranges. Batch deletes into chunks of 1,000–10,000 rows, or use row-level TTL features where available.
8. Schema Changes Without a Rollout Plan#
Both systems run online schema changes, but backfilling an index on a 5 TB table is hours of extra write and compaction load. Treat a schema change as a deploy: expand, backfill off-peak, verify, then contract. See Schema Design.
The Technology Landscape — Head-to-Head Comparison#
| Dimension | Spanner | CockroachDB | YugabyteDB | Hand-sharded Postgres/MySQL (or Vitess) | Aurora / single-primary Postgres |
|---|---|---|---|---|---|
| Delivery | Google Cloud managed only | Self-hosted or vendor cloud | Self-hosted or vendor cloud | You build or adopt the router | Managed single writer |
| Replication unit | Split, Paxos | Range, Raft | Tablet, Raft | Whole shard, async or semi-sync replicas | Whole database |
| Clock | TrueTime (hardware-bounded) | HLC, 500 ms default max offset | HLC | N/A | N/A |
| Isolation | External consistency | Serializable (default), Read Committed | Snapshot, Serializable, Read Committed | Per shard; cross-shard is your problem | Per database |
| Cross-shard transactions | Built in (2PC over Paxos) | Built in (parallel commits) | Built in | Sagas or app-level 2PC | Not applicable |
| Resharding | Automatic | Automatic | Automatic | A quarter-long project each time | Not applicable (scale up) |
| SQL dialect | GoogleSQL or PostgreSQL interface | PostgreSQL wire-compatible | PostgreSQL-compatible (reuses PG query layer) | Native | Native |
| Write latency in-region | 5–15 ms (design estimate) | 5–15 ms | 5–15 ms | 1–5 ms | 1–5 ms |
| Lock-in risk | Highest: proprietary, one cloud | Medium: vendor licence | Medium | Low | Low–medium |
| Who pays | Finance (premium per node) | Platform team or vendor contract | Platform team | Product teams, through every cross-shard feature | Nobody, until the single writer runs out |
Head-to-head guidance:
- Spanner vs CockroachDB: Spanner if you are on Google Cloud, want the operational burden fully outsourced and value the strongest consistency model. CockroachDB if you need multiple clouds or on-prem, per-row locality controls, or PostgreSQL compatibility for existing code.
- Distributed SQL vs DynamoDB: DynamoDB is cheaper per operation with single-digit-millisecond key access and no query planner to surprise you, but transactions are limited and secondary indexes are eventually consistent. If the access patterns are a dozen known key lookups, DynamoDB wins. If they are relational and evolving, distributed SQL. See DynamoDB.
- Distributed SQL vs Cassandra: Cassandra for write-heavy, multi-region, last-writer-wins data at the lowest cost per write. Distributed SQL when invariants span rows. See Cassandra.
Patterns#
Pattern 1: Regional Transactional Core#
A single-region cluster, 3 replicas across zones, serving the business's transactional core (orders, payments, inventory). UUID keys, parent-prefixed child keys, a retry wrapper in the data-access library. When: you've outgrown one primary, writers are in one region. This is the default for most designs that reach for distributed SQL.
Pattern 2: Home-Region Rows#
REGIONAL BY ROW tables where each user's data lives in their region, giving local writes and data residency by placement. When: users are spread across continents and their data is mostly their own. The trap: unique constraints that don't include the region column must be checked in every region, and joins between rows homed in different regions are cross-region reads. Include the region in unique constraints where the business allows it.
Pattern 3: Global Reference Data#
GLOBAL tables (CockroachDB) or a read-only replica per region with stale reads (Spanner) for config, catalogs, currency rates and permissions. When: read in every request, written a few times an hour. Not for anything written per user action.
Pattern 4: Stale-Read Offload#
Route reporting, search indexing and "my order history" pages to follower or bounded-staleness reads 5–15 s behind. They use any replica, don't contend with writers and keep leaseholders free for the transactional path. When: product can sign off on seconds of staleness for that screen.
Pattern 5: CDC Out of the Core#
Changefeeds (CockroachDB) or change streams (Spanner) publish committed changes to Kafka for search, caches, analytics and the outbox pattern. When: always, once more than one system needs the data. Treat the change stream as an API with a schema and an owner.
Pattern 6: Directory on Distributed SQL#
Even teams that hand-shard their bulk data often keep the shard directory, tenant metadata and global uniqueness (usernames, email addresses) in a small distributed SQL cluster, because it's the piece that needs strong consistency and cross-region availability but carries little volume. See Sharded Database.
Scaling#
The Numbers#
| Quantity | Value | Source / Basis |
|---|---|---|
| CockroachDB range size | 512 MiB max, 128 MiB min | Current docs default |
| CockroachDB replicas per range | 3 (5 for system ranges; 5 with region survival) | Current docs default |
| CockroachDB MVCC garbage collection TTL | 4 hours (gc.ttlseconds 14400) | Current docs default |
| CockroachDB max clock offset | 500 ms default; 250 ms suggested for multi-region with good sync | Docs |
| CockroachDB exact-staleness follower reads | ≥ 4.2 s in the past | Docs |
| Spanner compute unit | 1 node = 1,000 processing units; minimum 100 PU | Docs |
| Spanner storage per node | 10 TiB per node (1 TiB per 100 PU) | Docs |
| Spanner availability | 99.99% regional, 99.999% dual- and multi-region | Docs |
| TrueTime ε | About 1–7 ms sawtooth, ~4 ms typical | OSDI 2012 paper |
| In-region consensus write | 5–10 ms p50, 15–30 ms p99 | Design estimate |
| Writes per range (leaseholder ceiling) | ~1–5K/s depending on row size and contention | Design estimate |
| Cross-region RTT | US East↔West 60–70 ms, US↔EU 80–90 ms, US↔APAC 150–200 ms | Typical public-cloud figures |
The Spanner compute and storage figures come from its compute capacity documentation (Spanner compute capacity); the availability figures from its instance configuration documentation.
Capacity Math#
Workload: 8K writes/s peak, 60K reads/s, 6 TB data, one region, zone survival.
CockroachDB (self-hosted), sizing rule of thumb per node: ~1-2K writes/s and ~1.5 TB of
replicated data with headroom (design estimate; load-test your own schema):
replicated data = 6 TB x 3 replicas = 18 TB -> 12 nodes at 1.5 TB each
write capacity = 8K/s / 1.5K per node -> ~6 nodes, but every write lands on 3
reads = 60K/s, 70% served as stale follower reads -> spread across all nodes
=> 12-15 nodes across 3 zones (multiples of 3 so each zone holds one replica)
Spanner (regional): storage needs 6 TB / 10 TiB per node -> 1 node minimum for storage;
CPU drives it: plan 6-10 nodes for 8K writes/s with secondary indexes, then
run the managed autoscaler against a 65% high-priority CPU target (design estimate).
What Breaks First as You Grow#
- A hot range long before the cluster runs out of capacity. One sequential index, one popular tenant, one counter.
- Contention on a few rows: retries grow superlinearly with concurrency on the same key.
- Secondary index write amplification: at 5+ indexes the write path is mostly index maintenance.
- Schema change and backfill time on multi-terabyte tables: hours of extra load.
- Cross-region transactions that crept into the hot path when a second region was added.
- Cost: consensus-replicated storage at 3–5 copies, on premium SSD, for data that turned out to be logs.
Scaling Moves in Order#
- Fix keys and indexes (hash-shard or UUID; drop unused indexes).
- Move tolerant reads to stale or follower reads.
- Shard hot counters and batch hot-row updates.
- Add nodes; the rebalancer moves ranges in minutes to hours depending on data size.
- Move analytics and archives out through CDC; add row-level TTL.
- Only then: multi-region, table by table.
Failure Modes & Recovery#
1. Hot Range From a Sequential Key#
- Symptom: Write p99 rises from 15 ms to 400 ms; one node at 95% CPU, others under 30%; adding nodes doesn't help.
- Root cause: Primary key or secondary index leads with a timestamp or sequence; every insert targets the rightmost range.
- Detection: Per-range QPS (CockroachDB hot ranges page; Spanner Key Visualizer); leaseholder QPS skew across stores; CPU spread > 3× between nodes.
- Fix: Add a hash-sharded index or new UUID-keyed table and migrate writers; short term, split and scatter the range manually.
- Prevention: Schema review lint: no monotonic leading column on tables over 100 writes/s. See Hot Keys.
2. Retry Storm Under Contention#
- Symptom: A promotion starts;
40001errors climb from 0.1% to 30% of transactions; throughput falls as concurrency rises. - Root cause: Thousands of transactions write the same rows (inventory count, a shared balance); each abort retries immediately and contends again.
- Detection: Transaction restart rate by statement fingerprint;
txn.restartsmetrics; contention event tables; p99 latency of the hot statement. - Fix: Exponential backoff with jitter and a retry cap of 5–10; queue or batch writes to the hot row.
- Prevention: Load-test the hottest row at 3× peak concurrency; model hot counters as N sharded rows. See Concurrency Control and Backpressure.
3. Region Loss With Zone Survival#
- Symptom: A cloud region goes dark. Users homed there see errors on every write and strong read; other regions stay healthy except for transactions that touch those rows.
- Root cause: With zone survival, all voting replicas of a region's home rows are in that region. This is the configured trade, not a bug.
- Detection: Unavailable range count by region; node liveness by locality; synthetic writes per region.
- Fix: Wait for the region, or fail over from backups to another region and accept the RPO of the last backup. Stale reads may still be served from non-voting replicas elsewhere if configured.
- Prevention: Decide per data class, in writing: region survival for the ledger, zone survival plus frequent cross-region backups for the rest. Rehearse the restore.
4. Clock Skew Takes Nodes Out#
- Symptom: Several CockroachDB nodes crash in a short window; logs show clock offset errors; ranges briefly unavailable while leases move.
- Root cause: NTP misconfiguration, a VM pause, or a time-source change pushed offset past
--max-offset. Nodes self-terminate rather than serve inconsistent reads. - Detection:
clock-offset.meannanosper node against the configured limit; alert at 40% of the max offset. - Fix: Correct the time source; restart nodes once clocks agree.
- Prevention: Cloud provider's time-sync service with leap-second smearing configured identically on every node; alert well before the limit; in Spanner this risk is the provider's.
5. Analytics Query Starves the Transactional Path#
- Symptom: Checkout p99 doubles every morning at 9:00; CPU saturates on leaseholders; no change in traffic.
- Root cause: A dashboard runs full scans with strong reads against the OLTP tables; scans consume leaseholder CPU and I/O.
- Detection: Top statements by CPU and rows read; statement fingerprints from BI users; per-node CPU correlated with a schedule.
- Fix: Move the query to follower reads or kill it; give BI a separate user with admission-control priority set low.
- Prevention: BI goes through CDC to a column store; OLTP access by allow-listed services only. See OLAP Databases.
Operational Reality Matrix#
| Failure | Detection Signal | Blast Radius | Mitigation | Owner |
|---|---|---|---|---|
| Hot range | Leaseholder QPS and CPU skew | Every writer to that table | Hash-shard or rekey | Product team owns the schema; platform detects |
| Retry storm | Restart rate by statement | The contended feature, then shared nodes | Backoff, batch, shard row | Product team |
| Region loss (zone survival) | Unavailable ranges by region | Users homed in that region | Wait or restore elsewhere | Platform + business owner of RPO |
| Clock skew | Clock offset vs max | Nodes on bad clocks | Fix time source | Platform / infrastructure |
| Analytics overload | Top statements by CPU | Whole cluster | Kill, move to followers | Data team; platform enforces access |
| Leaseholder node loss | Lease transfers, p99 spike | Ranges on that node for seconds | Client retries with timeouts | Platform (automatic) |
| Long backfill | Schema job progress, IO | Cluster write latency | Pause, run off-peak | Product team + platform change review |
When to Use vs. Alternatives#
| Need | Pick | Why |
|---|---|---|
| Relational OLTP that fits one primary for 2+ years | Postgres (managed) | 1–5 ms writes, mature tools, lowest cost |
| Relational OLTP beyond one primary, cross-entity transactions | Distributed SQL, regional | Automatic sharding and serializable transactions |
| Users on several continents, data mostly per user | Distributed SQL, multi-region with home-region rows | Local writes plus residency by placement |
| Ledger that must survive a region with RPO 0 | Spanner multi-region or CockroachDB region survival | Synchronous cross-region quorum |
| Known key-value access, massive scale, minimal cost | DynamoDB / Cassandra | Cheaper per op, predictable latency |
| Append-heavy logs, events, metrics | Kafka, Cassandra, time-series DB | No need to pay consensus on every row |
| Analytics and dashboards | Column store fed by CDC | Scans are cheap there, expensive here |
| Existing sharded MySQL with a working router | Keep it (or Vitess) | Migration cost exceeds benefit unless cross-shard transactions are the pain |
When NOT to Use Distributed SQL#
- The workload fits one Postgres primary. A modern managed Postgres instance handles tens of thousands of writes per second and terabytes of data. Distributed SQL triples storage replicas, adds 5–10 ms to every write, and costs 2–5× more for the same workload. If your 3-year growth curve fits a single primary plus read replicas, buy the cheaper thing. See PostgreSQL.
- Single-region, latency-critical writes. An order-matching engine or a hot leaderboard needs sub-millisecond writes. Consensus on every write is the wrong tool. See Order Matching.
- The data has no cross-row invariants. Event logs, clickstreams, sensor data and chat messages need throughput and cheap storage, not serializable transactions.
- The access pattern is pure key-value at very high scale. DynamoDB or Cassandra deliver predictable single-digit-millisecond access at a fraction of the cost per operation.
- Every transaction is a hot row. Distributed SQL can't make one row faster. If the design is a global counter, fix the design.
- No one will own the cluster. Self-hosted CockroachDB is a fleet with upgrades, backups, clock hygiene and capacity planning. If there's no platform owner, use a managed service or a simpler database.
🎯 Staff Move: "Before I pick distributed SQL, I'll check whether we've actually outgrown Postgres. At 3K writes per second and 2 TB, one primary handles it for years at a third of the cost. The trigger for me is cross-shard transactions becoming the main source of bugs, or writers on two continents. Until then, I'd keep the schema distribution-friendly, UUID keys and parent-prefixed children, so the migration is cheap when it comes."
Operational Concerns#
What the On-Call Actually Does#
- Watches the skew, not the average. Per-node CPU and leaseholder QPS spread; a 3× gap is a hot range even when cluster CPU is 40%.
- Triages slow statements by fingerprint: full scans, missing indexes, contention events. Most "the database is slow" pages are one query.
- Watches retry rates per statement and owns the conversation with the team whose transaction is contending.
- Runs rolling upgrades (self-hosted): one node at a time, drain leases first, wait for under-replicated ranges to reach zero before the next. On a 30-node cluster this takes hours and is routine.
- Verifies backups by restoring them into a scratch cluster on a schedule; a backup that has never been restored is a hope.
- Guards clock hygiene (self-hosted CockroachDB) and the admission-control priorities that keep batch jobs behind user traffic.
Key Metrics & Alerts#
| Metric | Healthy | Alert |
|---|---|---|
| Statement p99 latency (OLTP fingerprints) | < 30 ms | > 100 ms for 10 min |
Transaction restart (40001) rate | < 0.5% | > 5% for 5 min |
| Node CPU spread (max / median) | < 1.5× | > 3× for 15 min (hot range) |
| Under-replicated ranges | 0 | > 0 for 10 min |
| Unavailable ranges | 0 | > 0 (page) |
Clock offset vs --max-offset | < 10% | > 40% (page) |
| Disk usage per node | < 60% | > 75% (rebalancing needs room) |
| Spanner high-priority CPU | < 65% regional, < 45% multi-region (design targets) | Above target for 30 min |
| Last successful restore test | < 30 days | > 45 days |
Upgrades and Schema Changes as Ongoing Cost#
Self-hosted CockroachDB ships major versions on a regular cadence, and each has a support window; falling behind ends in a forced multi-step upgrade. Spanner removes this entirely, which is part of what its premium buys. Schema changes, on both, need the same discipline as deploys: expand, backfill, verify, contract, with backfills scheduled off-peak and watched for write-latency impact.
Interview Application — Staff-Level Plays#
Which Case Studies Use Distributed SQL#
| Case Study | How Distributed SQL Is Used | Key Pattern |
|---|---|---|
| Ride Hailing & Delivery | Trip, order and courier state with cross-entity transactions | Uber's fulfillment platform on Spanner: consistency instead of app-level coordination |
| Multi-Region Active-Active | Home-region rows and synchronous global quorum as two of the three multi-region strategies | Classify data; region survival only for RPO-0 data |
| Sharded Database | The "buy" option against hand-sharding | Price latency and cost against the router and saga work |
| Database Selection | The relational choice once one primary is not enough | Distributed SQL vs Postgres vs DynamoDB decision |
| Payments | Ledger and payment state with serializable transfers | Single-range transactions per account; idempotency keys in the same transaction |
| Reservation Systems | Seat and room inventory with no double-booking | Serializable check-and-insert; hot inventory rows need batching |
| Replicated Key-Value Store | External consistency as the strong end of the spectrum | Commit wait and its latency cost |
| Auth & Identity | Globally unique usernames and emails across regions | Global uniqueness index; tokens validated locally, not per-request DB reads |
| ID Generation | Why time-ordered IDs hurt range-partitioned stores | UUIDv4 or bit-reversed sequences for primary keys |
| Ticket Drops | Order capture under extreme contention | Admission queue in front; never one inventory row for everything |
Every System Design Question Has a Distributed SQL Moment#
- Payments: "Each account's balance and its ledger entries share a key prefix, so a single-account debit is a one-range transaction at 5–10 ms. Transfers between accounts are cross-range but still serializable, and I keep the idempotency key in the same transaction so retries after
40001can't double-charge." - Global social app: "Profiles and posts are
REGIONAL BY ROW, homed in the author's region. The follow graph is read everywhere but written per user, so it's also per-row, and the feed is precomputed into a cache. The only global table is configuration." - Ticketing: "I won't put 50,000 buyers on one inventory row. Seats are rows; a buyer's transaction claims one seat row with a conditional update, so contention spreads across 50,000 keys instead of one."
- URL shortener: "I'd not use distributed SQL here. The access pattern is a key lookup at 100K reads per second, so a key-value store plus a cache is cheaper and faster. See URL Shortener."
What Interviewers Probe#
| After You Say... | They Will Ask... | What They're Evaluating |
|---|---|---|
| "Spanner gives us strong consistency" | "What does a write cost? From Europe?" | Commit latency formula and leader placement |
| "Auto-increment order IDs" | "Where does the insert land?" | Hot range awareness |
| "Multi-region for availability" | "What happens to writes when a region dies? What does it cost?" | Survival goals and their latency bill |
| "Serializable isolation" | "What does the client do on a conflict?" | Retry handling and idempotent transaction bodies |
| "Distributed SQL instead of sharding" | "Why not just Postgres? What does this cost?" | Knowing when not to use it |
| "TrueTime" | "How does CockroachDB do it without atomic clocks?" | HLC, max offset, read uncertainty restarts |
Common Interview Mistakes#
| What Candidates Say | What Interviewers Hear | What Staff Engineers Say |
|---|---|---|
| "Distributed SQL scales infinitely" | Doesn't know about per-range limits | "It scales across ranges. One hot range is still one leaseholder." |
| "Spread it across three regions" | Hasn't priced write latency | "Zone survival, regional by row; region survival only for the ledger, at +30 ms per write." |
| "It's Postgres-compatible, so we just migrate" | Will discover distributed query plans in production | "Wire-compatible; I'd replay production queries and review every plan that fans out." |
| "Strong consistency means no anomalies to handle" | Will ship code that fails on 40001 | "Serializable means aborts; the data layer retries with jitter and the body is idempotent." |
| "TrueTime makes it fast" | Confuses correctness with latency | "TrueTime buys ordering; commit wait costs a few milliseconds per write." |
| "We'll run dashboards on it" | Analytics will starve checkout | "Follower reads for light reporting, CDC to a column store for the rest." |
L5 vs L6 vs L7 Responses#
| Scenario | L5 Answer | L6 / Staff Answer | L7 / Principal Answer |
|---|---|---|---|
| "Our Postgres primary is at 80% CPU" | Migrate to CockroachDB | First check whether read replicas, query fixes or a bigger instance buy 2 years; distributed SQL when cross-entity writes outgrow one primary | Sets the company trigger for distributed SQL adoption and funds one platform team rather than five teams evaluating vendors |
| "Expand to the EU with data residency" | Add an EU region to the cluster | REGIONAL BY ROW for user tables, EU rows homed in EU, global tables for config, uniqueness checks priced | Residency policy per data class signed by legal; decides whether the EU gets its own cluster to contain regulatory blast radius |
| "Checkout p99 jumped to 400 ms" | Add nodes | Finds the hot range or contended row; rekeys or shards the counter; adds a schema lint | Makes key review part of the schema-change process for every team on the shared cluster |
| "Should we pick Spanner or CockroachDB?" | Whichever is faster | Spanner on Google Cloud for managed ops and strongest consistency; CockroachDB for multi-cloud or per-row locality | Prices 5-year TCO including exit cost and licence risk; keeps the access layer portable |
The Staff Distributed SQL Checklist#
- Justify it: "We need cross-entity transactions beyond one primary's write capacity, or writers on several continents. Otherwise Postgres."
- Keys: "UUIDv4 or hash-prefixed primary keys; no index leads with a timestamp; children prefixed with the parent key."
- Commit math: "Client to leaseholder, plus leaseholder to quorum, plus disk. In-region that's 5–10 ms."
- Locality per table: "Per-row home regions for user data, global for read-mostly config, regional for the rest."
- Survival goal, priced: "Zone survival by default; region survival only for the ledger, at a cross-region round trip per write and 5 replicas."
- Contention and reads: "Retries with jitter on
40001, sharded hot rows, follower reads for anything that tolerates 5 seconds, CDC for analytics."
🎯 Staff Insight: Don't use distributed SQL for event logs, analytics, blobs, queues, sessions, or any workload that fits a single Postgres primary for the next three years. The strongest distributed SQL signal in an interview is explaining why you aren't using it for half the data in the design.
Evaluation Rubric#
| Dimension | Senior (L5) | Staff (L6) | Principal (L7) |
|---|---|---|---|
| Justification | "It scales and it's consistent" | Names the trigger: write volume plus cross-entity invariants, or multi-region writers | Adoption trigger and paved road for the org; when to say no |
| Data model | Postgres-style schema | Keys for distribution, colocation for single-range commits, index budget | Key-design standard enforced in schema review |
| Latency | "Single-digit milliseconds" | Commit formula with RTT terms per configuration | Prices latency against revenue and survival requirements |
| Failure | "Consensus handles it" | Leaseholder failover in seconds; region loss under each survival goal | Org-wide survival policy per data class; restore drills |
| Cost | Not discussed | Replica count, node sizing, analytics offload | 5-year TCO, managed vs self-hosted, exit cost |
Strong hire signals
| Signal | What It Sounds Like |
|---|---|
| Knows who pays for clocks | "Spanner charges writers commit wait; Cockroach charges readers occasional restarts." |
| Key discipline | "No index leads with a timestamp. I'd hash-shard that one." |
| Prices multi-region | "Region survival is +30 ms on every write. Which tables need it?" |
| Expects aborts | "Serializable means retries. The transaction body must be idempotent." |
| Knows when not to | "At 3K writes per second this is a Postgres problem." |
Lean no-hire signals
| Signal | Why It Misses the Bar |
|---|---|
| Distributed SQL as the answer to every scaling question | Hasn't weighed latency or cost |
| Sequential keys in a range-partitioned store | Will build a single-node database with extra steps |
| "Multi-region" with no survival goal or latency number | Unpriced promise |
| No retry handling | Will fail under the first contention spike |
Common false positives
- Explaining TrueTime in detail ≠ designing with Spanner. Ask where the leader region is and what an EU write costs.
- "We migrated to CockroachDB" ≠ judgment. Ask what they'd keep in Postgres.
- Knowing Raft ≠ knowing hot ranges. Consensus correctness says nothing about where the inserts land.
Beyond Staff: The Principal View#
Why L7 Sees This Problem Differently#
At Staff level, distributed SQL is a database you configure well: keys, locality, survival, retries. At Principal level it is a bet on where the company's transactional data will live for the next decade, with a vendor, a cost curve and a set of one-way doors. The question changes from "Spanner or CockroachDB for this service?" to "Do we want one shared, strongly consistent OLTP platform that most teams default to, or a few dedicated clusters for the systems that genuinely need it, with Postgres as the default everywhere else? And what does it cost to leave?"
🧭 Principal Move: "Before choosing a vendor, I want three things written down: which data classes need serializable transactions across entities, which need to survive a region with zero data loss, and which need to live in a specific country. That list sets the cluster count, the survival goals and most of the bill. Anything not on it stays on Postgres, and we stop paying consensus prices for logs."
The Org-Level Fault Line#
One shared distributed SQL platform vs per-domain clusters vs Postgres-by-default.
| Option | What Works | What Breaks | Who Pays |
|---|---|---|---|
| One shared cluster for the company | One platform team, one upgrade path, cross-domain transactions possible | One team's hot range or analytics query hurts everyone; shared blast radius; cross-team schema coupling | Every product team during the incident; platform on-call always |
| Dedicated cluster per critical domain (payments, orders, identity) | Contained blast radius; survival goal and region set per domain | More clusters to upgrade; cross-domain transactions become sagas again | Platform headcount; domain teams for integration |
| Postgres by default, distributed SQL by exception | Lowest cost; simplest ops for most teams | Teams hit the single-primary ceiling and do emergency migrations | The team that outgrows Postgres, at the worst time |
| Hand-sharding everywhere | No vendor dependency | Every team rebuilds routers, resharding and sagas | Product velocity, for years |
The Principal default: Postgres is the paved road for new services. Distributed SQL is a platform product with an entry review: dedicated clusters per critical domain, owned by a database platform team, offered when a team shows the trigger (cross-entity writes beyond one primary, multi-region writers, or RPO 0 across regions). No shared cluster across unrelated domains, because a cross-domain transaction is usually a coupling to remove, not a feature to enable.
🧭 Principal Insight: "The expensive mistake isn't choosing the wrong vendor. It's letting distributed SQL become the default for data that never needed it, so the company pays 3–5 replicas of premium storage and a consensus round per write for audit logs and session state."
Cost Model#
Assumptions (directional only; check current pricing): self-hosted nodes 16 vCPU / 64 GiB with fast SSD ~$900/month each; managed regional node-equivalents ~$700–1,000/month and multi-region ~2.5–3.5× that; loaded engineer $250K/year ($21K/month). Storage is counted at 3 replicas (zone survival) or 5 (region survival).
| Scale | Setup | Infra/month | People | Total/month | Versus Postgres |
|---|---|---|---|---|---|
| Startup (1 TB, 1K writes/s) | Managed regional, minimum footprint | ~$3–5K | 0.25 FTE (~$5K) | ~$10K | Managed Postgres ~$2–3K: not yet worth it |
| Growth (8 TB, 10K writes/s, 1 region) | 15 nodes self-hosted, or managed equivalent | ~$14–20K | 1.5 FTE (~$32K) | ~$50K | Hand-sharded Postgres: ~$10K infra plus 3–4 FTE of router and saga work |
| Enterprise (100 TB, 3 regions, ledger with region survival) | 5 domain clusters, 150+ nodes | ~$200–300K | 6–8 FTE (~$150K) | ~$400K | The alternative is a bespoke multi-region data layer nobody wants to own |
Two levers dominate. First, data classification: moving logs, events and analytics out of the cluster often cuts its size by half. Second, survival goals: region survival adds roughly two-thirds more replicas than zone survival for the same data, so it should cover the ledger, not the product catalog.
The 3-Year Evolution Path#
One-Way Doors vs Two-Way Doors#
| Decision | Reversibility | Cost to Reverse |
|---|---|---|
| Vendor (Spanner's dialect and APIs vs CockroachDB's) | One-way at scale | Rewriting access code, re-validating transactions, months of dual writes |
| Primary key scheme on a multi-terabyte table | One-way-ish | New table, backfill, dual-write migration |
| Relying on serializable cross-entity transactions everywhere | One-way | Leaving means introducing sagas into every flow that relied on them |
| Shared cluster across domains | One-way-ish | Splitting a cluster is a migration per domain |
| Table locality (regional / per-row / global) | Two-way | ALTER TABLE ... SET LOCALITY plus rebalancing time |
| Survival goal | Two-way | Config change; replica movement takes hours |
| Adding or removing a region | Two-way-ish | Data movement and latency change for users there |
| Stale vs strong reads per query | Two-way | Code change |
The Standard I'd Write#
RFC-DATA-007: Using Distributed SQL
Scope: Any new or migrated datastore on the distributed SQL platform.
Entry criteria (one MUST apply):
a. Cross-entity transactional writes beyond one Postgres primary after tuning.
b. Writers in two or more regions with local-latency requirements.
c. Data class requires RPO 0 across a full region loss.
MUST
1. Primary keys and high-write index leading columns are non-monotonic
(UUIDv4, hash prefix, bit-reversed, or hash-sharded index). Lint enforced.
2. Every transaction runs through the shared retry wrapper (jittered backoff,
max 10 attempts) and has an idempotent body with no external side effects.
3. Each table declares a locality; each database declares a survival goal,
with region survival approved by the data owner and finance.
4. No analytics, logs, blobs or queues. Reporting uses stale/follower reads;
analytics uses CDC to the column store.
5. Schema changes follow expand -> backfill off-peak -> verify -> contract.
SHOULD
6. Children prefixed with the parent key so common transactions touch one range.
7. No more than 3 secondary indexes on tables above 1K writes/s.
Exceptions: Database platform lead plus owning director; time-boxed; on the dashboard.
Rollout: lint in warn mode for one quarter, then blocking in CI.
Success metrics: zero hot-range incidents per quarter; restart rate < 1% for
every service; restore drill per cluster each quarter; cost per transaction tracked.
What I'd Tell the VP#
"Our order and payment data has outgrown a single database server, and splitting it ourselves would cost us two years of engineering and a long tail of bugs. I recommend a distributed SQL database for that core data only, run by one platform team. It keeps our data consistent as we grow and lets us serve European customers from Europe. It is more expensive per transaction than what we run today, so everything that doesn't need it, like logs and analytics, stays on cheaper systems. Surviving the loss of a whole cloud region with no lost payments is possible, but it slows every payment by a few hundredths of a second and adds roughly two-thirds to that cluster's cost, so I'd apply it to the ledger only."
Principal Interview Signals#
| Signal | What It Sounds Like |
|---|---|
| Classifies data before choosing | "Three data classes need it; the other twelve stay on Postgres." |
| Prices survival | "Region survival is two more replicas and a cross-region round trip per write. The ledger gets it, the catalog doesn't." |
| Treats vendor choice as a one-way door | "The exit cost is the access code and every flow that assumes serializable cross-entity writes." |
| Contains blast radius | "Cluster per critical domain. A hot range in search metadata should never page payments." |
| Sets the adoption trigger | "Entry review with three criteria, so we stop re-arguing this per team." |
Staff answers that L7 interviewers find insufficient:
- "We'll use Spanner for global consistency" without saying which data needs it and what the rest of the data uses.
- "Multi-region with region survival" without the replica cost, the latency on every write, or who signed off.
- "Use hash-sharded keys" as advice for one service, rather than a lint in the schema-change pipeline for every team.
How Real Companies Built It#
Google Spanner — TrueTime and Commit Wait#
Google's OSDI 2012 paper describes Spanner as the first system to distribute data at global scale and support externally consistent distributed transactions. Data is sharded into Paxos-replicated groups, and the TrueTime API exposes clock uncertainty, which in production was a sawtooth of about 1–7 ms; read-write transactions commit-wait out that uncertainty so commit order matches real time (Spanner paper, OSDI 2012).
Staff insight: The design spends milliseconds on every write to buy global ordering and lock-free consistent snapshot reads. In an interview, say that trade out loud: commit wait is cheap relative to a cross-region quorum, which is why Spanner's real latency cost is replica placement, not clocks.
Google F1 — AdWords on Spanner#
Google's VLDB 2013 paper describes F1, the distributed SQL database built to support the AdWords business, layered on Spanner for synchronous cross-datacenter replication and strong consistency. The paper says the higher commit latency of synchronous replication is mitigated through a hierarchical schema model and application design (F1 paper, VLDB 2013).
Staff insight: The latency problem was solved in the schema: hierarchical (interleaved) tables keep a customer's data together so common transactions stay local. That is the "colocate what commits together" rule, proven on a revenue-critical system.
CockroachDB — Geo-Distributed SQL on Commodity Clocks#
The CockroachDB SIGMOD 2020 paper presents a distributed SQL database built for global OLTP workloads with high availability and strong consistency, covering its Raft-replicated ranges, its transaction model and its data placement and replication. Its current documentation describes the per-row home regions, global tables and survival goals that turn those mechanisms into schema-level choices (CockroachDB SIGMOD 2020 paper, survival goals).
Staff insight: Without hardware clocks, CockroachDB moves the cost of uncertainty to reads and exposes locality as SQL. That makes multi-region design a reviewable schema decision, which is exactly where a Staff engineer wants it.
Uber — Fulfillment Platform on Spanner#
Uber's engineering blog describes rebuilding its fulfillment platform, the system behind trips, orders and the entities that fulfill them, onto Google Cloud Spanner in a North America multi-region configuration. The previous design traded consistency for availability and coordinated multi-entity writes with best-effort mechanisms; the new one relies on cross-table, cross-shard ACID transactions and adds an internal component for at-least-once post-commit work (Uber blog).
Staff insight: The win was deleting application-level consistency code, not raw throughput. Notice the side component: post-commit work like notifications still needed its own at-least-once machinery, because a database transaction can't include an external side effect. See Ride Hailing & Delivery.
Practice Drill#
Prompt: "You run a marketplace on one Postgres primary: 6K writes/s at peak, growing 60% a year, 4 TB. Orders, payments and seller balances must stay consistent. Next year you launch in the EU with a data residency requirement. Design the data layer for the next three years."
Staff Answer
Start with headroom: at 60% growth the primary passes ~15K writes/s in two years, beyond what one primary sustains with margin, and the core invariant (an order debits a buyer, credits a seller, reserves inventory) spans entities with no common shard key. Hand-sharding would turn every order into a saga. So the transactional core moves to distributed SQL; everything else stays put. Schema first: UUIDv4 primary keys, the orders and order_items tables prefixed with buyer_id so checkout is mostly single-range, seller balances as their own rows with a 16-way sharded counter for the top 100 sellers, and at most 3 secondary indexes on the order table. Year one is regional: 3 replicas across zones, zone survival, about 5–10 ms per commit. All transactions go through a retry wrapper with jitter and a 10-attempt cap; payment calls happen after commit through an outbox, never inside the transaction. Order history and seller dashboards use follower reads at ~5 s staleness; analytics moves to a column store through a changefeed. For the EU: add the region, make users, orders and payments REGIONAL BY ROW so EU rows are homed and stored in the EU, keep product categories and currency rates GLOBAL, and include the region in unique constraints where possible so uniqueness checks don't cross the ocean. Zone survival for most data, with backups to a second EU region; region survival only for the payment ledger if finance signs off on the extra latency and replicas. Metrics: restart rate by statement, leaseholder skew, p99 commit by region, unavailable ranges.
Why this is L6:
- Justifies distributed SQL from the growth curve and the cross-entity invariant, not fashion.
- Designs keys for distribution and colocation before talking about regions.
- Prices each multi-region choice in milliseconds and replicas, and keeps region survival narrow.
- Moves side effects, analytics and history reads off the transactional path.
What L7 adds:
- Classifies data for the whole company and sets an entry review, so other teams don't default onto the cluster.
- Prices the move: cluster cost against the 3–4 engineers hand-sharding would need, and the exit cost of the vendor's dialect.
- Treats EU residency as a policy owned with legal, and decides whether the EU deserves its own cluster for regulatory blast radius.
Quick Reference Card#
Model: SQL over a sorted, transactional KV store, split into ranges,
each range replicated by its own consensus group
Spanner: splits + Paxos; TrueTime eps ~1-7 ms (2012 paper); commit wait;
external consistency; 1 node = 1,000 PU, 10 TiB per node;
99.99% regional, 99.999% multi-region
CockroachDB: ranges + Raft; HLC; --max-offset 500 ms default (250 ms suggested
for multi-region); range 128-512 MiB; 3 replicas (5 system);
SERIALIZABLE default, Read Committed available; retries = 40001
Commit latency: RTT(client->leaseholder) + RTT(leaseholder->quorum) + disk
(+ commit wait in Spanner); in-region ~5-10 ms
Survival: ZONE = 3 voters in home region; REGION = 5 replicas, >= 3 regions,
+RTT to nearest region on every write
Localities: REGIONAL BY TABLE | REGIONAL BY ROW (crdb_region) | GLOBAL
Follower reads: exact staleness >= 4.2 s; no contention with writers
Parallel commits: cross-range commit in one consensus round (CockroachDB >= 19.2)
Hot-key fixes: UUIDv4, hash prefix, bit-reversed sequence, hash-sharded index
RED FLAGS
- Auto-increment or timestamp-leading primary keys or indexes
- No retry handling for serialization errors
- Region survival on every table with no latency number
- Dashboards and full scans on the OLTP cluster
- Side effects (payments, emails) inside a retryable transaction
- One shared cluster for unrelated domains
- Distributed SQL for data that fits one Postgres primary for 3 years