Storage Engines and Data Structures

Topics Covered

Log-Structured Storage

SSTables and Memtables

LSM-Trees and Compaction

Bloom Filters

Real-World LSM Engines

Page-Oriented Storage

B-Tree Structure

Page Splits and In-Place Updates

Write-Ahead Log (WAL)

Concurrency and Latching

B-Tree Optimizations

Clustered vs. Non-Clustered Indexes

Column-Oriented Storage

How Column Storage Works

Compression in Column Stores

Sort Order in Column Stores

Materialized Aggregates

Column Stores and Write Performance

In-Memory Databases

Durability Without Disk as Primary Storage

Crash Recovery

Anti-Caching and Hybrid Approaches

Memory Management Considerations

Choosing a Storage Engine

LSM-Trees vs. B-Trees

Row vs. Column Storage

When to Go In-Memory

Decision Framework

Common Anti-Patterns

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.

A write hitting the log and the memtable, flushing out to an SSTable, and later merging away during compaction.

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:

python
1from bisect import bisect_right
2
3# one entry per ~64 KB block, loaded into RAM when the SSTable opens
4SPARSE = [("aardvark", 0), ("kiwi", 65536), ("panther", 131072), ("zebra", 196608)]
5KEYS = [k for k, _ in SPARSE]
6
7def lookup(key, sstable):
8    i = bisect_right(KEYS, key) - 1
9    if i < 0:
10        return None                      # key sorts before the first block
11    start = SPARSE[i][1]
12    end = SPARSE[i + 1][1] if i + 1 < len(SPARSE) else sstable.size
13    block = sstable.read(start, end - start)   # exactly one disk read
14    return scan_block(block, key)              # linear scan inside the block

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.

Four small SSTables merged under size-tiered and leveled strategies, with the amplification each one produces.

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 TT between levels and LL levels, each byte is rewritten roughly once per level and once per merge round within a level, giving

WA≈T×LWA \approx T \times L

so a default RocksDB configuration with T=10T = 10 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 LL, 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%.

Key Insight

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.

A key checked against a Bloom filter, with the three bits that skip a disk read and the false positive that does not.

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 mm bits, nn keys and kk hash functions, the false positive rate is

p=(1−e−kn/m)kp = \left(1 - e^{-kn/m}\right)^{k}

which is minimized at k=mnln⁡2k = \frac{m}{n}\ln 2, and inverting it gives the bits you have to buy for a target rate:

python
1import math
2
3def bloom_params(n, p):
4    """bits and hash functions needed for n keys at false positive rate p"""
5    m = -n * math.log(p) / (math.log(2) ** 2)
6    k = (m / n) * math.log(2)
7    return math.ceil(m), round(k)
8
9bloom_params(1_000_000, 0.01)    # (9_585_059, 7)   ->  9.6 bits per key
10bloom_params(1_000_000, 0.001)   # (14_377_588, 10) -> 14.4 bits per key
11bloom_params(1_000_000, 0.0001)  # (19_170_117, 13) -> 19.2 bits per key

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.