Sorted Runs

LSM trees, from first principles to production systems

Why So Many Database Engines? The RUM Conjecture Explained

ep2lsm treerum conjecturecompactionbloom filterstorage engines

Corrections to the video

The article below has these right. The video does not.

  • high The Bloom filter probability is backwards. The video says the filter is wrong about 1% of the time, so "probably yes" usually means yes. The 1% applies to keys that are not in the run: about 1 in 100 of them still get "probably yes".
  • medium RUM is called "the proof" and "math says you can't". It is a conjecture, not a proven theorem.
  • high Tiered compaction is described as "no cascade". Runs still merge down to the next level. What tiering avoids is rewriting the runs already sitting there.
  • medium "Cassandra 5.0 moved to a unified strategy": UCS is recommended in 5.0, but the default is still size-tiered (STCS).
  • low "Before that, the strategies had no shared language" overstates it. The leveled and size-tiered names already existed, and Cassandra shipped both.

LevelDB, RocksDB, Pebble, Cassandra, ScyllaDB, TigerBeetle, TiKV, YugabyteDB: dozens of storage engines are built on the same idea, the LSM tree. If the structure is the same, why hasn’t somebody built the best one and called it a day?

Because every engine on that list is paying a different bill. The shared idea, in short: an LSM tree buffers writes in memory, flushes sorted runs to disk and compacts them in the background. Writes are fast, reads can be slow, and compaction costs something. Those costs trade off against each other; this article gives that trade-off its research name and looks at who chose which side of it.

The RUM Conjecture: pick two

Picture a triangle of write, read and space amplification: you cannot push all three down at once. That claim has a name. Athanassoulis, Kester, Maas, Stoica, Idreos, Ailamaki and Callaghan call it the RUM Conjecture: the three overheads are Read, Update and Memory. In the paper’s words:

designing access methods that set an upper bound for two of the RUM overheads, leads to a hard lower bound for the third overhead which cannot be further reduced.

A triangle with corners labeled update (write) cost, read cost and memory (space) cost, and the word RUM in the middle: pick two to keep low, the third gets worse
The write, read and space amplification triangle, with its research name.

Two details matter. First, it is a conjecture, not a proven theorem: the paper argues it from the lower bounds of access methods tailored to each single overhead, and it is the kind of constraint you hit in practice, a bit like CAP for storage (that comparison is an analogy, not something the paper claims). Second, you do not engineer around it. You choose which penalty to accept. The paper’s own Figure 1 places popular structures in that space: hash indexes, B-trees, tries and skiplists on the read-optimized side, LSM trees and other differential structures on the write-optimized side, and Bloom filters, bitmaps and sparse indexes on the space-optimized side. Mark Callaghan, whose blog post Read, write & space amplification: B-Tree vs LSM is the clearest practical treatment, is a co-author, listed with Facebook.

For an LSM tree the question becomes concrete: how aggressively do you compact? Compact a lot and reads get cheap but writes get expensive. Compact lazily and writes get cheap but reads pay. Two strategies sit at the ends of that spectrum, and both are described side by side in the Dostoevsky paper: leveling and tiering.

Leveled compaction: organized shelves

In leveled compaction, every level except level 0 holds a single sorted run, split into files whose key ranges do not overlap. LevelDB’s implementation notes say it directly: files in the young level (level 0) may contain overlapping keys, but files in other levels have distinct, non-overlapping key ranges. Dostoevsky puts the policy in one line: with leveling, “we merge runs within a level whenever a new run comes in”.

That structure makes reads predictable. To find key G, you check level 0, then at most one file in each deeper level, because only one file can own that key range.

Leveled compaction read for key G: level 0 misses, one file A-M in level 1 misses, and the single file G-I in level 2 holds the key; only one file is checked per level
Looking up G in a leveled tree: one file per level, three checks.

The write cascade

The price is on the write side. When level 0 fills up, its files merge into level 1, and to keep level 1 clean and non-overlapping, they must be merged with the existing data there. In RocksDB, L0 files usually overlap, so all of them are picked, and then at least one L1 file is merged “with the overlapping range” of L2 when L1 outgrows its target. Each level is about ten times larger than the one above it (LevelDB’s notes: 10 MB for level 1, 100 MB for level 2, and so on).

ScyllaDB’s engineers walk through the arithmetic: one L1 file covers roughly a tenth of the key space, one L2 file a hundredth, so compacting a single L1 file means merging it with around ten overlapping L2 files. A small overflow at the top becomes a large rewrite at the bottom.

Three stacked levels: L0 with 4 files, L1 of 10 MB, L2 of 100 MB, with arrows showing L0 merging into overlapping L1 files and L1 overflowing into roughly ten L2 files
Each level is about 10x the one above, so each overflow drags more settled data into the merge.

One thing the cascade is not about: safety. The write-ahead log already made the data durable. Compaction exists to make reads fast.

How bad is the write amplification? It depends on the size ratio, the number of levels and the workload. In the experiment ScyllaDB describes, the leveled strategy wrote 111 GB for the dataset they were loading, “13-fold write amplification, and over twice the write amplification of STCS”, and they note that individual writes can be amplified 50-fold. So the picture is: you write 1 GB of your own data, and the disk may move ten, twenty or more times that over its lifetime. Callaghan adds that leveled is usually worse than tiered here but competitive in a few cases, such as key-order inserts and skewed writes.

Reads cheap, writes expensive. This is RocksDB’s default compaction style, and LevelDB uses it too. If your workload is read-heavy (a user-facing app, a serving layer), leveled is the natural bet.

Tiered compaction: the messy desk

Tiered compaction flips the rule. Multiple sorted runs can coexist at each level, and their key ranges are allowed to overlap. Per Dostoevsky: “with tiering, we merge runs within a level only when the level reaches capacity.” When a level fills, all its runs merge into one run that moves down.

Writes get cheap. New data lands next to what is already there. When a level fills, its runs merge and move down as a new run, without rewriting the runs already waiting in the next level. A data point from the same ScyllaDB experiment above: leveled’s write amplification was over twice that of size-tiered. RocksDB’s own docs say the same in general terms: tiered “provides far better write amplification with worse read amplification”. (RocksDB calls its tiered style universal compaction; Callaghan notes every other LSM calls it tiered.)

The bill comes on reads. A key could be in any run, so a lookup may have to check every run at every level.

Tiered compaction read for key G: four runs in each of three levels, numbered one to twelve, with overlapping key ranges such as B-W and A-T, all of which may have to be checked
Four runs per level, three levels: up to 12 checks for one key, and it gets worse with more levels.

There is a space cost too. During a merge, the old runs and the new merged run exist at the same time; Cassandra’s docs warn that the disk needed for both the old and new SSTables during size-tiered compaction can outstrip the space a node has. ScyllaDB measured almost 8-fold space amplification for size-tiered in its own overwrite experiment. The “double the data during a merge” rule of thumb is the mild version of that.

Cassandra’s size-tiered strategy is documented as the default if you specify nothing and as recommended for write-intensive workloads (the same page now recommends the Unified Compaction Strategy for most workloads starting with Cassandra 5.0, which blends both approaches). For years, if you were ingesting logs, telemetry or sensor data and reads were rare or latency-tolerant, tiered was the answer.

Two panels: leveled (one run per level, cheap reads, expensive writes, default in LevelDB and RocksDB) and tiered (many runs per level, expensive reads, cheap writes, long-time default in Cassandra size-tiered)
Same data, same LSM structure, opposite bills.

Bloom filters: the closest thing to a free lunch

Both strategies pay for reads that touch files which do not contain the key. Is there a way to make reads faster without making writes slower? There is a partial one: a Bloom filter attached to each sorted run.

Think of a bouncer with a guest list. If the filter says “not on the list”, the key is definitely not in that run, and that answer is always right. Dostoevsky describes it exactly this way: a Bloom filter cannot return a false negative, though it can return a false positive. When it says “probably”, you still have to read the run, because the answer can be wrong: for a key that is not in the run, the filter still says “probably” about one time in a hundred.

Lookup of Shalders across four runs: three filters answer not here and their runs are skipped, one answers probably, that run is read, and the key is not found; footnote 10 bits per key, 1 billion keys, 1.25 GB of RAM
The filter lets the read skip most runs; the one false alarm costs one wasted read.

The cost is memory. Dostoevsky notes that key-value stores in industry use 10 bits per entry, giving a false positive rate of about 1%, and the RocksDB Bloom filter wiki agrees: its example configuration uses about 10 bits per key, which “works well for many workloads”, and 9.9 bits per key gives a 1% false positive rate. Do the arithmetic for a billion keys: 10^9 keys times 10 bits is 10^10 bits, which is 1.25 GB of RAM. Call it 1.2 GB: enough to skip nearly every wasted disk read across hundreds of gigabytes of data.

It is not free. More keys means more filter memory, and at some point the RAM budget runs out. The RUM Conjecture still holds: there is always a bill somewhere. Here it is paid in Memory.

The system tour: same structure, different bets

With the two strategies and the filter in hand, the engines stop looking like duplicates.

LevelDB is the blueprint. Its authors are Sanjay Ghemawat and Jeff Dean, it has seven levels (kNumLevels = 7) and uses leveled compaction. Cockroach Labs describes LevelDB’s original size as 30k lines of code. Facebook found its single-threaded compaction did not work well for certain server workloads. See how an LSM tree works for the history.

RocksDB is what Facebook built on top of it. It builds on LevelDB to scale on many-core servers and use fast storage efficiently, and it grew from LevelDB’s original 30k lines to 350k+ lines. It is embedded in a lot of other databases: TiKV uses it, the founders of YugabyteDB helped build it and run a modified version, and CockroachDB used it until Pebble. The price of all that power is a configuration surface so large that the official tuning guide admits:

Even we as RocksDB developers don’t fully understand the effect of each configuration change.

That line is from the RocksDB Tuning Guide, under “Final thoughts”.

Pebble is Cockroach Labs’ answer. They wrote it in Go, focused on what CockroachDB needs, at a bit over 45k lines of code, and announced it as the replacement for RocksDB as the default engine in the 20.2 release. Its README lists the RocksDB features it does not implement: column families, universal compaction, FIFO compaction, sub-compactions, transactions and more. A general-purpose engine can lose to one stripped down to what a single database requires. (The full rewrite story: Why CockroachDB Replaced RocksDB with Pebble.)

Cassandra and ScyllaDB run the LSM design across many machines. Instagram described itself as running one of the largest Cassandra deployments in the world in a 2018 F8 talk. ScyllaDB is a distributed database written in C++ that offers API compatibility with Cassandra. Samsung’s benchmark, published by ScyllaDB, reports 10x to 37x better performance than Cassandra on the same 2 TB dataset over two hours. It is a vendor-published benchmark on high-end hardware (the same post says the gap is 1.5x to 3x on smaller machines), so treat the range as a data point and not a promise.

TigerBeetle goes the other way: it does one thing. It stores only accounts and transfers, and an account record is only 128 bytes. Jepsen’s 2025 analysis of versions 0.16.11 through 0.16.30 found seven client and server crashes and only two safety issues, and concluded that as of 0.16.30 TigerBeetle appeared to meet its promise of Strong Serializability (the report was funded by TigerBeetle). Joran Greef describes extending TigerBeetle’s LSM engine with a connector to object storage for the lower levels, keeping hot data close to the CPU and cold data cheap.

Who named all this

The basic terms were already in use: Cassandra shipped size-tiered and leveled strategies, and Dostoevsky compares tiering and leveling. What was missing was agreement on what exactly they meant. In Name that compaction algorithm, Callaghan set out leveled, tiered, tiered+leveled, leveled-N and time-series (the two hybrids were new names), and argued that “LSM tree” should be widened to include more than leveled compaction. He also notes he is not sure these have ever been formally defined, so the names are a vocabulary, not a standard.

What every one of these engines assumes

Every system above shares an assumption so basic that none of them question it: the data lives on a local disk. SSD or hard drive, it is physically attached to the machine running the database.

What if that assumption breaks? What if disk were close to infinite and cheap, but every read cost real money? That is what object storage looks like, and it changes which corner of the triangle hurts: the old trade-offs get re-priced, and a design tuned for local disks has to be rethought. S3 and object storage for databases covers what that looks like.

Sources and further reading

Sources

  1. Production LSM engines (LevelDB, RocksDB, Pebble, Cassandra, ScyllaDB, TigerBeetle, TiKV, YugabyteDB), as cited individually below
  2. Athanassoulis et al., Designing Access Methods: The RUM Conjecture, EDBT 2016
  3. Mark Callaghan, Read, write & space amplification: B-Tree vs LSM, 2015
  4. Dayan, Idreos, Dostoevsky: Better Space-Time Trade-Offs for LSM-Tree Based Key-Value Stores, SIGMOD 2018
  5. Mark Callaghan, Universal Compaction in RocksDB and Me, 2023
  6. Cockroach Labs, Introducing Pebble, 2020
  7. TiKV, RocksDB in TiKV
  8. Unite.AI, Interview with Karthik Ranganathan, co-founder of YugabyteDB
  9. Dostoevsky, Section 2 (tiered compaction), same paper as 4
  10. Apache Cassandra docs, Size Tiered Compaction Strategy
  11. Dostoevsky, Section 2 (Bloom filters), same paper as 4
  12. Arithmetic: 1 billion keys x 10 bits per key = 1.25 GB
  13. Google, LevelDB
  14. Facebook Engineering, Under the Hood: Building and open-sourcing RocksDB, 2013
  15. RocksDB Wiki, Tuning Guide
  16. Cockroach Labs, Pebble announcement, same post as 6
  17. Pebble README, RocksDB features not implemented
  18. Instagram, Cassandra on RocksDB at Instagram, F8 2018
  19. ScyllaDB, Technology
  20. TigerBeetle docs, Account
  21. Joran Greef, A Trillion Transactions, TigerBeetle, 2026: the talk and the written version
  22. Mark Callaghan, Name that compaction algorithm, 2018
  23. ScyllaDB, Compaction Series: Write Amplification in Leveled Compaction, 2018
  24. RocksDB Wiki, RocksDB Bloom Filter

Further reading

← all posts