Storage engines, write-ahead logs and durability
What happens inside a database when you write: write-ahead logs and fsync, B-trees versus LSM trees, memtables, SSTables and compaction, read, write and space amplification, group commit, checkpoints and crash recovery.
Reading is half of it. See this used in a real interview: walk through Design a Distributed Key-Value Store →
"Design a key-value store" or "design a message queue" eventually asks what happens on a single node when a write arrives, and how it survives a crash. The answer is a storage engine: a write-ahead log for durability, plus an on-disk structure (a B-tree or an LSM tree) for fast lookups. Knowing how these work explains the performance of almost every database you will name in an interview.
Durability and the write-ahead log
Writing a change directly into a data file in place is slow (random I/O) and unsafe (a crash halfway leaves a corrupt page). Instead, engines first append the change to a write-ahead log (WAL):
- Append the record to the log (sequential I/O, fast).
fsyncso it is really on durable storage, not just in the OS cache.- Acknowledge the write.
- Apply the change to the main data structure later, in memory and eventually on disk.
After a crash, the engine replays the log from the last checkpoint to rebuild state. Every serious database (Postgres, MySQL InnoDB, RocksDB, Kafka's log itself) works this way.
Group commit: fsync is expensive (tens of microseconds to milliseconds), so engines batch many concurrent writes into one fsync. Throughput rises sharply at a small latency cost. Some systems relax durability (fsync every second) for speed, risking the last second of writes.
B-trees
The classic structure for relational databases:
- Data lives in fixed-size pages (for example 8 to 16 KB) arranged in a balanced tree; lookups walk a few levels from root to leaf.
- Updates modify pages in place (logged in the WAL first).
- Reads are fast and predictable: one path down the tree.
- Writes cause random I/O and can rewrite a whole page for a small change.
Great for read-heavy workloads and range scans. See database indexing.
LSM trees
Log-structured merge trees optimise for writes (RocksDB, LevelDB, Cassandra, ScyllaDB, HBase):
- Writes go to the WAL and to an in-memory sorted memtable.
- When the memtable is full, it is flushed to disk as an immutable sorted file (SSTable).
- Background compaction merges SSTables, discarding overwritten values and deleted keys (tombstones).
- Reads check the memtable, then SSTables from newest to oldest, using Bloom filters to skip files that cannot contain the key and sparse indexes to find blocks.
All disk writes are sequential, so write throughput is very high. Reads may touch several files, and compaction uses background I/O and CPU. See probabilistic data structures.
The amplification trade-off
| B-tree | LSM tree | |
|---|---|---|
| Write amplification | high (page rewrites) | medium to high (compaction rewrites data) |
| Read amplification | low | higher (several levels), reduced by Bloom filters |
| Space amplification | some (fragmentation) | depends on compaction strategy |
| Write throughput | good | excellent |
| Read latency | predictable | mostly good, spikes during heavy compaction |
| Typical use | OLTP, relational | write-heavy, time series, key-value at scale |
Compaction strategies trade these against each other: leveled compaction keeps reads and space low at the cost of more writing; size-tiered writes less but uses more space and makes reads slower. Time-windowed compaction suits time series, where whole old files are dropped. See time-series data.
Checkpoints and recovery
The WAL cannot grow forever. A checkpoint writes dirty pages (or flushes memtables) to the data files and records a position; the log before that point can then be deleted or archived. Recovery time depends on how much log follows the last checkpoint. Archived WAL segments also enable point-in-time restore. See backups and disaster recovery.
Logs everywhere
The same append-only log idea powers:
- Replication: replicas receive and replay the leader's log. See replication and consistency.
- Message queues: Kafka is a partitioned, replicated append-only log with segment files and an index. See Design a Message Queue.
- Change data capture: reading the database log to stream changes. See change data capture.
In the interview
For Design a Key-Value Store: writes go to a WAL (with group commit) and a memtable, flush to SSTables, compact in the background, Bloom filters per SSTable for reads, tombstones for deletes, and replication on top. Explain why LSM suits a write-heavy store and what it costs on reads.
Checklist
- Append to a WAL and fsync before acknowledging; group commit for throughput.
- B-trees for read-heavy, range-scan workloads; LSM trees for write-heavy ones.
- Memtables, SSTables, Bloom filters and compaction in LSM engines.
- Read, write and space amplification as the core trade-off.
- Checkpoints to bound recovery time; archived logs for point-in-time restore.