LSM Trees and SSTables - Complete Deep Dive
Prerequisites: Database Indexing, Write-Ahead Log, Bloom Filters Used in: Cassandra, RocksDB, LevelDB, HBase, ScyllaDB, and the storage engine behind most time-series databases
What is an LSM Tree?
A Log-Structured Merge tree is a storage engine that never updates data in place. Writes land in a sorted in-memory buffer, that buffer is flushed to disk as one immutable sorted file, and background jobs merge those files together. Every disk write is a sequential append of a whole file. Nothing is ever patched.
Real-world analogy: A busy kitchen does not re-alphabetise the pantry after every delivery. New stock goes on a staging shelf, sorted. When the shelf fills, someone boxes it up and labels the box with its contents. Finding an ingredient means checking the staging shelf, then the newest box, then older boxes. Once a week, someone merges several small boxes into one big tidy box and throws the old ones out. That weekly merge is compaction, and it is where all the real work happens.
B-tree write: find page 4,712 โ modify 100 bytes โ write 8KB page back (random)
LSM write: append 100 bytes to log โ insert into RAM โ ack (sequential)
Why B-Trees Struggle With Writes
A B-tree index is excellent at reads and genuinely bad at small random writes. Three reasons compound.
In-place page updates. A B-tree stores rows inside fixed-size pages, 8KB in Postgres and 16KB in InnoDB. Updating a 100-byte row means reading that page, changing the bytes, and writing the whole page back. You asked for 100 bytes of work and the device did 8,192 bytes of work.
Random placement. Rows arrive keyed by user ID, order ID, or a UUID. Those keys scatter across the whole tree, so consecutive writes touch unrelated pages. On a spinning disk each one costs a seek, roughly 5-10ms, which caps you at low hundreds of writes per second per drive. On flash there is no seek, but the drive cannot overwrite a 8KB page in isolation either - the flash translation layer has to shuffle data to free an erase block, which is often 256KB to several MB. The amplification moved, it did not disappear.
The WAL writes the page too. Postgres has full_page_writes on by default, so the first modification of a page after a checkpoint copies that entire page into the WAL to guard against torn writes. Your 100-byte update now costs an 8KB page image in the WAL plus the 8KB data page at checkpoint time. Add a page split when a leaf is full and you write two pages and touch the parent.
So the honest number for a small random update on a B-tree is tens of kilobytes of physical write per hundred bytes of logical write, paid synchronously, on the critical path.
๐ก Write amplification means bytes the device actually wrote divided by bytes your application asked to write. Everything below is an argument about where to pay it.
The Write Path
flowchart LR
W["Write request"]:::client
WAL[("Commit log<br/>append only")]:::data
MEM["Active memtable<br/>skip list in RAM"]:::service
IMM["Frozen memtable<br/>awaiting flush"]:::service
SS0[("SSTable L0<br/>immutable sorted file")]:::data
SS1[("SSTable L1<br/>non overlapping run")]:::data
W -->|"1. Append and fsync"| WAL
W -->|"2. Insert in sort order"| MEM
MEM -->|"3. Freeze when full"| IMM
IMM -->|"4. Sequential flush"| SS0
SS0 -->|"5. Background compaction"| SS1
classDef client fill:#4c3a5e,stroke:#818cf8,color:#e2e8f0
classDef edge fill:#1e3a5f,stroke:#60a5fa,color:#e2e8f0
classDef service fill:#1a3a2a,stroke:#4ade80,color:#e2e8f0
classDef async fill:#3b1f5e,stroke:#c084fc,color:#e2e8f0
classDef data fill:#3b3520,stroke:#fbbf24,color:#e2e8f0
| Color | Meaning |
|---|---|
| ๐ฃ Purple | Clients |
| ๐ข Green | Services |
| ๐ก Yellow | Data stores |
| ๐ต Blue | Edge / CDN |
The memtable is a sorted in-memory structure, usually a skip list. It needs three things: ordered iteration so it can be written out as a sorted file and scanned by range, cheap concurrent inserts, and no rebalancing storms. A skip list gives all three and is what RocksDB uses by default. Red-black trees also appear. A plain hash map does not work - you would lose the ordering that makes the flush a single sequential pass.
The commit log is the write-ahead log, and it is the only reason this is safe. The memtable is volatile. If the process dies, everything buffered in RAM is gone, so the write is appended to the commit log and fsynced before the client gets an acknowledgement. On restart, the engine replays the log segments belonging to memtables that were never flushed. Once a memtable becomes an SSTable on disk, its log segment is dead weight and gets dropped. Cassandra calls this the commit log, RocksDB calls it the WAL, and they do the same job.
The flush happens when the active memtable crosses a size threshold, 64MB being the common default. It is frozen, a fresh memtable takes over immediately so writes never block on the flush, and the frozen one is streamed to disk in key order as a single file. One sequential write of 64MB, which any modern device does at full bandwidth.
SSTable Structure
An SSTable - Sorted String Table - is an immutable file holding key-value pairs in key order. The layout is the same idea everywhere, with naming differences:
SSTable on disk
โโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโ
โ Data block 0 keys aaa..amz compressed 4-64KB โ
โ Data block 1 keys ana..bqz compressed โ
โ Data block 2 keys bra..dzz compressed โ
โ ... โ
โโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโค
โ Sparse index aaa โ off 0 โ
โ ana โ off 4096 โ
โ bra โ off 9312 one per block โ
โโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโค
โ Bloom filter ~10 bits per key, about 1 pct FPR โ
โโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโค
โ Footer min key, max key, count, checksums โ
โโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโ
The index is sparse on purpose: one entry per block, not per key. At 64KB blocks and 100-byte rows that is roughly one index entry per 650 keys, so the index for a 1GB file is a couple of megabytes and lives comfortably in memory. Finding a key inside the file is a binary search over the sparse index to identify the one block that could hold it, one read of that block, decompress, then a binary search inside the block. One disk read per file, not one per key.
Immutability buys several things at once. Blocks are packed full and compressed hard because nothing will ever be inserted into them, so LZ4 or Zstd on a fully packed block beats what you can do to a B-tree page that must keep free space for future in-place updates. Readers need no locks, because a file that cannot change cannot be read inconsistently. And backups are a file copy.
The Read Path
flowchart LR
R["Read key K"]:::client
MEM["Active memtable"]:::service
IMM["Frozen memtables"]:::service
BF["Bloom filter probe<br/>one per SSTable"]:::service
SS[("SSTables<br/>newest to oldest")]:::data
MISS["Return not found"]:::async
R -->|"1. Check RAM first"| MEM
MEM -->|"2. Miss"| IMM
IMM -->|"3. Miss"| BF
BF -->|"4. Maybe present"| SS
BF -->|"5. Definitely absent"| MISS
classDef client fill:#4c3a5e,stroke:#818cf8,color:#e2e8f0
classDef edge fill:#1e3a5f,stroke:#60a5fa,color:#e2e8f0
classDef service fill:#1a3a2a,stroke:#4ade80,color:#e2e8f0
classDef async fill:#3b1f5e,stroke:#c084fc,color:#e2e8f0
classDef data fill:#3b3520,stroke:#fbbf24,color:#e2e8f0
Order matters absolutely: newest source first, and the first hit wins. There is no โlatest versionโ field to consult, because recency is encoded in position. A value written an hour ago sits in an older file and is shadowed by anything newer. Stop at the first match or you will return stale data.
Without help, a read would have to touch every file. That is what the Bloom filter per SSTable prevents - it answers โis this key definitely not in this file?โ with no disk I/O at all. At the usual 10 bits per key you get about a 1% false positive rate, so a file that cannot contain the key is skipped 99 times out of 100.
๐ก A Bloom filter never says โdefinitely yesโ. A โmaybeโ costs you one wasted block read; a โnoโ is always true.
Two numbers worth holding on to. A point read for a key that exists, in a well-compacted levelled store with a warm block cache, is typically one disk read. A point read for a key that does not exist anywhere costs roughly zero disk reads, because every Bloom filter says no - except for the few percent of the time one lies, and you pay a block read to find nothing.
Range scans get none of this. A Bloom filter answers point membership and nothing else, so WHERE k BETWEEN x AND y has to open an iterator on the memtable and on every overlapping file, then merge-sort the streams while discarding shadowed versions. Scan cost is proportional to how many overlapping runs exist, which is exactly what compaction controls.
Compaction
Compaction is the whole topic. Everything above is a buffering trick; compaction is what makes it a database.
Why it is not optional:
- Reads degrade as files accumulate. Every flush adds a file. More files means more Bloom probes, more candidate blocks, and more streams in a range-scan merge.
- Overwritten keys occupy space forever. Write a key ten times and ten copies exist in ten files. Nine are garbage that nothing will ever read.
- Deleted keys do not free space. A delete writes a marker, not a hole. Reclaiming the bytes requires merging the marker with every older copy.
- Scan latency tracks overlap, not data size. A thousand small overlapping files make a range scan a thousand-way merge.
The mechanism is boring and that is the point: read N SSTables, merge them by key, keep only the newest version of each key, drop expired markers, write new SSTables, delete the inputs. Sequential in, sequential out, no file ever modified. The only hard question is which files to merge, and there are two serious answers.
Size-Tiered Compaction
Group files into tiers by size. When about four similarly sized files accumulate in a tier, merge them into one file that lands in the next tier up.
Tier 0 [64MB] [64MB] [64MB] [64MB] โ merge โ one 256MB file
Tier 1 [256MB] [256MB] [256MB] [256MB] โ merge โ one 1GB file
Tier 2 [1GB] [1GB] ...
Key "user:42" may live in any file in any tier. Ranges overlap freely.
Each byte is rewritten only when its tier fills, so a byte gets rewritten a handful of times over its life. Cheap writes. The cost is that nothing constrains key ranges - a single key can appear in a file in every tier, so a point read may probe them all, and a range scan merges across all of them. Merging a tier also needs free disk space equal to the size of the files being merged, which is why operators running size-tiered Cassandra keep something like 50% headroom and still get surprised.
This is Cassandraโs historical default and roughly what HBase does. RocksDB calls its equivalent universal compaction.
Levelled Compaction
Enforce an invariant: within each level, SSTables have non-overlapping key ranges and together form one sorted run covering the keyspace. Each level is about 10x the size of the one above.
L0 [overlapping files flushed straight from memtables] ~4 files
L1 [a..f][g..m][n..s][t..z] one sorted run ~256MB
L2 [a..b][c..d][e..f] ... [y..z] one sorted run ~2.5GB
L3 ... one sorted run ~25GB
Key "user:42" appears in at most ONE file per level below L0.
Compaction picks a file in Ln, finds the files in Ln+1 whose ranges overlap it - typically about ten of them, given the 10x size ratio - merges, and writes the result back into Ln+1. Point reads get very good: at most one file per level can hold the key, Bloom filters eliminate almost all of those, and the answer is usually one disk read. Space overhead is around 10%, because the bottom level holds roughly 90% of the data and the levels above it hold at most one extra copy.
The bill arrives on writes. Merging one file against ten overlapping files rewrites eleven filesโ worth of bytes to move one fileโs worth of data down a level, and the byte repeats that at every level. Total write amplification in the 10-30x range is normal. This is LevelDB and RocksDBโs default, and Cassandraโs LCS.
| ย | Size-Tiered | Levelled |
|---|---|---|
| Key ranges | Overlap freely across all tiers | Non-overlapping within each level |
| Write amplification | Low, roughly 4-10x | High, roughly 10-30x |
| Read amplification | High, may probe many files | Low, at most one file per level |
| Space amplification | ~2x, needs large free headroom | ~1.1x |
| Range scans | Merge across every tier | Merge across levels only |
| Best for | Write-heavy, append-mostly, immutable records | Read-heavy, update-heavy, space-constrained |
| Defaults in | Cassandra STCS, HBase, RocksDB universal | RocksDB, LevelDB, Cassandra LCS |
There is a third option worth naming because time-series workloads hit it constantly: time-window compaction, where files are bucketed by time and only compacted within their bucket. Data written together expires together, so a whole file can be deleted rather than merged row by row. If your rows have a TTL and you never update them, this is the right answer and the other two are not.
The Three Amplifications
Pick any two. The third gets worse. This is the shape of the RUM conjecture - read, update, and memory overheads cannot all be minimised at once - and it is the most useful thing to say in an interview about storage engines.
| Amplification | Definition | What makes it bad |
|---|---|---|
| Write | Bytes written to device รท bytes written by app | Compaction rewriting the same data at every level |
| Read | Disk reads per logical read | Many overlapping files, weak or missing Bloom filters |
| Space | Bytes on disk รท bytes of live data | Shadowed versions and undeleted markers awaiting compaction |
Where the real configurations land:
- Levelled LSM - low read, low space, high write. You chose to burn disk bandwidth to keep reads and footprint tight.
- Size-tiered LSM - low write, high read, high space. You chose cheap ingestion and will pay on every lookup and in disk purchase orders.
- B-tree - low read, moderate space, high write for small random updates. Read amplification is about one page per level and most levels are cached, space sits around 1.3x from the fill factor, and the writes are expensive and synchronous.
The tuning knob that moves all three at once is compaction aggressiveness. Compact harder and reads and space improve while write amplification and device contention get worse. There is no setting that wins everywhere, so the question is always which one your workload can afford to sacrifice.
Tombstones and the Delete Problem
You cannot delete from an immutable file. So a delete writes a new record - a tombstone - carrying the key, a timestamp, and nothing else. A read that encounters a tombstone before any value returns not-found, because first hit wins.
Two consequences that surprise people:
A delete makes the database bigger. You added bytes. The old value is still sitting in an older file. Nothing is reclaimed until a compaction merges the tombstone with every older copy of that key, and if those copies live in the bottom level, that will not happen soon.
The tombstone cannot be dropped as soon as it is merged. In a replicated store this is a correctness issue, not an optimisation. If a replica was down when the delete was issued, it still holds the old value. Drop the tombstone too early and the next anti-entropy repair sees one node with a value and one with nothing, concludes the value is newer, and resurrects deleted data. Cassandraโs gc_grace_seconds defaults to 864000 - ten days - which is how long a tombstone must survive so that repair has every chance to carry the delete everywhere first.
๐ก Deleted data coming back from the dead after a repair is the single most common Cassandra correctness bug, and it is always someone lowering gc_grace_seconds.
The classic Cassandra failure mode
Model a work queue as a table. Insert rows, process them, delete them. The partition now holds a handful of live rows and a very large number of tombstones. Every read of that partition is a range scan, and a range scan gets no help from Bloom filters, so the engine walks tombstone after tombstone, materialising and discarding each one, to find the few live rows underneath.
Latency climbs, heap pressure climbs, and it gets worse every day because deletes keep arriving faster than compaction can merge them away. Cassandra ships guard rails that tell you exactly what is happening: tombstone_warn_threshold logs a warning at 1,000 tombstones scanned in a single read, and tombstone_failure_threshold aborts the query outright at 100,000. Those thresholds are not the bug. The bug is the data model.
The fix is to stop asking an LSM store to delete things. Give rows a TTL and use time-window compaction so entire files expire and are dropped as units. If you need queue semantics, use a queue - see message queues.
Comparison Table
| Dimension | B-Tree | LSM Tree |
|---|---|---|
| Write path | Locate leaf, modify in place, write page back. Random I/O, synchronous | Append to log, insert into memtable, ack. Sequential I/O, flush later |
| Read path | One page read per level, 3-4 levels, mostly cached. Predictable | Memtable, then one Bloom probe and one block read per level. Usually one disk read |
| Write amplification | Page-granularity. A 100-byte update can write 8KB plus a full-page WAL image | Compaction rewrites data once per level. 10-30x on levelled, less on size-tiered |
| Space amplification | ~1.3x, stable, from fill factor and fragmentation | 1.1x levelled, 2x or worse size-tiered. Spikes during compaction |
| Range scans | Leaves are linked in key order. One sequential pass | Merge iterator across every overlapping run. No Bloom filter help |
| Concurrency | Page latches, writer contention on hot pages, parent locked during splits | Concurrent memtable inserts, and immutable files mean readers take no locks |
| Compression | Poor. Pages keep free space for in-place updates, fill factor ~70% | Good. Blocks are packed once and compressed hard |
| Latency profile | Boring and flat, which is a feature | Good median, tail hostage to compaction scheduling |
Common Interview Questions
Q: โIf an LSM writes more total bytes than a B-tree, how is it faster for writes?โ A: Because of when and how the bytes are written. The acknowledged path is one sequential log append plus an in-memory insert, with no seek and no read-modify-write. Compactionโs bytes are paid later, in the background, in large sequential batches that a device handles at full streaming bandwidth. A B-tree pays random page I/O synchronously, on the request. Total bytes is the wrong metric - sequential bandwidth is often 100x random IOPS on a spinning disk and still several times better on flash.
Q: โA read could have to check every SSTable. Why doesnโt it?โ A: Three filters, applied in order. Each fileโs footer carries its min and max key, so a key outside that range is skipped immediately. Each file has a Bloom filter, which rules out roughly 99% of the files that pass the range check. And under levelled compaction, at most one file per level can contain any given key, so the candidate set is bounded by the level count rather than the file count. Net effect is usually one disk read.
Q: โWhat happens when compaction cannot keep up with writes?โ A: The engine applies backpressure, and you want it to. RocksDB slows writes when L0 file count passes a soft threshold and stalls them completely at a hard one, because the alternative is unbounded read amplification and eventually a full disk. The fix depends on the cause: more compaction threads if you have spare I/O, a compaction rate limiter if compaction is starving reads, a different compaction strategy if the workload shape is wrong, or more nodes if the box is genuinely out of write bandwidth. Watch pending compaction bytes and L0 file count - both move well before latency does.
Q: โWhy does deleting rows make the database use more disk?โ A: Immutable files cannot be edited, so a delete appends a tombstone rather than removing anything. Live data, old versions, and the tombstone all coexist until a compaction merges them, and in a replicated store the tombstone must additionally be retained for the grace period so repair can propagate the delete without resurrecting the old value. Delete-heavy workloads on an LSM grow before they shrink.
Q: โWhen would you pick a B-tree over an LSM tree?โ A: Read-heavy with a tight latency SLO, especially at P99. Workloads doing read-modify-write on single keys, where an LSM has to run the full read path then append. Anything needing mature multi-key ACID transactions, which Postgres and InnoDB give you and a raw LSM does not. And small datasets, where the three-layer lookup and a background compactor buy nothing.
Q: โHow do Bloom filters help a range scan?โ A: They do not, at all. A Bloom filter answers point membership only - it cannot tell you whether a file holds any key in a range. Range scans must open an iterator on every overlapping run and merge. This is why an update-heavy workload that is mostly range scans is the worst case for size-tiered compaction, and often the signal to use levelled or to reach for a B-tree instead.
When NOT to Use an LSM Tree
- Read-heavy with strict latency tails. Your P99 is hostage to the compaction scheduler. A compaction burst saturates disk bandwidth and read latency doubles while it runs. B-tree read latency is flat and dull, which is what you want behind a 10ms budget.
- In-place updates and read-modify-write on single keys. Incrementing a counter means a full read path plus a new append, and the old value sticks around. Cassandra counters have a long history of being awkward for exactly this reason.
- Strong single-key read latency guarantees. You can make the median excellent. Guaranteeing the tail requires controlling compaction, which means over-provisioning I/O.
- Delete-heavy or queue-shaped data. Tombstones. See above.
- Multi-key ACID transactions out of the box. A bolt-on transaction layer over an LSM exists in RocksDB and FoundationDB, but you are building, not configuring.
- Datasets that fit in RAM. Use a hash map or Redis. Compaction is machinery for data that does not fit.
- Mostly-range-scan workloads over frequently overwritten rows. Every scan merges shadowed versions it will throw away.
| โ Back to Fundamentals | Next: Bloom Filters โ |
Related Concepts
- Write-Ahead Log โ โ the commit log that makes the volatile memtable safe to acknowledge from
- Bloom Filters โ โ the one structure that stops a read from touching every SSTable
- Database Indexing โ โ B-tree internals, which this page exists to contrast against
- Merkle Trees โ โ anti-entropy repair, and why tombstones must outlive the grace period
- Query Complexity โ โ reasoning about scan cost when it depends on file layout
Where This Shows Up
This concept is load-bearing in these designs - each link goes straight to the design that leans on it:
- Design a Distributed Key-Value Store - RocksDB per node, with the memtable, SSTable levels, and Bloom-filtered read path spelled out
- Design a Metrics Monitoring System - 1M writes per second of time-series points, the workload LSM storage exists for
- Design Instagram - Cassandra absorbing fan-out feed writes that would melt a B-tree
- Design Twitter Feed - append-only tweet storage partitioned by user, read as a single-partition head scan
- Design Google Docs - an append-only operation log with periodic snapshot compaction
| All Concepts | All HLD Designs |