CockroachDB looks like PostgreSQL from the client side: it speaks the PostgreSQL wire protocol and most of its SQL. Underneath, it is a sorted, transactional key-value store cut into ranges, and each range is kept consistent by its own Raft group. Once you understand that layering, its behaviour becomes predictable: why a transaction touching rows on two ranges costs more than one touching a single range, why your client must retry SQLSTATE 40001, why a sequential primary key creates a hotspot, and why a node with a bad clock shuts itself down.
This article walks the stack from SQL to disk, follows one money transfer through it, and ends with retry, schema and operations guidance. The defaults quoted were checked against the current documentation; they change between versions, so confirm them for yours.
The five layers
Every statement passes through five layers. The SQL layer parses and plans the query and turns rows into key-value operations. A row in table 53 with primary key 1000 becomes a key shaped like /Table/53/1/1000, with the columns encoded into the value, and every secondary index entry is another key. The transaction layer gives those operations a timestamp from a hybrid logical clock and makes them atomic and isolated. The distribution layer finds the range holding each key and sends the request there. The replication layer runs Raft for that range, so the write survives losing a node. The storage layer, Pebble (an LSM-tree engine), stores each key with its MVCC timestamp, so several versions can coexist.
So SQL concepts map to key spans. A primary-key range scan is one contiguous span; a secondary-index lookup plus a row fetch is two, often on different ranges and nodes.
Ranges: how the keyspace is cut and found
The whole keyspace is one sorted map, cut into contiguous ranges. In current versions a range splits when it grows past range_max_bytes (512 MiB by default). Adjacent ranges can merge when they shrink below range_min_bytes (128 MiB by default). Ranges also split on load: a range taking a disproportionate share of requests is split, so the halves can live on different nodes. Both thresholds and the replica count are zone configuration settings, which you can set per database, table, index or partition.
To find a key, a node consults the meta ranges, a two-level index of range descriptors that is itself stored in ranges. Nodes cache descriptors, so a request normally goes straight to the right node; a stale entry costs a redirect and a refresh.
-- Where does this table live, and which node holds each lease?
SHOW RANGES FROM TABLE bank.accounts;
-- Five replicas for a table that must survive two simultaneous node failures.
ALTER TABLE bank.accounts CONFIGURE ZONE USING num_replicas = 5;
-- Pre-split a UUID-keyed table before a bulk load, so the load is not one hot range.
ALTER TABLE bank.accounts SPLIT AT VALUES
('40000000-0000-0000-0000-000000000000'::UUID),
('80000000-0000-0000-0000-000000000000'::UUID),
('c0000000-0000-0000-0000-000000000000'::UUID);
Replication: Raft per range and the leaseholder
Each range is replicated, three times by default, and each set of replicas forms its own Raft group. A write commits once a majority (two of three) has appended it to the Raft log; each replica then applies it to its storage engine. Three replicas survive one node failure; five survive two.
Reads do not go through Raft. One replica holds the range lease, and that leaseholder serves consistent reads from its local data and coordinates the range's writes. Because it holds the lease, it knows no other replica can have accepted a newer write, so a local read is safe. In current versions the default is the leader lease, which keeps the lease on the Raft leader except briefly during transfers. Epoch-based leases, which tied lease validity to node liveness, are now disabled by default. Expiration-based leases, which time out after a few seconds, are used for system ranges such as the meta ranges.
When a leaseholder's node dies, the survivors elect a new Raft leader, which acquires the lease; the documentation says this should take a few seconds, during which requests to those ranges stall rather than fail.
Transactions: HLC, intents and the transaction record
Every transaction gets its timestamp from a hybrid logical clock (HLC): a physical component close to wall time plus a logical counter, advanced whenever a node sees a higher timestamp, which keeps causally related events ordered.
A write does not overwrite a value. It lays down a write intent: a provisional MVCC value plus an exclusive lock, replicated through Raft and pointing at the transaction's record. The transaction record holds the status: PENDING, STAGING, COMMITTED or ABORTED. A reader that hits an intent looks up the record. If the transaction committed, the intent counts as the value; if it aborted, the intent is ignored. If it is still pending, the reader waits for the writer or pushes it, depending on priorities and timestamps. Committed intents are resolved into ordinary values asynchronously.
A naive commit takes two consensus rounds: write the intents, then flip the record to COMMITTED. Parallel commits cut this to one. The coordinator writes the record as STAGING, listing the in-flight writes, at the same time as the final writes. Once all of them succeed, the transaction is implicitly committed and the client gets its answer; the record is marked COMMITTED afterwards. Any transaction that finds a STAGING record can check the listed writes itself to decide the outcome.
Isolation defaults to SERIALIZABLE; READ COMMITTED is also supported. The leaseholder enforces serializability with a timestamp cache, which records the latest read timestamp for each key span. A write below that timestamp would invalidate someone's read, so the writer is pushed above it. A pushed transaction must then refresh its reads, checking that nothing it read changed between its old and new timestamps. If the refresh succeeds, it commits at the new timestamp. If not, the client gets a retryable error.
Clocks and the uncertainty interval
Node clocks are only approximately synchronised. A value timestamped slightly above a reader's timestamp may, in real time, have been written before the read. CockroachDB handles this with an uncertainty interval. A read at time t treats any value between t and t plus the maximum clock offset as uncertain, and restarts at a higher timestamp. If that restart cannot happen transparently, the client sees ReadWithinUncertaintyIntervalError. The maximum offset defaults to 500 ms and is set with --max-offset on cockroach start. Keep it the same on every node.
The same reasoning explains why a node with a bad clock crashes on purpose. The documentation says that if observed offsets exceed the limit, servers crash to minimise the chance of reading inconsistent data. Run NTP or chrony on every node, and alert on clock-offset metrics well before the limit. Never fix the problem by raising the maximum offset: that widens every uncertainty window and increases retries for every workload.
Worked example: one transfer across two ranges
Take a transfer of 100 from account A to account B. A's row lives in range r12, whose leaseholder is node 1. B's row lives in range r40 on node 3. The client is connected to node 2.
- Node 2 plans two UPDATEs and becomes the transaction coordinator. Its HLC assigns timestamp t=100.
- The first write goes to node 1. Nothing in r12's timestamp cache is above 100, so node 1 replicates an intent on A through Raft. The transaction record will live on r12, the range of the first key written.
- The second write goes to node 3. A concurrent analytics query read B at t=105, so the timestamp cache pushes this write above 105, and the intent is written at the pushed timestamp.
- At COMMIT, the coordinator sees that its timestamp moved and refreshes its reads. It checks that the spans it read on A and B have no new writes between 100 and 105. They have none, so the refresh succeeds.
- The record is written as STAGING, listing both intents. Once the writes are confirmed, the client gets success after a single consensus round on the commit path.
- Asynchronously, the record becomes COMMITTED and the intents on A and B are resolved into ordinary MVCC values.
Had another transaction committed a write to B at 103, the refresh would fail and the client would receive SQLSTATE 40001 with RETRY_SERIALIZABLE in the message, and would simply run the transaction again.
Client retries are part of the contract
All retryable transaction errors use SQLSTATE 40001 and contain the text "restart transaction". Common causes listed in the error reference include RETRY_WRITE_TOO_OLD (a newer write landed on a key you are writing), RETRY_SERIALIZABLE (your timestamp was pushed and the refresh failed) and ReadWithinUncertaintyIntervalError (clock uncertainty). The server retries some of these itself, but explicit multi-statement transactions need a retry loop in the application. With plain psycopg 3 it looks like this:
import random
import time
import psycopg
from psycopg import errors
def run_txn(conn, fn, max_attempts=8):
"""Run fn(conn) in a transaction, retrying on SQLSTATE 40001.
Connect with autocommit=True so conn.transaction() issues BEGIN/COMMIT."""
for attempt in range(1, max_attempts + 1):
try:
with conn.transaction():
return fn(conn)
except errors.SerializationFailure: # SQLSTATE 40001
if attempt == max_attempts:
raise
time.sleep(random.random() * min(1.0, 0.02 * 2 ** attempt))
def transfer(conn, src, dst, amount):
conn.execute("UPDATE accounts SET balance = balance - %s WHERE id = %s", (amount, src))
conn.execute("UPDATE accounts SET balance = balance + %s WHERE id = %s", (amount, dst))
with psycopg.connect("postgresql://app@crdb:26257/bank", autocommit=True) as conn:
run_txn(conn, lambda c: transfer(c, "3f1c9a2e-7b4d-4e8a-9c61-2d5f0b8e1a47",
"b70e4d19-52a3-4c6f-8e2b-91d7a0c3f655", 100))The transaction function can run more than once, so keep side effects such as sending email out of it. Keep transactions short, because the longer one stays open, the more likely its timestamp is pushed and its refresh fails. Watch for a connection that drops during COMMIT: the outcome is then ambiguous. Make writes idempotent, for example with an idempotency key in a unique column, and check that key before retrying.
Schema design for a range-partitioned store
- Avoid monotonically increasing primary keys. A SERIAL or timestamp-first key sends every insert to the last range, so one leaseholder takes all the writes however many nodes you have. Use
gen_random_uuid(), or a hash-sharded index, which prefixes a computed shard column to spread sequential values across several ranges. - Group rows that are read together under one key prefix. A composite key such as
(tenant, id)keeps a tenant's rows contiguous, so a tenant query hits one span on one or a few ranges. - Count secondary indexes as writes. Each index entry is another key, often on another range, so every index adds an intent, and possibly another leaseholder, to each commit.
CREATE TABLE app.events (
id UUID NOT NULL DEFAULT gen_random_uuid(),
ts TIMESTAMPTZ NOT NULL DEFAULT now(),
tenant INT8 NOT NULL,
payload JSONB,
PRIMARY KEY (tenant, id) -- one tenant's rows stay contiguous
);
-- Time-ordered lookups without every insert landing on one tail range.
CREATE INDEX events_by_ts ON app.events (ts) USING HASH;
-- A dashboard that tolerates a few seconds of staleness reads from the nearest replica.
SELECT count(*) FROM app.events AS OF SYSTEM TIME follower_read_timestamp()
WHERE tenant = 7 AND ts > now() - INTERVAL '1 hour';The last query is a follower read. Each range keeps a closed timestamp trailing the present by a few seconds, below which no new writes can appear, so any replica, including a nearby one, can serve reads at follower_read_timestamp() without contacting the leaseholder.
Failure modes
- Hot range: one node's CPU is pegged while the others idle. The cause is usually a sequential key or one very popular row. Find the range on the DB Console's hot ranges page and fix the key design.
- Retry storms: contention on a few rows produces many 40001 errors, and the retries collide again. Add jittered backoff, shrink the transaction, and use
SELECT ... FOR UPDATEso contenders queue instead of failing at commit. - Clock trouble: a node exits with clock-offset errors. Fix time synchronisation rather than raising the offset.
- Lost quorum: with two of three replicas down, a range is unavailable. Use node locality so one zone cannot take out two.
- Huge transactions: a bulk UPDATE in one transaction holds intents for a long time, blocks readers and is likely to be pushed. Work in chunks of a few thousand rows.
Trade-offs
| Choice | You gain | You pay |
|---|---|---|
| 5 replicas instead of 3 | Survive two simultaneous failures | More write fan-out and 5/3 the storage |
| SERIALIZABLE (default) | No anomalies; simpler reasoning | More retries under contention |
| READ COMMITTED | Fewer retries | Application must tolerate anomalies |
| UUID primary keys | Writes spread across ranges | No natural order; time scans need an index |
| Hash-sharded index | Sequential writes spread | Range scans fan out to every shard |
| Follower reads | Local, cheap reads | Data a few seconds stale |
What to do next
- Run
SHOW RANGES FROM TABLEon your three busiest tables and note where the leaseholders sit relative to your clients. - Find every multi-statement transaction in your code and confirm it has a 40001 retry loop with jitter and no side effects inside.
- List tables with sequential primary keys and plan a move to UUIDs or a hash-sharded index.
- Confirm NTP or chrony on every node and add an alert on clock offset well below 500 ms.
- Verify that replica placement survives losing a whole zone, not just one node.
- Move read-only dashboard queries to
follower_read_timestamp()where a few seconds of staleness is acceptable. - Go deeper on the mechanisms: MVCC, serializable isolation, hybrid logical clocks, the Raft algorithm and two-phase commit, which parallel commits improves on.