Log-structured storage (LSM trees)
Storage engines that turn every write into a sequential append and merge files in the background: very fast writes, at the cost of compaction, tombstones and more expensive reads.
Storage & state
Learn it
Updating data in place, as B-tree databases do, means random disk writes, the slowest thing a disk does. Write-heavy systems (chat history, metrics, event logs) want writes as cheap as an append.
An LSM (log-structured merge) tree never updates in place.
- A write is appended to a commit log (for durability) and inserted into an in-memory sorted table, the memtable. Nothing on disk is modified.
- When the memtable fills, it's flushed to disk as an immutable sorted file, an SSTable.
- A read checks the memtable, then possibly several SSTables, newest first. Bloom filters skip files that can't contain the key, but reads still cost more than writes.
- Compaction merges SSTables in the background, discarding overwritten values, using disk and CPU that traffic also needs.
Check
In an LSM store, which is usually cheaper: a write or a read?- A delete can't erase a value from an immutable file, so it writes a tombstone. Tombstones stay until compaction can remove them safely (after a grace period, so replicas that missed the delete don't resurrect the value). A read across many tombstones must scan them all.
So designs on LSM stores keep partitions bounded and avoid reading through lots of deleted rows.
Think first
You delete a million rows from an LSM store. Does disk usage go down right away?
Quick reference
The same ideas, condensed for revision.
How it goes wrong
- Compaction falls behind
- File counts grow, reads touch more files, and latency climbs.
- Tombstone scans
- A read over a range of deleted rows scans every tombstone and can stall the node.
- Oversized partitions
- Huge partitions make compaction and repair slow and memory-hungry.
- Resurrected deletes
- Tombstones dropped before every replica saw them let deleted data come back.
Instead, consider
- B-tree storage (Postgres, MySQL InnoDB)
- Reads dominate, updates are in place, and you want predictable read latency.
- Append-only object storage
- Data is written once in large batches and rarely read.
In practice
- RocksDB / LevelDB
- Embedded LSM engines inside many databases and services.
- Cassandra / ScyllaDB
- Distributed wide-column stores built on SSTables.
- HBase, Bigtable
- Wide-column stores on LSM storage.
It assumes
- The workload is write-heavy or append-mostly.
- Reads mostly fetch recent data or single partitions.
- Background compaction has spare I/O to run.
Explain it in your own words
Where you practise it
Further reading
Engineers describing it in systems they run.
- How Discord Stores Trillions of Messages
Discord · Bo Ingram · Post, Mar 2023
The same data model six years on: hot partitions, a service layer that merges identical reads, and a migration of the full history to a new database.
- How Discord Stores Billions of Messages
Discord · Stanislav Vishnevskiy · Post, Jan 2017
Choosing a database and a partition key for chat history, and the surprises that followed: tombstones, and an edit racing a delete.
Related concepts
- Durability
What has to have happened before a system may say "saved": which failures the data must survive, and where that guarantee is actually made.
- Append-only logs
Recording changes as an ordered, immutable sequence of facts, from which current state, history and replicas can be derived.
- Partitioning
Splitting data or work by key so each part is handled independently: scaling out, and giving each key a single owner.
- Columnar storage
Storing each column of a table separately, so analytical queries read only the columns they use and compress them well, at the cost of slow single-row lookups and updates.