Databases / Working draft
Keeping a transaction
together across machines
A ticket service keeps a wallet of credits for each customer. Alice has 100, and Ben has 40. Alice sends Ben 10 credits. We want Alice to end up with 90 and Ben with 50. The total should remain 140, even if a machine fails halfway through.
On one database server, we'd use a transaction to keep the debit and credit together. Now suppose the service's data is spread across machines. Alice's row is on one machine and Ben's on another. Each has other customers too. Sending two SQL updates to two machines doesn't by itself create the transaction we wanted.
Distributed SQL systems combine relational queries and transactions with partitioned, replicated storage. The SQL interface still describes the change in familiar terms. Underneath it, the database has to find each row, keep copies in agreement, and decide whether all the transaction's changes belong in committed history.
BEGIN;
UPDATE wallets SET credits = credits - 10 WHERE id = 'Alice';
UPDATE wallets SET credits = credits + 10 WHERE id = 'Ben';
COMMIT;The balance check must belong to the same transaction, with isolation that protects it from competing transfers. A successful commit keeps these database changes together. Sending an email or charging a payment service isn't included merely because the application does it between BEGIN and COMMIT.
A table becomes a set of ranges
CockroachDB encodes table and index entries as ordered keys and values. It divides this key space into ranges. A range owns a contiguous interval of keys. Alice's encoded wallet key might fall before a range boundary and Ben's after it. The actual boundaries depend on key encoding and how ranges have split. For our example, we'll deliberately put them in separate ranges.
The SQL layer plans the transfer, and requests travel to the ranges holding the affected keys. A small transaction might stay inside one range; another can touch several. Adding a secondary index can introduce more affected keys too. The database tracks range locations so the application doesn't need to hard-code which machine holds Alice. Ranges can split and move as the cluster changes. CockroachDB's distribution layer describes that routing.
Splitting spreads different keys across more storage and compute. It doesn't split one contested balance into independently writable pieces. If thousands of requests all spend Alice's credits, they still need a common ordering for that row. Distributed storage can increase the work done across many customers while one very busy customer remains a bottleneck.
Agreeing inside one range
Keeping three copies protects a range against some failures. But copying alone leaves a question: if two machines disagree, which history should a future writer continue?
CockroachDB uses a Raft consensus group for each range. A leader orders proposed changes in a log and replicates entries to followers. With three voting replicas, a majority is two. The leader counts its own durable entry and waits for another voter's response before the entry can commit under the protocol's rules. The third replica can catch up later. The replication layer explains the implementation.
Two majorities among three voters must overlap. Raft combines that overlap with rules for elections and logs so a new leader preserves committed history. Counting responses is only one part of consensus: replicas must also agree on the ordering and who can lead. If an isolated former leader could accept new committed writes by itself, the two connected replicas could continue a different history. Its local disk is still useful, but it is insufficient authority for that new write. The Raft paper gives the safety rules.
Interactive · Distributed SQL
Which reply completes a range write?
One log entry records a seat reservation. A leads a three-voter replication group; its own durable entry counts as one vote. Two stored entries form a majority. The clock starts when A receives the request, after the outbound trip from the client. The expected client request-to-reply time also includes that outbound trip and the return trip after commit. Both legs are shown in the result.
Choose Three regions, then move the clock to 71 ms. A and B have the entry; C is still waiting. Isolate A to see why its local copy cannot establish a new commit alone.
Controls
Result Invented message and persistence times
What this model leaves out
This is one schematic consensus log entry with a fixed leader and three voting replicas. Entry-and-reply delays are 1, 5 and 7 ms locally, or 1, 71 and 131 ms across regions. The clock starts after the request’s outbound client leg; acknowledgement waits for the return leg. Total client latency includes both legs, assumed symmetric. Placement and network changes start a fresh write. No elections, terms, leader leases, clock uncertainty, contention, batches or transaction protocol are simulated. B and C may form a majority elsewhere when A is isolated; this view follows A’s attempt only. These timings are not CockroachDB benchmarks.
A client response is another event. The entry may commit, then the connection can break before the success reaches the caller. A timeout doesn't prove the write failed. Retrying a transfer without an operation identifier can transfer another 10 credits. The application needs a way to recognise its earlier attempt and discover its outcome.
A fresh read also needs a safe route into the committed history. Reading any nearby follower's disk can return an older value. Systems use mechanisms such as leases, timestamp bounds and safe follower reads to establish what a replica can serve. Majority replication doesn't make every replica instantly current.
Two successful writes can still make a broken transfer
Suppose Alice's range commits the debit. It is stored safely on a majority, and Alice now has 90. Before we ask Ben's range to commit the credit, the coordinating process fails. Ben still has 40. We have durably stored a total of 130.
Nothing went wrong with replication. It preserved exactly what we asked each range to store. The missing part was a common transaction outcome: either both changes commit, or neither does.
An atomic commit protocol supplies that outcome. In a simple version, each participant first records that it is prepared to make its change, then a coordinator records the commit decision durably. The participant's provisional change isn't yet an ordinary committed balance. If the coordinator disappears, another process must discover the recorded outcome or safely decide the attempt can abort. It cannot just guess that the debit succeeded and forget the credit.
CockroachDB represents provisional writes as write intents associated with transaction state. Its Parallel Commits protocol overlaps work and can establish commitment without a separate sequential wait for every phase in that simple account. A reader encountering an intent resolves its status rather than silently presenting an unfinished transfer as committed. The transaction layer explains intents, recovery and the optimized commit protocol.
Interactive · Distributed SQL
Two replicated changes need one outcome
Account A has 100 credits and B has 40. Transfer 10 from A to B; their combined balance should stay 140. They live in different ranges, each with a working majority.
Advance once to replicate the debit, then Lose coordinator. Compare one transaction with two independent writes. Recovery in the atomic model aborts a transaction that has no commit decision.
Controls
Result Committed snapshot and provisional changes
What this model leaves out
This is a deliberately sequential atomic-commit teaching model, not CockroachDB’s optimized Parallel Commits protocol. Each advance completes replication to a majority for that step. Committed values show a snapshot; a current read encountering unresolved provisional changes waits. We omit competing transactions, safe timestamps, election, actual recovery records and ambiguous client timeouts. Recovery aborts before a decision and preserves commit afterwards. Independent writes are explicitly two separate operations, not an SQL transaction.
Atomicity must be paired with appropriate reads. Reading Alice before a completed transfer and Ben afterwards, in two separate requests, can produce a mixed view even though the transfer itself was atomic. A read transaction with a suitable consistent snapshot can see the balances together as they stood before or after the transfer. Isolation also governs what happens when concurrent transactions read and write the same records. CockroachDB's isolation-level discussion compares read committed and serializable contracts.
A conflict may require a transaction retry. Retrying means repeating its reads and checking its assumptions again, not blindly replaying just the failed statement. A newer balance may change whether the transfer is allowed. Application code also needs to keep effects outside the database, such as sending receipts, from happening again on every retry. Retry error guidance explains the cases callers must handle.
Moving a copy can move the wait
Our replicas can be in different zones within one region or spread across regions. Local copies can communicate quickly, but a failure affecting the whole region may take them all away. Spreading voters across regions changes which failures the group can survive and how far messages must travel.
With one of three voters in each of three regions, a leader normally needs a response from another region to commit a new entry. A client near that leader still waits for that trip. A client far from the leader adds travel to and from the database on top. If a transaction touches another range with its leader somewhere else, it can introduce further communication.
Placement therefore follows the workload and the failure you need to survive. Putting a customer's frequently updated data together near its writers can avoid remote participants. Keeping a quorum nearby reduces normal write latency, but losing that region may remove the quorum. Distributing enough voters to retain a majority after a region loss usually brings more distant acknowledgements into the normal write path. Exact outcomes depend on voter placement and the system's failover rules. Topology patterns makes those tradeoffs explicit.
Spanner makes related choices with partitioned data, replicated groups and distributed transactions, but its protocols aren't CockroachDB's. It uses Paxos for replication and TrueTime's bounded clock uncertainty to support transaction timestamp ordering. At serializable isolation, it provides external consistency: an ordering consistent with transactions completing in real time. Having accurate clocks alone doesn't supply this contract; replication and transaction protocols are part of it. See Spanner replication and TrueTime and external consistency.
During a network partition, a group without enough communicating voters cannot keep accepting new consensus writes. Availability isn't a property the label “SQL” removes. A reachable majority may continue elsewhere; unavailable ranges can still block transactions that need them. This is a concrete failure condition, not a claim that any database is simultaneously available and strongly consistent through every partition. Google's discussion of Spanner and CAP distinguishes ordinary availability from that guarantee.
Which work actually needs the distribution?
A service with many independent customers, substantial transactional data, and an explicit requirement to survive regional failures can benefit from distributing its rows while keeping relational transactions. The benefit follows from sharing the workload and arranging surviving copies. It comes with coordination, placement decisions, operational complexity and retry behaviour.
A small application whose writes come from one region may find a conventional relational database simpler. A globally shared counter updated on every request won't become independent just because it lives in a larger cluster. And a reporting query that scans the whole dataset can move substantial data between nodes even when it is expressed as one short SQL statement.
Now imagine Alice and Ben normally belong to different regions. Would putting both wallet rows near Alice solve the transfer's latency? It can reduce coordination for Alice's request, but Ben's local requests become farther from his data, and the desired region-failure tolerance still determines replica placement. We'd need the actual mix of requests and failure requirements to choose. The transaction guarantee tells us which results are allowed; the placement tells us what communication it takes to get one.
← Return to the database field guideSources and model notes
Working draft. Primary documentation and the Raft paper linked beside their claims were accessed on 2 October 2026. The models use invented balances, stages and message delays. They establish no measured latency, capacity or production-readiness claim. The transfer experiment deliberately separates replication and atomic commitment; it does not reproduce CockroachDB Parallel Commits or Spanner's commit protocol.
Teaching references: Sam Who's controlled experiments, Bartosz Ciechanowski's progressive models, and Julia Evans's concrete scenes. The prose and code-native diagrams are original.