Sorted Runs

LSM trees, from first principles to production systems

What Is an LSM Tree? The Data Structure Behind Most Modern Databases

ep1lsm treerocksdbcompactionstorage engines

Corrections to the video

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

  • medium Run #2 is flushed as dani, edu, fabi, but at the compaction step it shows cron 115, dani, edu. The article uses the second version throughout.
  • high The WAL is introduced as protection "in case the power goes out". RocksDB's default WAL only guarantees recovery from a process crash. Surviving a power cut needs synced writes.
  • medium "Sequential writes are 10 to 100 times faster, even on a modern SSD" has no primary source, and the paper cited next to it (Mohan et al.) measures file system I/O amplification, not that ratio.
  • low "Back in 1991 a team at DEC developed it": the 1991 version is a UMass Boston technical report by the same four authors. Only Edward Cheng was at DEC.

Every time you write to a database, your data goes through some structure that decides where the bytes land on disk. For a large share of the databases built in the last two decades (RocksDB, Cassandra, CockroachDB, TiKV, ScyllaDB, YugabyteDB and many more) that structure is a Log-Structured Merge tree, or LSM tree.

This article walks through how an LSM tree works, from the first write to compaction, with diagrams you can stare at and links to go deeper.

What happens after db.put()?

You call db.put(key, value) and the call returns. Somewhere between that function and the physical disk, the database has to decide where the value lives and how to find it again. The answer to that question shapes everything else: how fast writes are, how fast reads are, and how much disk you pay for.

The baseline: B-trees and random writes

For decades the default answer was the B-tree. A B-tree keeps keys sorted in fixed-size pages, arranged as a shallow tree. Reads are excellent: a few page hops from the root and you are at the right leaf. Most relational databases still index data this way, and for good reason.

The cost shows up on writes. To update a key in place, the engine has to find the exact leaf page that owns that key, modify it, and write the page back. Consecutive writes for unrelated keys land on unrelated pages, so each one is a random I/O.

A B-tree with three writes for keys zool, bob and kim landing on three different leaf pages far apart from each other
Three writes, three different leaf pages. On disk, those pages are nowhere near each other.

On a spinning disk that means physically moving the head. On an SSD there is no head, but the engine still writes whole pages: changing a 128-byte row can mean writing back an entire 4 KB page (Mark Callaghan’s example).

The filing cabinet

Picture a filing cabinet with ten thousand folders in alphabetical order. Each time someone hands you a document, you find the right drawer, pull the folder, insert the page in order, and push it back. One document is fine. A hundred people handing you documents at once, each for a different drawer, and the queue grows faster than you can file. You are working at full speed and still falling behind, because the walking is the bottleneck, not the filing.

That walking is random I/O. Writing one item after another in a straight line, sequential I/O, is far cheaper. The LSM-tree paper itself puts it at about ten to one on the disks of its day: a page moved as part of a large multi-page block cost roughly a tenth of a random single-page I/O (O’Neil et al., section 3.1). B-trees leave that difference on the table.

The LSM idea: stop filing immediately

What if you did not file each document on arrival? Put it on your desk instead. When the stack is tall enough, sort it once and file the whole stack in a single trip.

That is the whole trick, described by O’Neil, Cheng, Gawlick and O’Neil in the 1996 LSM-Tree paper: defer and batch the work so the disk only ever sees large sequential writes.

LSM write path: put goes first to the write-ahead log on disk, then into the in-memory MemTable kept sorted; when full, the MemTable is flushed to disk as a new sorted run, newest on top
The write path: log it, buffer it, flush it as a sorted run.

MemTable: the desk

Incoming writes go into an in-memory structure kept sorted by key, called the MemTable. In RocksDB it is a skip list by default; other engines use balanced trees or similar. No disk access happens here, so inserting is cheap.

Flush: one sequential burst

When the MemTable reaches its size limit, the engine freezes it, starts a fresh one for new writes, and writes the frozen one to disk as a single sorted, immutable file. Different engines call it a sorted run or an SSTable (Sorted String Table, a name that comes from Bigtable). Because the data is already sorted, the file is written front to back in one pass. That is the sequential write we wanted.

WAL: in case the power goes out

The MemTable lives in memory, so a crash would lose it. Before a write touches the MemTable, the engine appends it to a write-ahead log on disk. Appending is sequential too. On restart, the engine replays the log to rebuild the MemTable. Once a MemTable is flushed, its part of the log can be deleted. In its default configuration RocksDB guarantees recovery from a process crash; surviving a machine crash or power cut needs the write synced to disk, which the caller opts into. See the RocksDB WAL docs for the details.

So the write path is: log, buffer, flush. Nothing on disk is ever modified in place. New data is only ever appended.

A worked example

Three writes arrive: cron → 100, garu → 200, zool → 150. They land in the MemTable in key order. The MemTable fills up and flushes to disk as Run #1.

Three more writes arrive: cron → 115 (an update), dani → 300, edu → 250. Another flush, and now there is Run #2 next to Run #1. Note that cron now exists in both files with different values. Nothing was updated in place; the new value was simply written somewhere newer.

Writes stay fast no matter how many runs pile up. Reads are a different story.

The read problem

To read garu, the engine has to find the most recent value for that key. It checks the places where data can live, newest first:

Reading garu: the empty MemTable misses, Run #2 with cron, dani and edu misses, Run #1 contains garu 200 and the read stops there
Newest first. The first match is the current value, so the search stops there.
  1. The MemTable: maybe it was written recently. Not there.
  2. Run #2, the newest file on disk. Not there.
  3. Run #1. Found: garu → 200.

Each run is sorted, so looking inside one is a binary search, not a scan. The problem is how many runs there are. With two, this is fast. After enough flushes there might be two hundred, and a key that lives in the oldest run, or does not exist at all, means two hundred lookups. Writes are fast because we only append. Reads get slower over time because there is more and more to search through.

Real engines add tricks to skip files quickly (Bloom filters, per-file key ranges, caches), which are worth their own deep dive. But they reduce the cost per file; they do not remove the growing number of files. Something has to clean up.

Compaction: the background janitor

Compaction is a background job that takes two or more sorted runs and merge-sorts them into one. Because each input is already sorted, this is a streaming merge, the same step as the second half of merge sort, and it is sequential on both the read and write side.

When the same key shows up in more than one input, the newest version wins and the older one is dropped. Deletes work the same way: a delete is written as a special marker (a tombstone), and compaction eventually drops both the marker and the values it shadows (LevelDB’s implementation notes describe exactly this).

Compaction merges Run #2 (cron 115, dani, edu) and Run #1 (cron 100, garu, zool) into Run #3 with cron 115, dani, edu, garu, zool; the stale cron 100 is dropped
Two runs in, one run out. The stale cron → 100 disappears.

After compaction, a read for garu checks one file instead of two, and the disk no longer stores the dead cron → 100. That is the full LSM loop: buffer, flush, compact.

What compaction costs

Compaction is not free. To produce Run #3, the engine read both inputs from disk and wrote every surviving byte again. A value you wrote once may get rewritten several times over its lifetime as it is merged into bigger and bigger runs. The RocksDB Tuning Guide calls this write amplification: write 10 MB/s to the database, see 30 MB/s hit the disk, and your write amplification is 3. For a 500 GB database with leveled compaction, the same guide works out about 33, and much of the guide is about keeping that number in check.

That is the tradeoff: faster reads, paid for with extra write work in the background.

The amplification triangle

Every LSM engine is balancing three costs at the same time:

A triangle with corners labeled write amplification, read amplification and space amplification, with a dot for your engine inside
Every engine picks a point inside this triangle.

You cannot minimize all three at once. Push one corner down and at least one of the others goes up. This is not a code quality problem; it is a structural constraint. Mark Callaghan’s Read, write & space amplification: B-Tree vs LSM is the clearest short treatment, and Athanassoulis et al. formalized it as the RUM Conjecture in 2016. The trade-off is what drives the design of every LSM engine.

Where it came from, and where it went

Timeline: 1996 LSM-Tree paper, 2006 Bigtable, 2011 LevelDB, 2013 RocksDB, 2010s YugabyteDB, TiKV and CockroachDB build on RocksDB, 2020 Pebble
Thirty years from paper to the engine under your database.

The four authors first described it in a 1991 UMass Boston technical report, which the 1996 Acta Informatica paper cites; co-author Edward Cheng was at Digital Equipment Corporation. Google’s Bigtable (2006) brought the memtable plus SSTables plus compaction design to production at scale. Sanjay Ghemawat and Jeff Dean then wrote LevelDB, open-sourced in 2011, a small embeddable engine with the same design (its implementation notes are short and worth reading).

Facebook built RocksDB on LevelDB and open-sourced it in 2013, because LevelDB could not keep up with fast flash storage and many-core servers (its single-threaded compaction caused write stalls). RocksDB became the default embedded engine for a generation of distributed databases: TiKV builds on it, engineers who helped build it went on to co-found YugabyteDB on a modified RocksDB, and CockroachDB ran on it from its inception until rewriting the engine in Go as Pebble (the default from release 20.2 in 2020), partly because they could not profile across the Go to C++ boundary.

Same structure, different bets

All of these engines start from the same three steps. They end up in very different places on the triangle because their workloads are different: some bet on fast reads, some on cheap writes, some on predictable latency, some on using less disk. The tuning surface is large enough that the RocksDB team itself writes, in the tuning guide’s closing notes:

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

If the structure is the same, why are there so many engines? That is the subject of why there are so many LSM engines, which takes the three-way trade-off above and shows how each engine picks its corner.

Sources and further reading

Sources

  1. O’Neil, Cheng, Gawlick, O’Neil, The Log-Structured Merge-Tree (LSM-Tree), Acta Informatica, 1996
  2. ByteByteGo, The Secret Sauce Behind NoSQL: LSM Tree
  3. Facebook Engineering, Under the Hood: Building and open-sourcing RocksDB, 2013
  4. Cockroach Labs, Introducing Pebble: A RocksDB-inspired key-value store written in Go
  5. TiKV, RocksDB in TiKV
  6. ScyllaDB, Scaling ScyllaDB storage engine with state-of-art compaction
  7. Andy Pavlo, CMU Intro to Database Systems, #04 Database Storage: Log-Structured Merge Trees & Tuples, Fall 2024
  8. Mohan, Kadekodi, Chidambaram, Analyzing IO Amplification in Linux File Systems
  9. Google, LevelDB
  10. RocksDB Wiki, Write Ahead Log (WAL)
  11. RocksDB Wiki, Tuning Guide
  12. Athanassoulis et al., Designing Access Methods: The RUM Conjecture, EDBT 2016
  13. Unite.AI, Karthik Ranganathan, Co-Founder and Co-CEO of Yugabyte, Interview Series

Further reading

← all posts