Skip to content

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

0 of 2 checks done
  1. 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.

    1. 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.
    2. When the memtable fills, it's flushed to disk as an immutable sorted file, an SSTable.
    3. 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.
    4. Compaction merges SSTables in the background, discarding overwritten values, using disk and CPU that traffic also needs.
  2. Check

    In an LSM store, which is usually cheaper: a write or a read?

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

Write at least 60 characters (0 so far). Write it as you would say it in a design review. You will compare it against the points a strong answer makes.

Where you practise it

Further reading

Engineers describing it in systems they run.

  • 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.