CockroachDB from First Principles: The Database That Refuses to Die
Why a global SQL database is really a problem of agreement and time, how Google solved it with TrueTime, GPS, and atomic clocks, and how CockroachDB reaches a practical design on ordinary infrastructure.

Suppose the requirement is this:
- keep accepting transactions when a machine or availability zone disappears;
- survive the loss of an entire region if the topology was designed for it;
- let applications use SQL and multi-row transactions;
- add capacity without manually assigning customers to new shards;
- keep two concurrent buyers from purchasing the same last item;
- change the schema while the system is serving traffic.
That requirement is not asking for a larger PostgreSQL server. It is asking for many machines to behave like one correct database while machines, disks, clocks, and networks keep disagreeing.
CockroachDB is one answer. It presents a PostgreSQL-compatible SQL interface, converts SQL work into operations on a distributed key-value store, divides that keyspace into small ranges, and replicates each range with Raft consensus. A range can move, split, lose a replica, or elect a new leader without turning the whole database into a manual failover project.
The impressive part is not that it stores three copies. Replication is easy until two copies both believe they are allowed to accept the next write. The hard part is establishing one history: which transaction happened first, which value a read is allowed to observe, and whether an acknowledged write can ever disappear.
To understand why CockroachDB looks the way it does, we need to start with PostgreSQL, then time, then Google Spanner. Otherwise terms such as TrueTime, hybrid logical clock, range lease, and uncertainty interval are just expensive vocabulary with no problem attached.
PostgreSQL is not the villain
PostgreSQL is usually the correct place to start. One primary server owns the writable database. Its local storage, buffer cache, write-ahead log, indexes, and lock manager are close together. A transaction can modify several rows and commit without asking machines on another continent for permission.
That architecture goes a very long way. A team can scale the machine vertically, tune queries and indexes, pool connections, partition large tables, move analytical work elsewhere, and add read replicas. Managed PostgreSQL can automate backups, patching, replica promotion, and much of the undifferentiated operations work.
Read replicas solve a specific problem: more places can serve suitable reads. They do not turn the database into several independent writers. Inserts, updates, deletes, and the authoritative transaction order still pass through the primary.
This is a feature. One writer makes coordination understandable. The difficult questions begin when the requirements exceed that boundary.
Failure changes the meaning of replication
With asynchronous replication, the primary can acknowledge a commit before a replica has received its latest WAL records. If the primary disappears in that gap and an older replica is promoted, the system may lose transactions that clients were already told had committed.
Synchronous replication can require confirmation from one or more standbys before acknowledging a commit. That improves durability, but the write now pays for another machine and its network path. If the required standby is unavailable, either writes wait or the durability policy must change. There is no setting that creates both zero coordination cost and zero data-loss risk.
Failover also needs an authority that can distinguish a dead primary from an isolated primary. If the network partitions and both sides accept writes, there are now two histories. Production PostgreSQL platforms prevent that with orchestration, quorum-backed coordination, fencing, and careful promotion rules. PostgreSQL itself is not defective because this machinery exists. The machinery is the price of changing ownership safely.
Read scaling is not write scaling
Adding replicas can increase read capacity, but one primary remains the write coordination point. Eventually a workload may exceed one machine’s CPU, memory, storage throughput, storage capacity, or acceptable regional latency.
The conventional next step is sharding:
customer_id 1..9,999 -> shard A
customer_id 10,000..19,999 -> shard B
customer_id 20,000..29,999 -> shard C
The application or a routing layer decides which shard owns a row. This distributes writes, but it also distributes responsibility into application code.
Now the team owns the shard map, rebalancing, hot-shard detection, cross-shard queries, cross-shard uniqueness, backups across several independent databases, and transactions that touch more than one shard. Moving customer 9,999 while orders are still arriving is no longer a storage copy. It is a correctness migration.
Sharding is not impossible. Some of the largest systems in the world use it successfully. The issue is who owns it. CockroachDB’s central product decision is that partitioning, replication, and transaction coordination should be database responsibilities instead of an application framework built separately by every company.
The real problem is agreement
Imagine three replicas of one account balance. The value is 100 on all three. A client submits withdraw 20 while one replica is temporarily unreachable.
Can the reachable replicas commit? With three voting replicas, two form a majority, so they can. If the isolated replica also tried to accept a competing update alone, it would have only one vote and could not commit. When it reconnects, it learns the majority’s log.
That majority rule is more important than the number of copies. Three replicas tolerate one unavailable voter. Five tolerate two. Two replicas tolerate none while preserving a majority, because one out of two cannot prove that the other side is not also accepting work.
Consensus protocols such as Raft establish an ordered log for a replicated unit of data. A leader proposes a command. Followers persist the proposal. Once a majority has accepted it, the entry is committed and surviving replicas can apply the same state transition in the same order.
Consensus answers questions such as:
- which replica is currently allowed to lead;
- which commands are committed;
- what happens when the leader disappears;
- how a returning replica catches up;
- why two network partitions cannot both commit conflicting histories.
It does not make a wide-area network fast. If a quorum crosses regions, a write crosses regions. Physics remains on the invoice.
It also does not by itself provide a SQL transaction across many independent consensus groups. That requires transaction coordination and a trustworthy ordering model. This is where time enters the story.
The enemy is not the clock; it is pretending the clock is exact
Every machine has a physical clock, and every physical clock is wrong by some amount.
Quartz oscillators drift. Virtual machines can pause. A hypervisor can expose unstable time. NTP corrects clocks over a network whose delay changes. An operator can misconfigure a time source. Leap seconds and clock steps create edge cases. Even healthy machines can disagree by milliseconds, and unhealthy ones can be much further apart.
Suppose node A says an order committed at 10:00:00.120 and node B says a payment committed at 10:00:00.080. The numeric timestamps suggest that the payment happened first. But if B’s clock was 100 milliseconds ahead, the conclusion is wrong.
For logs, a little disorder is irritating. For a serializable database, ordering decides what a transaction was permitted to see. If transaction B began after transaction A completed in the real world, users reasonably expect B to observe A. A timestamping system that violates that relationship can produce a history that looks valid locally but contradicts the order clients experienced.
Distributed databases therefore need to separate two ideas:
- physical time: approximately where an event happened on a wall clock;
- logical order: which event must come after another because one caused or observed the other.
Google’s solution made physical-time uncertainty explicit. CockroachDB’s solution combines approximate physical time with logical ordering and enforces a bound on clock disagreement. They share goals, but they are not the same design.
What Google Spanner actually is
Spanner is Google’s globally distributed transactional database. Google published its architecture in 2012 after years of operating systems such as Bigtable and a manually sharded MySQL estate behind its advertising platform.
Bigtable could distribute enormous datasets, but it did not provide the relational model and general transactions that many applications wanted. Sharded MySQL provided familiar relational behavior inside a shard, but changing the shard layout and coordinating work across shards became a large engineering burden. F1, Google’s advertising database, was built on Spanner to combine a richer SQL system with synchronous replication and global transactions.
Spanner has two ideas that are often mixed together:
- replicated data is coordinated through Paxos groups;
- a clock service called TrueTime exposes bounded uncertainty and lets transactions receive timestamps with useful real-world ordering properties.
The atomic clocks do not replace consensus. Paxos still decides the replicated history. TrueTime helps Spanner assign and reason about timestamps across those groups.
TrueTime returns an interval, not a magic timestamp
An ordinary clock API returns something like:
10:00:00.123
TrueTime returns an interval:
[earliest = 10:00:00.119, latest = 10:00:00.127]
The contract is that the real time is somewhere inside that interval. The half-width is often written as epsilon. A narrow interval means the system is highly certain. If synchronization has been unavailable and oscillators have been drifting, the interval grows.
That honesty is the important part. Spanner does not claim perfect clocks. It turns uncertainty into data the transaction protocol can use.
Why there are GPS antennas on data-center roofs
GPS is not only a navigation system. Its satellites broadcast precisely timed signals derived from onboard atomic clocks. A receiver that hears several satellites can calculate position, but a fixed receiver can also use those signals as a highly accurate time source.
The antenna sits where it can receive satellite signals—typically with a clear view of the sky. Cabling carries the signal to a GPS receiver or time master inside the facility. Database servers do not each query a satellite before committing a row. Time masters derive reference time, and other machines synchronize from those masters.
GPS has failure modes: antenna damage, receiver faults, interference, bad satellite data, and environmental or operational mistakes. Depending on only GPS would turn one impressive clock source into a shared risk.
Why atomic clocks are inside the data center
An atomic clock measures the frequency of transitions in atoms and provides a highly stable local oscillator. It can keep time with very low drift even when an external radio signal is unavailable.
Google described a set of time masters in each data center. Most used GPS receivers with dedicated antennas. The remaining masters used atomic clocks. These sources have deliberately different failure modes. GPS ties the facility to an external global reference. An atomic clock can hold stable time locally when GPS is unavailable or untrusted.
Servers called timeslaves poll several masters. They account for communication delay and local oscillator drift, then expose the result as the TrueTime interval. As time passes without a trustworthy synchronization, uncertainty expands. When contact returns, it can narrow again.
The architecture is not “buy an atomic clock and the database becomes consistent.” It is a fleet-level time service with independent references, monitoring, drift assumptions, and an API that carries uncertainty all the way into transaction logic.
Commit wait turns uncertainty into external consistency
Spanner chooses a commit timestamp s for a transaction. Before reporting success, it waits until TrueTime can prove that real time is later than s. In simplified form:
choose commit timestamp s
replicate the commit through consensus
wait until TT.after(s) is true
acknowledge commit to the client
This commit wait means that when the client receives success, the chosen commit timestamp is definitely in the past. If another transaction begins after that response, its timestamp can be ordered after the first transaction. Spanner calls the resulting property external consistency; it is closely related to strict serializability.
Smaller uncertainty means a shorter potential wait. This is why the quality of the time infrastructure matters to database latency.
CockroachDB did not rebuild Spanner in pure software
The popular version of the story is that Google needed atomic clocks, three engineers removed the hardware, and CockroachDB reproduced Spanner in software. It is memorable and wrong in the places that matter.
CockroachDB learned from Spanner, F1, Bigtable, Raft, MVCC systems, and decades of relational database work. It pursued a similar product boundary—distributed SQL with strong transactions—but made different engineering choices:
- Raft rather than Spanner’s Paxos implementation;
- a PostgreSQL-compatible wire protocol and SQL dialect;
- hybrid logical clocks on ordinary machines;
- an explicit maximum clock-offset assumption;
- its own transaction protocol, range leases, optimizer, and storage engine;
- deployment on commodity infrastructure across clouds or data centers.
There is still clock infrastructure. Nodes should run a reliable synchronization service such as chrony or another supported NTP implementation. CockroachDB has not escaped clocks; it has chosen a clock model that operators can provide without a private GPS-and-atomic-clock service.
The correct comparison is not hardware versus no hardware. It is bounded uncertainty as a dedicated time API versus ordinary synchronized clocks plus logical causality and an enforced offset limit.
Hybrid logical clocks preserve causality
Each CockroachDB timestamp combines two values:
(physical wall time, logical counter)
The physical component keeps timestamps close to real time, which is useful for MVCC, historical reads, expiration, observability, and operational reasoning. The logical component orders events when physical time alone cannot.
Assume node A creates timestamp (1000, 0) and sends it with a request to node B. B’s physical clock only reads 995. If B used its wall clock blindly, the caused event would appear to happen before the message that caused it. Instead B advances its hybrid logical clock to at least the received physical value and increments the logical component, producing something like (1000, 1).
If several events happen within the same physical tick, the counter continues:
(1000, 0) -> (1000, 1) -> (1000, 2)
If the physical clock later advances beyond the stored physical component, the logical counter can return to zero.
The important guarantee is causal: if one observed event causes another, the second can receive a greater HLC timestamp even when the receiving machine’s wall clock lags.
An HLC does not prove that every unrelated event across the cluster is perfectly ordered by real time. CockroachDB still needs transaction rules for concurrent work and still needs physical clocks to remain within an assumed bound.
The maximum clock offset is a correctness boundary
CockroachDB’s default maximum clock offset is 500 milliseconds. That is not a claim that clocks normally drift by half a second. Healthy clocks should be much closer. It is the configured outer bound the database uses when reasoning about whether a timestamp from another node might be in the local node’s future.
Nodes exchange clock-offset measurements. If a node determines that it is too far from at least half of the other nodes—CockroachDB documents the shutdown threshold as 80 percent of the configured maximum—it terminates itself rather than continue with a violated correctness assumption.
That behavior can look aggressive until the alternative is stated plainly: an unavailable replica is recoverable; a replica silently inventing an invalid transaction order can corrupt the history the database promised.
Operators therefore still own time:
- use supported, monitored time synchronization on every node;
- alert on offset and synchronization health before the shutdown boundary;
- avoid casually increasing maximum offset to hide a broken clock service;
- validate virtualization and suspend/resume behavior;
- treat clock faults as failure-injection scenarios, not theoretical trivia.
The architecture in one request
A useful view of CockroachDB is a stack of translations:
PostgreSQL client / driver
|
| SQL over PostgreSQL wire protocol
v
SQL gateway: parse -> optimize -> plan -> execute
|
| KV operations and transaction metadata
v
DistSender: find ranges and route requests
|
+----------+----------+
v v v
range A range B range C
lease lease lease
| | |
Raft group Raft group Raft group
| | |
Pebble/MVCC replicas on local stores
Any SQL-capable node can accept a client connection and act as the gateway. That does not mean every node can independently modify every copy. The gateway plans the statement, converts it into key-value operations, finds the relevant ranges, and routes each operation to the replica currently responsible for coordinating that range.
This distinction matters. CockroachDB is often described as active-active or multi-active because applications can connect through multiple nodes. Internally it avoids conflicting multi-writer histories by assigning narrow authority per range and using consensus.
One sorted keyspace, divided into ranges
At the storage boundary, CockroachDB represents data as ordered key-value pairs. A SQL row is encoded under keys that include identifiers such as tenant, table, index, and indexed values. A secondary index becomes another ordered span of keys. System metadata lives in the keyspace too.
Conceptually:
/table/53/primary/1001 -> customer row 1001
/table/53/primary/1002 -> customer row 1002
/table/53/email/alex@example.com -> secondary index entry
The exact binary encoding is more sophisticated, but the consequence is simple: all data lives in one sorted address space.
CockroachDB divides that space into contiguous ranges:
[a, f) range 11
[f, m) range 42
[m, t) range 87
[t, z) range 91
Ranges split as they grow and can also split according to load so one hot portion can move or receive independent lease placement. Small adjacent ranges can merge. Replicas can move between stores as the allocator balances capacity, load, constraints, and locality.
This is the mechanism behind horizontal scale. Adding a node does not make one existing B-tree magically parallel. It gives the allocator another place to put range replicas and leaseholders. Capacity improves as independent ranges and workload spread across nodes.
Finding a key without a central router
The gateway needs to know which range contains a requested key and where that range’s leaseholder lives. CockroachDB stores range descriptors in its own keyspace and uses a hierarchy of meta ranges to locate them. Nodes cache descriptors and routing information.
If a cached descriptor is stale because a range split or its lease moved, the request can be redirected and the cache repaired. There is no single external shard-map service that must be consulted for every row.
This also explains why a split is more than a storage-engine detail. It creates two independently addressable replication and scheduling units from one key span.
Distribution does not remove hotspots
Automatic splitting cannot make a single logical row accept unlimited serialized updates. If every request increments the same counter, reserves the same inventory record, or appends beneath a key prefix that stays in one narrow range, one leaseholder and one consensus group remain busy.
The fix is usually in the data model:
- distribute counters and aggregate them;
- partition work by a real business key;
- use hash-sharded indexes when a sequential indexed key creates a write hotspot;
- shorten transactions and reduce overlapping write sets;
- avoid global singleton rows in a supposedly distributed design.
CockroachDB automates data movement. It cannot infer a contention-free business model.
Each range is its own Raft group
By default, a range is commonly replicated to three voting replicas placed according to constraints and locality. Those replicas form a Raft group. One replica is the Raft leader, and one holds the range lease. Current leader-lease work increasingly aligns those responsibilities, but the concepts remain useful:
- Raft leader: coordinates the replicated log;
- leaseholder: owns the authority needed to serve current consistent reads and coordinate writes for the range.
A write roughly follows this path:
- A client sends SQL to any gateway node.
- The SQL layer produces one or more KV writes.
- The gateway routes each operation to the relevant range leaseholder.
- The leaseholder evaluates concurrency and transaction state.
- The command is proposed through the range’s Raft group.
- A majority accepts the log entry.
- The committed command is applied to each replica’s local MVCC storage.
- The transaction protocol decides when the SQL transaction can report success.
With three voters, a quorum is two. Losing one voter still leaves two. Losing two leaves only one, so that range cannot safely commit new writes. Other ranges may remain available if their quorums are intact.
There is no database-wide Raft group. A cluster can contain hundreds of thousands of ranges and therefore a huge number of small consensus groups. Their leaders and leaseholders can be spread across the fleet, allowing different data to make progress independently.
This is the architectural break from a single writable primary: not “every replica writes anything,” but “many small pieces each have a safe, movable authority.”
What happens when a node dies
Suppose node 2 holds replicas for thousands of ranges and is the leader or leaseholder for some of them.
When it becomes unavailable:
- ranges whose authority was elsewhere continue normally;
- affected Raft groups detect the missing leader and elect a leader from surviving voters;
- leases are acquired or transferred safely;
- in-flight requests may wait, redirect, or retry;
- under-replicated ranges receive replacement replicas when capacity and placement permit;
- the allocator gradually restores the configured replication state.
The data does not fail over as one monolithic database. Authority changes range by range.
No committed entry is accepted on a minority partition. That prevents split brain, but it also means CockroachDB chooses consistency over accepting writes everywhere during an arbitrary network partition. “Survives a region failure” is only true when replica placement leaves a quorum on the surviving side.
If all three voters for a range were placed in one region, losing that region loses the quorum. If voters are split across three regions, losing one can leave two, but every consensus write now includes regional network latency. Topology is part of the guarantee.
Reads, follower reads, and closed timestamps
The leaseholder can serve an authoritative current read without running a fresh consensus round for that read. Its lease establishes that no other replica is simultaneously allowed to act as the current authority for the same range.
Sending every read to a remote leaseholder would make geographically distributed reads expensive. CockroachDB can also serve sufficiently historical reads from follower replicas.
The key is a closed timestamp. A range communicates that no future write will be accepted below a particular timestamp. Once a follower has applied all replicated work through that timestamp, it can safely answer a read at or below it from local state.
The trade-off is explicit:
- current, strongly consistent reads normally follow leaseholder authority;
- follower reads can be local and scalable, but they read at a safe timestamp behind the absolute present;
- placing leases near writers improves write latency for that data;
- placing replicas near readers enables local historical reads;
- none of these settings make a transoceanic quorum round trip disappear.
Data locality should follow an access pattern, not a slide that says “global.”
A transaction is more than several Raft writes
Consider a transfer:
BEGIN;
UPDATE accounts SET balance = balance - 100 WHERE id = 7;
UPDATE accounts SET balance = balance + 100 WHERE id = 9000007;
COMMIT;
The two accounts may live in different ranges with different leaseholders and Raft groups. Committing each update independently would allow a failure between them to destroy money.
CockroachDB coordinates them as one transaction using MVCC timestamps, write intents, and transaction metadata.
MVCC gives values a time dimension
Multi-version concurrency control stores versions of a key at timestamps rather than overwriting one value in place:
account/7 @ (1000,0) -> 500
account/7 @ (1030,2) -> 400
A read at timestamp (1020,0) can see the first version. A later read can see the second. This supports transaction snapshots, historical reads, changefeeds, backup behavior, and concurrency without making every reader block every writer.
Old versions are not retained forever. Garbage-collection settings and protected timestamps determine how long history remains usable. A four-hour-old MVCC version is not a backup merely because it once existed on three replicas.
Write intents are provisional facts
During a transaction, a write is stored as an intent associated with the transaction. Other operations that encounter it do not treat it as an ordinary committed value. They inspect transaction state, wait, push, or help resolve the intent according to concurrency rules.
The transaction also has a record whose state can include values such as PENDING, STAGING, COMMITTED, or ABORTED.
For the transfer, CockroachDB can:
- write an intent subtracting from account 7;
- write an intent adding to account 9000007;
- replicate each range’s commands through its own quorum;
- stage or commit the transaction record once the required writes are proven;
- expose both writes as committed together;
- clean up the intents asynchronously.
CockroachDB’s Parallel Commits protocol reduces the number of serial coordination rounds by allowing the transaction to enter STAGING while proving that all in-flight writes are durably replicated. If recovery encounters a staged transaction, it can determine whether those writes succeeded and finalize the correct result.
The implementation is more nuanced than classic two-phase commit diagrams, but the invariant is familiar: observers must not see only half of a committed transaction.
Serializable does not mean retry-free
CockroachDB uses serializable isolation by default. Concurrent transactions must produce a result equivalent to some valid serial order.
When two transactions race, a timestamp may be pushed forward. CockroachDB attempts to refresh earlier reads and continue automatically when it can prove their results would not change. If it cannot, the transaction must retry at a new timestamp. Some retries happen inside the server; others reach the client as a serialization error.
Applications therefore need a transaction-retry boundary:
retry loop
begin transaction
read and write database state
commit
on retryable serialization failure:
back off and run the complete transaction again
External side effects do not automatically roll back. Sending an email, charging a card, or publishing to another system inside a retryable transaction can happen more than once. Use an outbox, idempotency key, or another deliberate coordination pattern.
This is not a CockroachDB footnote. It is part of the application contract created by optimistic distributed serializability.
Online schema changes are distributed jobs
Adding an index on a large table is not one metadata write. Every range containing relevant rows may need to participate in a backfill while reads and writes continue.
CockroachDB represents schema state with versioned descriptors and runs schema changes as resumable jobs. New and old application work may temporarily observe adjacent descriptor versions. Mutations move through states so concurrent writes can maintain the structure being created while historical rows are backfilled.
A simplified index addition looks like this:
- publish a descriptor version announcing the new index mutation;
- make concurrent writes maintain the new index as required;
- scan existing rows and backfill index entries;
- validate the result;
- publish the index as public;
- resume or roll back safely after failures according to job state.
This avoids a database-wide maintenance window, but it does not make schema changes free. A backfill consumes CPU, storage bandwidth, Raft traffic, and MVCC space. Large changes should be observed, rate-limited where appropriate, tested with production-shaped data, and coordinated with application compatibility.
Online means the system has a protocol for change. It does not mean the change has no operational cost.
Multi-region is a placement decision
CockroachDB lets teams express database regions, survival goals, and table locality. Those abstractions eventually become replica and lease placement.
Three broad workload shapes matter:
Regional data
If most reads and writes for a row belong to one region, placing its leaseholder and a useful quorum close to that region can keep the common path fast. Additional replicas elsewhere provide survival according to the chosen policy.
Global reference data
Data such as product catalogs or configuration may be read everywhere and changed infrequently. Replicas and follower reads can make reads local while writes pay wider coordination costs.
Truly global writes
If users in several continents constantly update the same rows, there is no perfect placement. The system must order conflicting writes somewhere. Moving the lease helps one side and hurts another; scattering replicas changes quorum latency; weakening consistency changes the product semantics.
This is the physics tax. A correct architecture names which operations can tolerate it.
Region survival must be tested per range
A label such as “three-region cluster” is not proof of region survival. The questions are:
- where are the voting replicas for each critical table or partition;
- where is the leaseholder during normal operation;
- which two voters remain after each region loss;
- whether the application endpoint and SQL gateways also survive;
- how long leader and lease transitions take;
- what latency the surviving topology produces;
- whether capacity is sufficient after losing a third of the fleet.
Database replication cannot rescue a DNS record, load balancer, identity dependency, or application tier that exists only in the failed region.
What CockroachDB solves
CockroachDB moves several difficult responsibilities into one database architecture:
- automatic partitioning of an ordered keyspace;
- replication and consensus per partition;
- safe authority transfer without a database-wide primary promotion;
- transactional SQL across partitions;
- replica placement and rebalancing;
- online, resumable schema changes;
- PostgreSQL-wire access through familiar drivers;
- consistent recovery of transaction state after partial failures.
The operational change is significant. A team no longer has to build a custom shard router and a cross-shard transaction story before adding the fourth database server.
What CockroachDB does not solve
It does not provide infinite scale. A hot row is still a hot row. A poor index still creates an expensive query. A transaction that touches hundreds of ranges still creates coordination and failure surface.
It does not provide zero-downtime under every partition. A range without a quorum becomes unavailable because accepting a minority write would violate consistency.
It does not provide zero-latency global writes. Consensus waits for networks and durable storage.
It does not provide complete PostgreSQL compatibility. CockroachDB speaks the PostgreSQL wire protocol and implements a broad SQL surface, but it is not a PostgreSQL fork. Extensions, procedural behavior, system catalogs, locking assumptions, administrative tools, and edge-case semantics must be tested against the compatibility documentation.
It does not replace backups. Raft faithfully replicates an accidental DELETE. Point-in-time recovery and tested restores solve a different problem.
It does not make application transactions idempotent. Retry-safe side effects remain an application responsibility.
It does not eliminate operators. It changes their work from promoting one primary and moving manual shards toward topology, locality, capacity, clock health, range health, hot-key design, upgrades, and recovery testing.
The PostgreSQL compatibility treadmill
CockroachDB’s adoption strategy was pragmatic: keep the client interface familiar. psql, PostgreSQL drivers, and many ORMs can connect because CockroachDB implements the wire protocol and compatible SQL behavior.
The internals are not PostgreSQL. SQL is planned over distributed KV operations; replicas use Raft; local MVCC data is stored through Pebble, CockroachDB’s LSM-based storage engine.
Compatibility is therefore a continuing engineering program rather than inherited behavior. PostgreSQL keeps evolving, and its ecosystem includes extensions that assume PostgreSQL internals. A migration assessment should inventory actual queries, data types, transaction behavior, extensions, ORM-generated SQL, administrative tools, and error handling. “Our driver connected” proves connectivity, not compatibility.
The same directness applies to licensing. Current CockroachDB releases use the CockroachDB Software License, a source-available license rather than an OSI-approved open-source license. Teams evaluating self-hosted deployment should review the current terms and commercial requirements as carefully as technical compatibility.
Failure tests I would run before production
A three-node demo that survives kill -9 is useful, but it proves only the smallest failure.
I would test the system at the boundaries it promises to own.
Kill a node under write load
Measure error rate, p95 and p99 latency, leader and lease recovery, transaction retries, under-replicated duration, and the time required to restore healthy replica counts.
Isolate the old leader instead of killing it
A network partition is more revealing than a stopped process. Confirm the minority cannot commit and that the majority continues only where it has quorum.
Remove a zone and then a region
Validate critical ranges, not only cluster health. Verify the application path, connection recovery, surviving capacity, and actual recovery time objective.
Introduce clock skew safely in a test environment
Observe offset metrics, node self-protection, replacement behavior, and alerts. The first time the team learns the maximum-offset rule should not be during a production hypervisor incident.
Create a hot key and a hot range
Run the business access pattern that concentrates writes. Inspect contention, range QPS, lease placement, retries, and whether load-based splits can help. Then change the data model and compare.
Run a large schema backfill
Watch foreground latency, CPU, disk bandwidth, LSM compaction, range movement, and job recovery after interruption.
Restore from backup
Restore into an isolated environment, validate row counts and application queries, and measure recovery time. Replica count is not a restore drill.
Force client-visible serialization retries
Prove that the transaction wrapper retries the complete unit, applies backoff, and does not duplicate external effects.
Availability is a measured behavior of the whole system, not a property inherited from a database product page.
When I would keep PostgreSQL
I would keep PostgreSQL when one regional primary with managed failover satisfies the recovery objective, the workload fits comfortably on one write node, PostgreSQL extensions matter, and the team benefits more from simplicity than from transparent distribution.
That describes a large number of serious production systems.
I would consider CockroachDB when the requirements are concrete:
- write capacity or data size must spread across nodes;
- manual sharding has become an application and operations burden;
- transactions must remain consistent across those partitions;
- node or zone failures should recover without database-wide promotion;
- regional survival is a tested business requirement;
- data locality and global access are worth the latency and operational cost;
- PostgreSQL-compatible access is sufficient after a real compatibility assessment.
“We may become global” is not enough. Name the failed boundary in the current architecture. If there is no failed boundary, distribution may be complexity purchased before demand.
The complete mental model
CockroachDB starts with one sorted keyspace. It cuts that keyspace into ranges. It replicates each range. Each set of replicas uses Raft to agree on an ordered log. A lease gives one replica authority for current reads and write coordination. Any SQL node can receive a query, but the work is routed to those range authorities.
Hybrid logical clocks keep causal events ordered while remaining close to physical time. A maximum-offset rule bounds how wrong physical clocks may be. MVCC stores timestamped versions. Write intents and transaction records make changes provisional until a distributed transaction can commit atomically. Serializable isolation rejects or retries executions that cannot fit one valid order. Closed timestamps let followers serve safe historical reads. Locality rules decide where the replicas and authority live.
Google Spanner matters because it demonstrated that global distribution, SQL transactions, and strong consistency could coexist. Its TrueTime service uses GPS-backed and atomic-clock-backed time masters to expose bounded uncertainty, then commit wait converts that uncertainty into externally ordered transactions.
CockroachDB did not simply delete the clocks from that design. It chose a different clock contract, a different consensus protocol, a different transaction implementation, and a PostgreSQL-compatible product surface that can run on infrastructure normal companies can operate.
The result is not a database that escaped trade-offs. It is a database that moved the hardest ones into explicit machinery:
- quorum instead of wishful replication;
- leases instead of multiple uncoordinated writers;
- HLCs and clock bounds instead of pretending wall clocks agree;
- retries instead of silently accepting impossible serial histories;
- placement rules instead of pretending geography is free;
- automated ranges instead of an application-owned shard map.
That is why the architecture is interesting. CockroachDB does not refuse to die through branding. It survives the failures its topology has quorum for, stops when continuing would make the history unsafe, and makes the remaining costs visible to the engineers responsible for the system.
Primary technical references
- CockroachDB source repository and architecture reading list
- CockroachDB architecture overview
- CockroachDB SQL layer
- CockroachDB transaction layer
- CockroachDB distribution layer
- CockroachDB replication layer
- CockroachDB storage layer
- CockroachDB: The Resilient Geo-Distributed SQL Database, SIGMOD 2020
- Spanner: Google’s Globally-Distributed Database, OSDI 2012
- F1: A Distributed SQL Database That Scales, VLDB 2013
- Logical Physical Clocks and Consistent Snapshots in Globally Distributed Databases


