When building applications that ingest massive write streams (such as high-frequency telemetry, financial transaction logs, metrics collectors, or event queues), relational databases with standard B-Tree indexes eventually hit an I/O brick wall.
Even on ultra-fast NVMe SSDs, write throughput collapses, disk queues back up, and I/O latency spikes into hundreds of milliseconds.
To solve this exact bottleneck, modern distributed systems like RocksDB, Apache Cassandra, ScyllaDB, ClickHouse, TiKV, and BadgerDB abandoned B-Trees entirely. Instead, they run on a storage architecture called the Log-Structured Merge-Tree (LSM-Tree).
Here is what actually happens inside an LSM-Tree storage engine, why it turns expensive random disk writes into high-speed sequential streaming, and how it balances the brutal trade-offs between read, write, and space amplification.
The Root Problem: Why B-Trees Choke on Writes
To understand why LSM-Trees exist, you first have to see where classical B-Trees fail under heavy write loads.
In a traditional B-Tree (used by PostgreSQL, MySQL InnoDB, and SQLite), data is organized into fixed-size disk pages (typically 4KB, 8KB, or 16KB). When you insert a new record or update an existing row:
- The database traverses the tree in memory to find the exact leaf page containing that key.
- The page is modified in the buffer pool.
- The page is marked dirty.
- A background flush (or checkpoint process) writes the entire 8KB page back to its specific physical offset on disk.
Classical B-Tree (In-Place Updates):
[ Client PUT: key="user_912", val="active" (32 bytes) ]
│
▼
[ Buffer Pool: Locates Page #4189 (8,192 bytes) ]
│ Modifies 32 bytes in memory
▼
[ Disk: Overwrites Page #4189 at Random Offset 0x7FA0200 ]
This model is called In-Place Update (Update-in-Place). It is fast for reads because every key lives in exactly one predictable disk page. But for writes, it creates two major bottlenecks:
- Massive Write Amplification (WA): If you modify a 50-byte record, the engine must write the entire 8KB/16KB page to disk. That is a 160x to 320x write amplification penalty at the storage layer.
- Random I/O Thrashing: Even on modern NVMe drives, random writes force the Flash Translation Layer (FTL) to execute continuous garbage collection and block erases. Flash memory cannot overwrite bytes in place; it must erase entire NAND blocks (often 2MB to 8MB) to write updated 4KB pages.
Under sustained write ingestion of 50,000+ records per second, B-Tree engines spend almost all their I/O budget flushing scattered pages across random disk blocks.
The LSM-Tree Philosophy: "Never Modify Disk in Place"
First introduced by Patrick O'Neil, Edward O'Neil, and Gerhard Weikum in 1996, the LSM-Tree operates on one simple, radical principle:
All writes, updates, and deletes are strictly sequential appends to immutable files.
An LSM-Tree engine never seeks to a random file offset to update existing bytes. It accepts writes in memory, sorts them in batches, and streams them to disk sequentially.
To make this work without making reads impossibly slow, an LSM storage engine relies on four core components:
+──────────────────────────────────────────────────────────────+
│ INCOMING WRITES │
+──────────────────────────────────────────────────────────────+
│ │
▼ ▼
+───────────────+ +────────────────+
│ Write-Ahead │ (Disk) │ Active MemTable│ (RAM: SkipList)
│ Log (WAL) │ Append-Only │ Sorted in-RAM │
+───────────────+ +────────────────+
│
▼ (When Full: ~64MB)
+────────────────+
│ Immutable │
│ MemTable (RAM) │
+────────────────+
│
▼ (Background Flush)
================================================================
DISK STORAGE (Immutable SSTables)
================================================================
Level 0 (L0): [ SSTable 1 ] [ SSTable 2 ] [ SSTable 3 ]
(Keys overlap across files)
│
▼ Compaction (Multi-way Merge Sort)
Level 1 (L1): [ SSTable A ] ─── [ SSTable B ] ─── [ SSTable C ]
(Keys strictly partitioned, non-overlapping)
│
▼ Compaction
Level 2 (L2): [ SST 10 ] ────── [ SST 11 ] ────── [ SST 12 ] ...
Let us break down each component under the hood.
1. The Write-Ahead Log (WAL)
Because incoming writes are initially stored in RAM, a sudden power failure or OS crash would lose all uncommitted data.
To guarantee ACID durability (the "D" in ACID), the engine first appends the raw operation to an append-only disk log called the Write-Ahead Log (WAL) using sequential writes.
Sequential writes to disk do not require random seeks or NAND block reallocations. The OS writes data straight into page cache buffers and streams it out in contiguous chunks. This allows the WAL to achieve near-saturating disk bus bandwidth with sub-millisecond latency.
2. The MemTable (In-Memory Sorted Buffer)
Simultaneously with the WAL append, the key-value pair is inserted into an in-memory data structure called the MemTable.
Unlike a hash map, the MemTable must keep all keys sorted at all times. Most production engines (such as RocksDB and LevelDB) implement the MemTable as a SkipList.
Why a SkipList Instead of a Red-Black Tree?
A Red-Black Tree or AVL Tree requires pointer rebalancing and node rotations during insertions, which locks large portions of the tree and kills multi-threaded concurrency.
A SkipList uses probabilistic multi-level linked lists with lock-free atomic pointer swaps (CAS - Compare-And-Swap), allowing multiple writer and reader threads to operate concurrently without coarse-grained mutex contention:
Level 3: [10] ─────────────────────────────> [90] ──> NULL
Level 2: [10] ─────────────> [45] ─────────> [90] ──> NULL
Level 1: [10] ──> [25] ────> [45] ──> [60] ─> [90] ──> NULL
Level 0: [10] ──> [25] ─[30]─[45] ──> [60] ─> [90] ──> NULL
When the active MemTable reaches a configurable threshold (typically 64MB or 128MB via write_buffer_size), the engine:
- Marks the active MemTable as Immutable (read-only).
- Instantly allocates a fresh, empty MemTable to accept new incoming client writes without pausing traffic.
- Spawns a background thread to flush the immutable MemTable to disk as an SSTable.
3. SSTables (Sorted String Tables)
An SSTable (Sorted String Table) is an immutable, ordered file on disk.
Because the MemTable was already sorted in memory, writing an SSTable to disk is a pure sequential stream (an O(N) linear write). Once written, an SSTable is never modified. It is read-only for its entire lifetime until it is eventually deleted during compaction.
An SSTable file is structured internally into distinct blocks:
+───────────────────────────────────────────────────────────────+
│ Data Block 0: [ "alice": {...}, "bob": {...}, "carol": {...} ]│
├───────────────────────────────────────────────────────────────┤
│ Data Block 1: [ "dave": {...}, "eve": {...}, "frank": {...} ] │
├───────────────────────────────────────────────────────────────┤
│ Data Block 2: [ "grace": {...}, "heidi": {...}, "ivan": {...}│
├───────────────────────────────────────────────────────────────┤
│ Filter Block: [ Bloom Filter Bit Array for all keys in SST ] │
├───────────────────────────────────────────────────────────────┤
│ Index Block: [ "alice" -> Block 0 Offset, │
│ "dave" -> Block 1 Offset, │
│ "grace" -> Block 2 Offset ] │
├───────────────────────────────────────────────────────────────┤
│ Footer: [ Magic Number, Index Offset, Filter Offset ] │
+───────────────────────────────────────────────────────────────+
- Data Blocks: Fixed-size chunks (typically 4KB to 64KB, compressed with ZSTD or Snappy) containing sorted key-value entries.
- Index Block: A directory mapping the start key of each data block to its exact byte offset inside the file.
- Filter Block: In-memory Bloom filter bit array for zero-I/O membership checks.
- Footer: Fixed-size trailer at the end of the file containing pointers to the Index and Filter blocks.
To read a key from an SSTable, the engine loads the Index Block, runs binary search on the keys (O(log N)), identifies the exact Data Block, and decompresses only that single block.
4. Bloom Filters: Preventing Disk Search Explosions
Because an LSM engine can have hundreds of SSTable files on disk, how does it avoid reading 100 files just to find out that a key does not exist?
This is solved by Bloom Filters stored in RAM.
A Bloom filter is a space-efficient probabilistic data structure. For every key added to an SSTable during flush:
- The key is passed through $k$ independent hash functions.
- The corresponding $k$ bit positions in a bit array are set to
1.
Key: "user_402"
Hash1("user_402") % 16 = 3 ──┐
Hash2("user_402") % 16 = 7 ──┼──> Set bits 3, 7, 12 to 1
Hash3("user_402") % 16 = 12 ──┘
Bit Array: [ 0 | 0 | 0 | 1 | 0 | 0 | 0 | 1 | 0 | 0 | 0 | 0 | 1 | 0 | 0 | 0 ]
^ ^ ^
Bit 3 Bit 7 Bit 12
When looking up a key:
- If any of the $k$ bit positions is
0, the key definitely does not exist in that SSTable. The database skips the file entirely with zero disk I/O. - If all bits are
1, the key might exist. The engine reads the SSTable index block.
With just 10 bits per key, a Bloom filter achieves a false positive rate under 1%. That means 99% of negative file lookups are eliminated in RAM before touching the disk.
The Complete Write and Read Paths
Now let us trace how PUT and GET operations traverse the engine.
The Write Path (Latency: Microseconds)
Client PUT(k, v)
│
├──> 1. Append record to Write-Ahead Log (Sequential Disk Append)
│
└──> 2. Insert into Active MemTable SkipList (In-Memory O(log N))
│
└──> 3. Return 200 OK to Client
The write is complete. No pages were searched, no disk blocks were rewritten, and no locks were held on tree nodes.
The Read Path (Hierarchical Traversal)
When a client requests GET(key):
- Search Active MemTable: Check the in-memory SkipList. If found, return the value immediately (most recent write).
- Search Immutable MemTables: Check any in-memory MemTables waiting to be flushed.
- Search Level 0 SSTables: L0 files are direct dumps of MemTables, so their key ranges can overlap. The engine checks Bloom filters from newest to oldest L0 SSTable.
- Search Level 1 to Level N SSTables: In Level 1 and beyond, key ranges within each level are strictly non-overlapping. The engine runs a binary search across level metadata to pick the single candidate SSTable, checks its Bloom filter, and reads the data block.
GET("order_883")
│
[ Active MemTable ] ──── Found? ───> Return Value
│ No
[ Immutable MemTables ] ─ Found? ──> Return Value
│ No
[ Level 0 SSTables ] ─── Bloom Check ──> Match? ──> Read Block
│ No
[ Level 1 SSTables ] ─── Binary Search Files ──> Bloom Check ──> Read Block
│ No
[ Level 2..N SSTables ]
The Tombstone Paradox: Why Deletes Increase Disk Usage
How do you delete a record in a storage engine where disk files are immutable?
You cannot go back and erase bytes from an existing SSTable. Instead, the engine executes deletions by appending a new record with a special flag called a Tombstone:
Client: DELETE("user_104")
Engine: PUT("user_104", DELETED_TOMBSTONE_MARKER)
This leads to the Tombstone Paradox:
In an LSM-Tree, deleting a million rows immediately increases disk usage and consumes memory buffers.
When a read query searches for "user_104", it encounters the tombstone in the newest SSTable and returns None (Key Not Found).
The old values living in older SSTables are not physically removed until a background process called Compaction merges the levels.
Compaction: Merging the Levels and Paying the Debt
If writes keep appending SSTables indefinitely, three problems emerge:
- Disk space fills up with obsolete versions of updated keys and expired tombstones.
- Read performance degrades because queries must search through dozens of SSTables.
- Level 0 file counts balloon.
To solve this, LSM engines run continuous background Compaction. Compaction reads multiple sorted SSTable files, merges them using an $O(N)$ Multi-Way Merge Sort, discards dead records, and outputs new, cleanly sorted, non-overlapping SSTable files.
SSTable 1: [ "apple": $1.00 (v1), "banana": $0.50 (v1) ]
SSTable 2: [ "apple": $1.20 (v2), "cherry": $2.00 (v1) ]
SSTable 3: [ "banana": TOMBSTONE, "date": $3.00 (v1) ]
│
▼ Multi-way Merge Sort
Output SST: [ "apple": $1.20 (v2), "cherry": $2.00 (v1), "date": $3.00 (v1) ]
(Discarded old "apple" v1 and purged deleted "banana")
Production engines primarily use two compaction strategies:
1. Size-Tiered Compaction Strategy (STCS)
Common in Cassandra and ScyllaDB for append-only workloads.
- When $N$ SSTables of similar size accumulate in a tier, the engine merges them into a single larger SSTable in the next tier.
- Pros: Fast compaction with low CPU and write overhead.
- Cons: High Space Amplification. Merging huge files requires up to 50% free disk space as temporary headroom.
2. Leveled Compaction Strategy (LCS)
Default in RocksDB, LevelDB, and TiKV.
- Disk is split into levels ($L_0, L_1, L_2, \dots, L_k$), where each level is typically 10x larger than the previous level (e.g., $L_1 = 10\text{MB}, L_2 = 100\text{MB}, L_3 = 1\text{GB}, L_4 = 10\text{GB}$).
- Within $L_1$ and above, key ranges are strictly partitioned across files so no two SSTables share overlapping keys.
- When $L_1$ exceeds its capacity, an $L_1$ file is merged with all overlapping files in $L_2$.
- Pros: Minimal space amplification (typically $< 1.1\text{x}$) and predictable read latencies.
- Cons: High write amplification because data is rewritten across multiple levels as it ages.
The RUM Trade-Off: LSM-Trees vs B-Trees
In storage engine design, the RUM Conjecture states that you cannot optimize all three metrics simultaneously:
- Read Amplification (RA): Number of bytes read from disk per logical byte requested.
- Write Amplification (WA): Number of bytes written to disk per logical byte inserted.
- Space Amplification (SA): Ratio of physical disk space used vs logical uncompressed data.
Metric B-Tree (InnoDB, Postgres) LSM-Tree (RocksDB, Cassandra)
────────────────────────────────────────────────────────────────────────────────
Write Ingestion Random I/O (Bottleneck) Sequential Append (Fast)
Write Amplification Very High (50x - 100x+) Moderate to High (10x - 30x)
Point Read Latency Predictable (1-3 page reads) Depends on Bloom filter / cache
Range Scan Speed Fast (Traverse leaf links) Requires Multi-Way Heap Merge
Space Overhead High (Page internal fragmentation) Low (Dense sequential packing)
NAND Flash Wear High (Random block rewrites) Low (Large sequential streams)
Practical Architectural Takeaways
Understanding the internal mechanics of LSM-Trees makes choosing and tuning databases straightforward:
- Pick B-Trees (PostgreSQL, MySQL, SQLite) when your workload is read-dominant ($>80\%$ reads), requires complex SQL queries with arbitrary secondary indexes, or demands low-latency random point lookups without background compaction jitter.
- Pick LSM-Trees (RocksDB, Cassandra, ClickHouse, TiKV) when your workload is write-intensive, deals with time-series or event ingestion, requires maximum storage compression density, or runs on distributed architectures where sequential appends simplify replication logs.
- Beware Tombstone Bloat: If you run bulk delete operations on an LSM engine (like Cassandra or RocksDB), point reads scanning those key ranges will experience sudden latency spikes until compaction merges the tombstones away.
-
Tune MemTable Buffers: In high-throughput ingest services, increasing
write_buffer_sizeand the number of background flush threads prevents write stalls when ingestion outpaces the disk flush rate.
By turning random in-place updates into structured sequential streams, LSM-Trees provide the high-throughput backbone that powers modern distributed infrastructure.
Top comments (0)