DILITH: ScramDB's high-performance distributed commit protocol
Say you move $100 from checking to savings. Checking drops by 100, savings grows by 100, and those two updates have to land together. In ScramDB the two accounts can sit on two different servers, even in two different countries. If one server dies after taking the $100 out and before the other puts it in, your money is gone. Stopping that is the commit protocol's job. Each distributed database ships its own, and the rest of the database can only be as fast and as safe as that one piece.
DILITH is ScramDB's new commit protocol, we invented it in house, and it's SERIALIZABLE on every shard. We raced it against the protocols that power today's distributed databases, CockroachDB's Parallel Commits among them, and DILITH was the cheapest and the fastest of the lot. When we crashed servers and cut networks on purpose, DILITH kept every answer right. Then we wrote machine-checked proofs that it's correct on clusters of any size.
We're publishing the paper with this post, so go check every number yourself: DILITH: Serializable Commit at the Speed of Consensus, Proved for Clusters of Any Size (PDF)
48.6x more transactions.
That's how many more transactions ScramDB commits on TPC-C, the industry's standard order-processing benchmark, with DILITH than with its previous commit protocol. Here's the rest of what we measured:
- 88% lower cost than CockroachDB's Parallel Commits
- 50x faster than Parallel Commits when the cluster spans three regions
- 21x less network traffic for each transaction
- 30x faster New-Order transactions, the heaviest write in TPC-C
- 17x less CPU for each commit
- The only protocol of the five we tested that passed every correctness check under crashes and network failures
- 25 of 25 properties proved in Lean, for clusters of any size
- SERIALIZABLE on all 126 runs, checked by replaying every committed transaction
SERIALIZABLE, the strictest level SQL has
Isolation decides what a transaction can see while others run next to it. At the weaker levels, two people can read the same $100 balance at the same moment, both decide there's enough, and both withdraw it. The database reports both withdrawals as fine, and the bank is $100 short. SERIALIZABLE rules that out: whatever happens, the result matches running the transactions one after another.
It's also the hardest level to deliver, and the slowest, so plenty of distributed databases default to something weaker. Even the TPC-C kit YugabyteDB publishes sets REPEATABLE READ in its config. We built DILITH around SERIALIZABLE before we wrote a line of it. In our harness a checker replayed every committed transaction of every run and compared each node against the replay. DILITH passed that check on all 126 runs.
We also lined DILITH up against 11 published commit designs. None of them combines serializability with the fastest possible commit, no clock and no central service, and DILITH has all four.
What DILITH gives you:
- Serializable transactions across shards. Fire a thousand transactions at once and every node ends up with a result you could have gotten by running them one at a time.
- Correct reads from any replica. A replica can lag the leader for a moment, and whatever it answers still fits that same order.
- No clock and no central service on the commit path. No timestamp oracle, no sequencer, no coordinator tier. Servers with drifting clocks don't change your results.
- Recovery from any node. Pull the plug on the server running your transaction and another server picks it up and finishes it.
- Bounded memory. The memory DILITH uses has a hard ceiling, backed by a proof.
- Fresh analytics while transactions commit. Queries see the rows that were just committed, with a freshness lag of 0.0 ms.
Before and after on TPC-C
We ran TPC-C twice on one machine. The first run used ScramDB's previous commit protocol; the second used DILITH. Data, workers and the two-minute window stayed identical.
Each segment equals the previous commit path's total for the run.
| Commit path | Committed, relative |
|---|---|
| Before | 1.0x |
| DILITH | 48.6x |
The old protocol was Percolator-style two-phase commit. Plenty of distributed databases still run on that design today. Latency and CPU dropped along with the jump in commits:
- 48.6xmore transactions committed
- 30.8xfaster New-Order, median
- 20.7xfaster New-Order, p99
- 16.8xless CPU per commit
Head to head with CockroachDB's Parallel Commits and three more
How does that stack up outside our own walls? We built four published protocols from their papers and ran them next to DILITH:
- Parallel Commits, the protocol CockroachDB uses
- Sundial, from VLDB 2018
- Classic two-phase commit, with both phases run in parallel
- Sequential two-phase commit, as used by Percolator
Each of the five got the same cluster harness, the same scenarios, the same seeds and one shared checker. Scoring combines three things you'd notice as a user: how long transactions take, how fast and fresh analytical queries are, and how many bytes every commit costs. Lower wins.
| Protocol | Cost |
|---|---|
| DILITH | 1.0x |
| Sundial | 1.7x |
| Sequential 2PC | 5.9x |
| Parallel Commits | 8.3x |
| 2PC, parallel phases | 8.5x |
CockroachDB's Parallel Commits ends up at 8.3 times DILITH's cost, and even Sundial, the closest of the four, costs 1.7 times as much.
Median latency per transaction. Flip between hot keys on a single rack and a cluster stretched over three regions:
| Protocol | Contention | Three regions |
|---|---|---|
| DILITH | 36 ms | 304 ms |
| Sundial | 40 ms | 320 ms |
| Parallel Commits | 1.6 s | 15.2 s |
| 2PC, parallel phases | 1.6 s | 16.5 s |
| Sequential 2PC | 1.1 s | 30.0 s |
With three regions, sequential two-phase commit mostly never finishes. Its clients hit the 30 second timeout and give up.
Here is the score broken into its parts, next to Parallel Commits:
| Term | Times lower |
|---|---|
| Wire bytes per commit | 21.1x |
| Log bytes per commit | 17.9x |
| Transaction tail, queries running | 9.5x |
| Transaction median | 6.2x |
| Transaction tail | 3.8x |
| Query tail | 2.1x |
| Query median | 1.8x |
| Query median, transactions running | 1.3x |
Per commit, 21.1x fewer bytes cross the network and 17.9x fewer land in the log. Two parts, freshness lag and bytes per query, are missing from the chart because both sit at the minimum already.
Under faults
Then we started breaking things. Regions got partitioned, servers crashed, leaders were killed, and in one suite a server went down every few seconds while logs were being compacted. A run only counted as a pass if every answer was correct, every transaction finished, and memory stayed flat.
| Suite | DILITH | Sundial | Parallel Commits | 2PC, parallel phases | Sequential 2PC |
|---|---|---|---|---|---|
| Every gate, seven scenarios | 19 of 19 | 15 of 19 | 11 of 19 | 5 of 19 | 14 of 19 |
| Every gate, 25 deployments | 70 of 75 | 59 of 75 | 34 of 75 | 7 of 75 | 35 of 75 |
| Correct history, 25 deployments | 75 of 75 | 68 of 75 | 47 of 75 | 43 of 75 | 51 of 75 |
| Correct history, restore stress | 24 of 24 | 14 of 24 | 3 of 24 | 2 of 24 | 2 of 24 |
| 3 regions, 27 nodes, partitions | passes | fails: unbounded state | fails: abandoned work | fails: abandoned work, unbounded state | fails: abandoned work, unbounded state |
| Crashes, leader kills, partitions, 9 nodes | passes | fails: abandoned work, wrong answers, unbounded state | fails: abandoned work | fails: abandoned work, unbounded state | fails: abandoned work, unbounded state |
Every run DILITH did ended with a correct history. Crashes, leader kills and partitions cost it zero transactions. Across 3 regions and on 9 nodes under faults, it's the only one of the five that passed every check.
Analytics on the same nodes
Transactions and analytics share the same nodes in ScramDB. On the scenario with extra read replicas, median query time was 12 ms under DILITH and 21 ms under the runner-up. Freshness lag: 0.0 ms. Your dashboard shows the row committed a moment ago.
Proved
- 25/25properties proved in Lean, for clusters of any size
- 34/37deliberately broken variants caught by the model checker
- 630runs behind the paper, 158.8 CPU hours
The specification was model checked with crashes, restarts, leader changes and partitions thrown in by the checker itself, and the Lean 4 proofs cover clusters of any size. We also fed the checker 37 deliberately broken copies of the protocol. It caught 34. The 3 it missed had each dropped a rule that the other rules already cover, so nothing actually broke.
Read the paper
The harness results above all come from the paper. It lists every scenario, baseline and gate:
DILITH: Serializable Commit at the Speed of Consensus, Proved for Clusters of Any Size (PDF)
When
DILITH is landing in ScramDB right now, and those TPC-C numbers at the top were measured on phase one. More is coming soon!

