What Is an LSM Tree? The Data Structure Behind Most Modern Databases
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.
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.
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:
- The MemTable: maybe it was written recently. Not there.
- Run #2, the newest file on disk. Not there.
- 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).
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:
- Write amplification: bytes written to disk per byte the application wrote. Compact more aggressively and this goes up.
- Read amplification: how many places a lookup has to check. Let runs accumulate before compacting and this goes up.
- Space amplification: disk used per byte of live data. Stale versions waiting for compaction, plus temporary copies while a compaction is running, push this up. A compaction that rewrites everything holds input and output on disk at once, so a 1 GB dataset can briefly need 2 GB (ScyllaDB).
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
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
- O’Neil, Cheng, Gawlick, O’Neil, The Log-Structured Merge-Tree (LSM-Tree), Acta Informatica, 1996
- ByteByteGo, The Secret Sauce Behind NoSQL: LSM Tree
- Facebook Engineering, Under the Hood: Building and open-sourcing RocksDB, 2013
- Cockroach Labs, Introducing Pebble: A RocksDB-inspired key-value store written in Go
- TiKV, RocksDB in TiKV
- ScyllaDB, Scaling ScyllaDB storage engine with state-of-art compaction
- Andy Pavlo, CMU Intro to Database Systems, #04 Database Storage: Log-Structured Merge Trees & Tuples, Fall 2024
- Mohan, Kadekodi, Chidambaram, Analyzing IO Amplification in Linux File Systems
- Google, LevelDB
- RocksDB Wiki, Write Ahead Log (WAL)
- RocksDB Wiki, Tuning Guide
- Athanassoulis et al., Designing Access Methods: The RUM Conjecture, EDBT 2016
- Unite.AI, Karthik Ranganathan, Co-Founder and Co-CEO of Yugabyte, Interview Series
Further reading
- Chang et al., Bigtable: A Distributed Storage System for Structured Data, OSDI 2006
- LevelDB, Implementation notes
- RocksDB Wiki, MemTable
- Mark Callaghan, Read, write & space amplification: B-Tree vs LSM, 2015
- Dong et al., RocksDB: Evolution of Development Priorities in a Key-value Store Serving Large-scale Applications, ACM TOS 2021
- Cockroach Labs, Adventures in Performance Debugging, 2016
- YugabyteDB, DocDB performance enhancements to RocksDB
- ScyllaDB, Compaction Series: Space Amplification
- Martin Kleppmann, Designing Data-Intensive Applications, chapter 3