0%
Data-Intensive Applications
Distributed Data
Encoding and Evolution
Batch Processing
Stream Processing
Data Quality and Governance
Operational Patterns
Storage Engines and Data Structures
Every database must answer two questions: how do I store data when a write arrives, and how do I find it again when a read arrives? The simplest possible answer is an append-only log. Every write appends a record to the end of a file. No seeking, no overwriting, no index maintenance. Writes are as fast as sequential disk I/O allows, which on modern SSDs is hundreds of megabytes per second.
The problem is reads. To find a key, you scan the entire file from end to beginning, looking for the most recent entry. A 10 GB log with 100 million entries means reading billions of bytes to answer a single point query. This is why nobody uses a raw append-only log in production. But the insight behind it — sequential writes are fast, random writes are slow — is the foundation of every log-structured storage engine.
You can improve the raw log with a hash index: maintain an in-memory hash map that maps every key to a byte offset in the log file. Writes append to the log and update the hash map. Reads look up the offset in the hash map and do a single seek. Bitcask (the default storage engine in Riak) uses exactly this approach. It works well when the number of distinct keys fits in memory, but it breaks down when the keyspace is larger than RAM — you cannot keep a hash map for a billion keys in memory. It also cannot do efficient range queries because the keys are not sorted. These limitations motivate the next step in the evolution: sorted segments.
Even with the hash index, the log file grows without bound as updates overwrite previous values. The solution is segment compaction: break the log into segments of a fixed size. When a segment fills up, start a new one. In the background, merge old segments by keeping only the latest value for each key and discarding earlier entries. After merging, the old segments are deleted. Each segment has its own hash index. A read checks the most recent segment's hash map first, then the next most recent, and so on. This is a simple but effective approach: writes stay sequential, reads check a small number of hash maps, and compaction reclaims space from deleted and overwritten entries.
SSTables and Memtables
An SSTable (Sorted String Table) fixes the read problem by keeping entries sorted by key within each segment. Instead of appending to one massive file, writes go to an in-memory balanced tree called a memtable (typically a red-black tree or AVL tree). When the memtable reaches a size threshold (usually 2-4 MB), it is flushed to disk as an immutable SSTable. Because the memtable is sorted, the resulting SSTable is sorted too, and writing it is a single sequential operation.
There is one critical durability issue: if the database crashes, the memtable contents (which are only in RAM) are lost. To prevent data loss, every write is also appended to a WAL (write-ahead log) on disk before being inserted into the memtable. This WAL is a simple sequential file — not sorted, not indexed, just raw writes in arrival order. On crash recovery, the engine replays the WAL to reconstruct the memtable. Once the memtable flushes to an SSTable, the corresponding WAL entries are no longer needed and can be discarded.
To read a key, the engine checks the memtable first (it holds the most recent writes). If the key is not there, it checks SSTables on disk from newest to oldest. Within each SSTable, binary search on a sparse in-memory index locates the right block, and then a short scan within that block finds the key. This is fast, but the number of SSTables grows over time, which brings us to compaction.
The sparse index is what keeps that read cheap. It holds one offset per block, not one per key, so an index for a 256 MB SSTable fits comfortably in memory:
The trade-off sits in the block size. Larger blocks mean a smaller index and fewer bytes of RAM per SSTable, but more bytes read and scanned per lookup. Compression pushes the same lever: the block is the unit that gets compressed, so a bigger block compresses better and costs more to decompress for a single key.
LSM-Trees and Compaction
An LSM-tree (Log-Structured Merge-tree) is the full architecture: memtable plus multiple SSTable levels with a background compaction process. Compaction reads two or more SSTables, merges them like merge-sort, discards deleted entries (tombstones) and overwritten values, and writes a single new SSTable. This keeps the number of on-disk segments manageable and reclaims space from obsolete data.
There are two major compaction strategies:
Size-tiered compaction groups SSTables of similar size and merges them when enough accumulate. Newer, smaller SSTables get merged into progressively larger ones. This is write-optimized because compaction happens less frequently, but it can temporarily use 2x disk space during a merge and may leave duplicate keys across levels until they are compacted.
Leveled compaction (used by LevelDB and RocksDB) partitions SSTables into levels with size limits. Level 0 holds recently flushed SSTables. When a level fills up, its SSTables are merged into the next level. Each level above L0 has non-overlapping key ranges, which means reads check at most one SSTable per level. This is read-optimized but causes more write amplification because data is rewritten every time it moves to a deeper level.
To make this concrete: imagine a leveled LSM-tree with 4 levels. A key written to the memtable is flushed to Level 0. When Level 0 fills up, its SSTables are merged into Level 1. When Level 1 fills, its data merges into Level 2, and so on. Each merge rewrites the data to a new SSTable. By the time a key reaches Level 3, it has been written to disk 4 times — once per level — even though the application only wrote it once. With a typical amplification factor of 10-30x, a workload writing 100 MB/s of user data may generate 1-3 GB/s of actual disk writes internally.
For leveled compaction the factor is not a mystery number. With a size ratio between levels and levels, each byte is rewritten roughly once per level and once per merge round within a level, giving
so a default RocksDB configuration with and four levels lands near 40x in the worst case and closer to 10-20x in practice, because not every level is full. Size-tiered compaction trades this for space: its write amplification is closer to , but it can hold two copies of the merged data at once, which is why size-tiered deployments are provisioned with 50% free disk rather than 20%.
Write amplification is the hidden cost of LSM-trees. A single user write may be rewritten 10-30 times as it is compacted through levels. This consumes disk I/O bandwidth that could otherwise serve reads. When evaluating LSM-based engines (RocksDB, Cassandra, LevelDB), check the write amplification factor for your workload — it determines whether the disk can keep up.
Bloom Filters
Even with sorted SSTables and sparse indexes, an LSM-tree must check multiple SSTables for a key that does not exist. A query for a missing key scans the memtable and every SSTable level before returning "not found." This is read amplification — the cost of confirming absence. For workloads with frequent lookups for nonexistent keys (e.g., checking whether a username is taken), read amplification can dominate query latency.
Bloom filters solve this. A Bloom filter is a compact bit array with multiple hash functions. When a key is written to an SSTable, it is hashed through each function and the corresponding bits are set to 1. To check membership, hash the query key and check the bits. If any bit is 0, the key is definitely not in this SSTable — skip the disk read entirely. If all bits are 1, the key might be there (false positives are possible but typically below 1%), so proceed to the actual SSTable lookup.
Every production LSM-tree uses Bloom filters. RocksDB, Cassandra, HBase, and LevelDB all attach a Bloom filter to each SSTable. The space cost is roughly 10 bits per key (about 1.2 bytes), a tiny overhead that eliminates most unnecessary disk reads.
That 10 bits is not folklore, it is the solution to a sizing equation. For bits, keys and hash functions, the false positive rate is
which is minimized at , and inverting it gives the bits you have to buy for a target rate:
Two things follow from the shape of that result. The cost is per key and independent of key length, so a Bloom filter over 40-byte keys costs the same as one over 4-byte keys. And each additional factor of ten in accuracy costs a flat 4.8 bits per key, which is why engines stop around 10 to 14 bits: below that the filter stops skipping reads, and above it the filter stops fitting in the memory you wanted to spend on block cache.
Real-World LSM Engines
Understanding where LSM-trees appear in practice helps connect the theory to systems you will encounter:
RocksDB (Facebook/Meta) is an embedded key-value store used as the storage engine inside MySQL (MyRocks), CockroachDB, and TiDB. It uses leveled compaction by default and provides fine-grained tuning for write amplification, compression, and bloom filter configuration.
Cassandra (Apache) uses LSM-trees with size-tiered compaction by default, optimized for high write throughput across distributed nodes. Each node independently manages its own LSM-tree, and compaction runs locally without cross-node coordination.
LevelDB (Google) is the original implementation that popularized leveled compaction. It is a library, not a standalone database — you embed it in your application. Chrome uses LevelDB for IndexedDB storage.
HBase (Apache) is a distributed LSM-tree store modeled after Google's Bigtable. It runs on top of HDFS (Hadoop Distributed File System), storing SSTables as HDFS files. HBase handles large-scale random read-write access patterns on datasets too large for a single machine, with automatic sharding across a cluster.
The common thread across all these systems is the same core architecture: memtable, immutable SSTables, compaction, and Bloom filters. The differences are in compaction strategy, distributed coordination, and the API exposed to applications. Understanding the LSM-tree fundamentals lets you reason about any of these systems without memorizing implementation details.
One practical consideration when operating LSM-tree engines: monitor space amplification. Because compaction creates new SSTables before deleting old ones, the database temporarily uses more disk space than the actual data size. In the worst case with size-tiered compaction, the database may need 2x the data size in free disk space to complete a compaction run. Running out of disk space during compaction can freeze writes and require manual intervention to recover. Production deployments typically provision 2-3x the expected data size in disk capacity and set alerts at 60-70% utilization. Some operators schedule compaction during off-peak hours by throttling compaction throughput during busy periods and allowing full-speed compaction overnight, trading temporarily higher read amplification for more predictable foreground write performance during peak traffic.