A distributed transaction is expensive because of what happens while its locks are held. In a classic partitioned database, a transaction that touches two machines takes locks on both, does its work, then runs two-phase commit so both sides agree to commit. Every lock is held across at least two network round trips, and a slow participant holds them longer. The Calvin paper (Thomson, Diamond, Weng, Ren, Shao and Abadi, SIGMOD 2012) calls this window the contention footprint. On a hot record, throughput is roughly one transaction per footprint, however many machines you add.
Calvin moves agreement to before execution. Replicas first agree on the order of transaction inputs. After that, every node executes the same transactions in the same order deterministically, so there is nothing left to agree about at commit time. This article builds the idea up, walks a two-partition example produced by running the code shown, and explains the costs. All figures are from the paper.
The architecture
The core idea: agree on inputs, not outcomes
Start with a single-node fact. If a database starts in state S and applies the same transactions in the same order, and each transaction is a deterministic function of what it reads, it always ends in the same state. Replication then only needs to ship the ordered input log, not the effects. Two replicas that read the same log are identical without exchanging a single page.
Distributing it raises three questions: who decides the order, how nodes execute concurrently yet match that serial order, and what happens when a transaction cannot finish. Calvin answers with three layers: a sequencer for ordering and replication, a scheduler for deterministic locking and execution, and an unmodified storage layer behind a CRUD interface. The storage engine knows nothing about transactions, which is why the paper can plug in both memory and disk engines.
The sequencer: epochs, batches and replication
Every node runs a sequencer. Time is cut into 10 ms epochs. During an epoch each sequencer collects client requests, and at the end it closes them into a batch. The batch is then replicated in one of two modes. In asynchronous mode, one replica is master: its sequencers forward each batch to the matching sequencers in other replicas. Latency is minimal, but failover is hard, because everyone must agree which batch was the failed master's last one and exactly what it held. In Paxos mode, the sequencers in a replication group agree on each epoch's batch. The paper's implementation used ZooKeeper for this, and the authors note a leaner Paxos would cut latency, not raise throughput.
After replication, each sequencer sends every scheduler in its replica a message with its node id, the epoch number and the transactions that scheduler takes part in. Each scheduler builds the global order by interleaving all sequencers' batches for the epoch in a fixed round-robin order. No coordinator hands out sequence numbers, and every scheduler, in every replica, computes the same order independently. The cost is latency. A request waits for its epoch to close, about 5 ms on average, plus replication. In the paper's measurements, replication mode did not change throughput, because it happens before any lock is taken.
The scheduler: deterministic locking and five phases
Each scheduler locks only the records stored on its own node. The protocol resembles strict two-phase locking with two extra invariants. First, if A precedes B in the global order and both want record R, A requests its lock first. Calvin enforces this by making lock requests from a single thread that walks the order and asks for every lock a transaction will ever need. Second, locks are granted strictly in request order. Together these make the outcome equal to the serial order, and they make deadlock impossible, because nobody ever waits for a lock requested after their own. The price is that every transaction must declare its full read and write sets before it is sequenced.
Once a transaction holds all its local locks, a worker runs five phases: (1) analyse the read and write sets and identify the active participants (nodes holding something it writes) and passive participants (nodes it only reads from); (2) perform local reads; (3) send those reads to every active participant; (4) on active participants, collect the remote reads; (5) run the transaction logic and apply only the local writes. Passive participants finish after phase 3. Every active participant runs the full logic on identical inputs, so each one computes the same result and applies its own share of the writes. Nobody asks anybody whether to commit.
Why there is no two-phase commit
Why is 2PC safe to drop? It exists to handle two cases: a participant that crashes mid-transaction, and a participant that wants to abort. Calvin handles both differently. A crash does not abort anything. The input log is replicated, so a node recovers by loading a checkpoint and replaying the log, or clients fail over to another replica that is executing the same order. A logic abort (insufficient funds, a constraint check) is a deterministic function of the reads, and every active participant has all the reads, so all of them reach the same verdict without exchanging votes. What remains of the footprint is local lock time plus one one-way message carrying reads, never a commit protocol. See two-phase commit in depth for the protocol being replaced.
Worked example: four transactions, two partitions
Take two partitions: P1 stores A and B, P2 stores C and D. Start with A=100, B=7, C=50, D=10. In one epoch, sequencer 1 receives T1 (move 30 from A to C if A has it) and T3 (increment B). Sequencer 2 receives T2 (C += D) and T4 (D += A). Round-robin interleaving gives the global order T1, T3, T2, T4 on every scheduler in every replica.
from collections import defaultdict, deque
PART = {"A": 1, "B": 1, "C": 2, "D": 2} # key -> partition
class Txn:
def __init__(self, tid, reads, writes, logic):
self.tid, self.reads, self.writes, self.logic = tid, set(reads), set(writes), logic
def lock_requests(order, partition):
"""One partition's single lock thread: walk the global order, enqueue every local lock."""
queues = defaultdict(deque)
for t in order:
for k in sorted(t.reads | t.writes):
if PART[k] == partition:
queues[k].append((t.tid, "X" if k in t.writes else "S"))
return queues
def run(order, db):
"""The result every replica must reach: the serial execution of the global order."""
for t in order:
db.update(t.logic({k: db[k] for k in t.reads}))
return db
T1 = Txn("T1", "AC", "AC", lambda r: {"A": r["A"] - 30, "C": r["C"] + 30} if r["A"] >= 30 else {})
T3 = Txn("T3", "B", "B", lambda r: {"B": r["B"] + 1})
T2 = Txn("T2", "CD", "C", lambda r: {"C": r["C"] + r["D"]})
T4 = Txn("T4", "AD", "D", lambda r: {"D": r["D"] + r["A"]})
order = [T1, T3] + [T2, T4] # sequencer 1's batch, then sequencer 2's, round-robin
print(run(order, {"A": 100, "B": 7, "C": 50, "D": 10}))
# {'A': 70, 'B': 8, 'C': 90, 'D': 80}| Partition | Record | Lock queue (in request order) |
|---|---|---|
| P1 | A | T1 X, then T4 S |
| P1 | B | T3 X |
| P2 | C | T1 X, then T2 X |
| P2 | D | T2 S, then T4 X |
Read the queues as a schedule. T1 and T3 get all their locks at once and run in parallel. T1 is multi-partition: both P1 and P2 are active, they swap A and C, and each runs the same logic, so P1 writes A=70 and P2 writes C=80. T3 is local to P1. T2 waits only for T1 to release C, then computes C=90 on P2. T4 is the interesting one. It writes only D, so P2 is active and P1 is passive: P1 waits for T1's lock on A, reads A=70 and forwards it, then is done. T4 also waits for T2's shared lock on D, then computes D=10+70=80. The final state is exactly the serial result, A=70, B=8, C=90, D=80, and a replica replaying the same log reaches it without any coordination. Change the order and the answer changes (T4 before T1 would give D=110), which is why agreeing on the order is the whole game.
Dependent transactions and OLLP
Declared read and write sets are the main restriction. A transaction that must read data to know what it will touch, such as "update the order whose secondary-index key is X", is a dependent transaction. Calvin handles these with Optimistic Lock Location Prediction (OLLP). The client first runs a cheap, low-isolation, unreplicated, read-only reconnaissance query to find the read/write set, then submits the real transaction with that set. At execution the transaction re-checks what it read. If the set has changed, it is deterministically restarted. This works when the lookup keys change rarely, as with an index on an item name rather than a stock quantity. The paper notes that TPC-C's Payment transaction is dependent but never restarts, because the benchmark never changes that index. On a workload where the indexed field churns, OLLP can restart repeatedly, so measure the restart rate.
Disk, prefetching and the latency trade
Determinism has a cost on disk. A conventional engine can let B and C overtake a transaction A that is stalled on I/O. Calvin cannot reorder, so everything behind A waits. Calvin's rule is to do heavy work before locks are taken. If a transaction may touch cold data, the sequencer delays forwarding it and asks storage to prefetch the records, so execution only touches memory. Too short a delay stalls the transaction under locks; too long costs latency and memory. In the paper, reaching 99% prefetch completion on its file-based engine needed a 40 ms artificial delay under high contention. At low contention the default batching delay was enough. On that low-end hardware, throughput held while no more than 0.9% of transactions touched disk, and local disk throughput, not contention, was the limit. The sequencer also needs to know which keys are in memory across all nodes, which the authors say does not scale as implemented.
Checkpoints and recovery
Because only inputs are logged, recovery means loading a checkpoint and replaying the log from it. Calvin offers three checkpoint modes. The naive mode freezes one replica and snapshots it, which the other replicas hide from clients, but the frozen replica falls behind and must catch up. The second is a variant of Cao et al.'s Zig-Zag algorithm. It picks a virtual point of consistency, a position in the global order. Transactions before it write "before" versions of records, and later ones write "after" versions. Once everything before the point has run, the before versions are immutable and a background thread writes them out. No quiescing is needed, only an order everyone agrees on. The third mode uses the storage engine's own multiversioning when it has it.
What the numbers show
On TPC-C New Order with 10 warehouses per node and about 10% of transactions spanning two machines, the paper reports roughly 5,000 transactions per second per node beyond 10 nodes, scaling linearly to "nearly half a million" per second on 100 commodity EC2 nodes. Each partition has at most 100 districts, so at most 100 New Orders can run per node at once, and throughput depends on how quickly locks are released. Holding them through 2PC would cap it.
Read these as evidence for the footprint argument, not a benchmark to match; Calvin is a research system. Compare Spanner and TrueTime, which keeps 2PC and makes it cheaper with synchronised clocks, and 2PC and 3PC for the classic baseline.
Failure modes
- Stragglers stall the order. A slow node delays every transaction queued behind it on any record. A conventional engine can reorder around it; Calvin cannot. Watch per-node lag in execution of the global order.
- Nondeterministic logic. Reading the clock, using random numbers or iterating an unordered map inside a transaction makes replicas diverge silently. Put such values into the input at sequencing time, and compare replica checksums at checkpoints.
- OLLP storms. A volatile index makes reconnaissance stale and transactions restart again and again. Alert on the restart rate.
- Async-mode failover. Losing the master sequencer before its batch reaches the other replicas loses those inputs. Use Paxos mode when you cannot tolerate that.
- Mis-tuned prefetch delay. Too short stalls under locks; too long inflates latency and memory.
Trade-offs
Calvin trades latency and flexibility for throughput under contention. Every transaction pays the epoch wait plus input replication, and declared read and write sets rule out interactive transactions. In return, hot multi-partition transactions stop collapsing throughput and replication of any strength costs no lock time. It fits many short, known-shape transactions over contended partitioned data; 2PC designs fit interactive or latency-critical work. Distributed transactions in practice compares the wider field.
What to do next
- Run the simulation above, then add a fifth transaction that conflicts with T4 and predict its lock queue position before running it.
- List your hot multi-partition transactions and check whether each one's read and write sets can be declared up front or need OLLP.
- Estimate the lock hold time each one has today, including 2PC round trips, to see how big your contention footprint really is.
- Audit transaction code for nondeterminism (clocks, randomness, unordered iteration) before considering any deterministic or log-replay design.
- If you prototype one, measure epoch latency, the OLLP restart rate and per-node execution lag from day one.
- Read the paper's microbenchmark section to see how throughput changes with contention and the share of distributed transactions.