Write-Ahead Log (WAL) & LSM-Trees
The Time-Series Database That Froze on Flush
Test your architecture intuition: Pitch a 7-axis solution, survive two aggressive reviewer objections, and inspect the staff-level Teacher Gold Answer.
1. What It Is & Why It Exists
The Core Problem: The B+ Tree In-Place Update Penalty
Traditional relational database storage engines (e.g., MySQL InnoDB's clustered B+ tree, PostgreSQL's heap pages with B-tree indexes) use B+ Trees and page files organized into fixed-size disk pages (8 KB in PostgreSQL, 16 KB in InnoDB). When an application updates or inserts a record:
- The engine traverses the B+ Tree index, loads the target page into memory (the buffer pool), and modifies the record.
- The dirty page is eventually written back to disk in-place at its fixed physical address.
- In high-velocity write workloads with random keys (such as UUIDs, user sessions, or sensor events), modifications scatter across thousands of disparate pages. Writing these pages forces the underlying storage controller to perform massive Random Disk I/O.
On modern NVMe Solid-State Drives (SSDs), random writes trigger extra internal write amplification due to NAND flash constraints (flash is written in pages but erased only in much larger blocks, so the drive's garbage collection copies live data around). On mechanical hard drives, random seeks incur a mechanical penalty per operation, capping throughput at about 100 to 200 random operations per second.
In contrast, Sequential Disk I/O streams data continuously, bypassing head seeks and maximizing SSD controller parallelism to reach the drive's full bandwidth.
The Breakdown: The In-Place Page Rewrite Cliff
Concrete Proof: B+ Tree Random Page Rewrites vs. LSM-Tree Sequential Appends
Consider a distributed storage system ingesting with record payloads ( raw data ingestion):
| Metric / Operational Dimension | B+ Tree Storage Engine (In-Place 16 KB Pages) | LSM-Tree Engine with WAL (Sequential Append) |
|---|---|---|
| Disk Write Pattern | 🚨 Scattered Random Writes across multi-GB page pools | ⚡ Pure Sequential Appends to WAL & MemTable |
| Physical Data Written to Disk | Updating a 1 KB record rewrites the entire 16 KB page: | Appends 1 KB record to sequential log: (with metadata) |
| Write Amplification (WA) | 🚨 up to (a 16 KB page written twice for crash safety, plus a 1 KB log record: ) | 🛡️ at initial ingest (often in total after leveled compaction, all sequential) |
| Disk IOPS Demand | 🚨 (Saturates storage queues) | 🛡️ About log syncs a second (group commit at a 1 ms sync) plus large sequential flush and compaction I/Os |
| P99 Write Latency | 🚨 Queues behind random page write-back and checkpoints | ⚡ One log sync per batch: between and with group commit ( at a 1 ms sync) |
| SSD Hardware Endurance | 🚨 Rapid flash wear-out (Exhausts Drive Writes Per Day) | 🟢 With random keys, fewer bytes written and all sequential (compaction still rewrites each byte many times) |
The First-Principles Solution: Append-Only LSM-Trees & WAL
The Log-Structured Merge-Tree (LSM-Tree), designed by Patrick O'Neil et al. in 1996, resolves the random write bottleneck by converting all incoming mutations (inserts, updates, and deletes) into strictly sequential disk appends:
- Write-Ahead Log (WAL): The mutation is immediately appended to an append-only log before being acknowledged. It provides ACID Durability () against power loss only if the log is synced (
fsync/fdatasync) before the acknowledgement; RocksDB's default (WriteOptions.sync = false) hands the record to the OS without syncing it (see the section Durability modes: what "acknowledged" means of Write-Ahead Log, fsync & Group Commit). - MemTable: Then the record is inserted into an in-memory sorted data structure (typically a Concurrent SkipList or Red-Black Tree), allowing immediate search and range scans.
- Immutable SSTable Flush: When the MemTable reaches its capacity threshold (e.g., 64 MB), it is frozen into a read-only MemTable and flushed sequentially to disk as a sorted, immutable Sorted String Table (SSTable).
- Background Compaction: Background threads periodically merge overlapping SSTables, discard obsolete versions, and remove tombstoned records (a tombstone only once nothing older can exist below it), maintaining logarithmic search performance.
Synthesizing vector architecture diagram...
Follow one PUT from the top. In the "Sequential Ingest Path" panel, the write is appended to the WAL on disk (for durability, once the log is synced) and then inserted into the in-memory MemTable, then acknowledged, so there are no random disk writes on the critical path. In the "Immutable SSTable Flush Tier" panel, once the MemTable reaches 64 MB it is frozen and written sequentially to disk as a Level 0 SSTable, whose key ranges may overlap with other L0 files. In the "Asynchronous Leveled Compaction Tier" panel, background merges push data down into levels about 10 times larger (256 MB, 2.5 GB, 25 GB) where files never overlap. Writes are fast because they are only appends; the cost is paid later by compaction, and by reads that may have to check several levels.
2. Core Mechanics & Algorithmic Architecture
The Fundamental Distributed Storage Triangle: R-W-S Amplification
Every persistent storage engine operates within the boundaries of the R-W-S Amplification Trade-off:
Synthesizing vector architecture diagram...
Each corner is one cost, and each arrow means that reducing one cost raises another. Write amplification (WA) counts how many bytes hit the disk per byte you write: lowest in LSM-trees (append now, merge later), high in B+ trees (rewrite a whole page for a small change). Read amplification (RA) counts bytes read per lookup: lowest in B+ trees (one path of pages), higher in LSM-trees (check several levels). Space amplification (SA) counts disk used per byte of live data: low in B+ trees (update in place), higher in LSM-trees (old versions wait for compaction). No engine wins all three corners, so choose by workload: write-heavy leans LSM, read-heavy leans B+ tree.
Leveled Compaction Sizing Mathematics
In Leveled Compaction (LCS) (used by RocksDB and TiKV), disk storage is organized into discrete exponential levels ():
- Level 0 holds SSTables flushed directly from MemTables. Because they are flushed independently, Level 0 SSTables have overlapping key ranges.
- Levels and higher are strictly sorted such that no two SSTables in the same level share overlapping key ranges.
- Each level has a maximum capacity that grows exponentially by a multiplier (typically ):
If and :
The total worst-case Write Amplification in Leveled Compaction is bounded by:
Where is the total database size.
Step-by-Step Deterministic Write & Read Pipelines
Synthesizing vector architecture diagram...
In-Memory & On-Disk State Structures
1. In-Memory SkipList Architecture
The MemTable is backed by a lock-free Concurrent SkipList. Elements are arranged in linked lists with probabilistic multi-level forward pointers, providing search, insertion, and range iterations without global lock contention.
2. Immutable SSTable File Anatomy
When a MemTable flushes to disk, it produces an immutable SSTable file structured into fixed-size blocks (typically 4 KB):
| Block Name | Format / Contents | Purpose | Memory Location |
|---|---|---|---|
| Data Blocks | Compressed key-value pairs sorted lexicographically | Holds raw payload data (compressed via Zstandard or Snappy) | Loaded into Block Cache in RAM on-demand |
| Filter Block | Bit arrays representing Bloom filters for all data blocks | Evaluates whether a key definitely does not exist in this SSTable | Pinned in RAM off-heap |
| Index Block | Sparse index: one entry per data block (a separator key and the block's file offset); optionally partitioned into two levels | Enables binary search to locate the exact 4 KB data block for a key | Cached in RAM |
| Footer Block | Fixed-size trailer (48 bytes in LevelDB and RocksDB's legacy format, 53 bytes in its newer formats) containing magic number and handles to Index & Meta blocks | Anchor point read by storage engine when opening an SSTable | Loaded on file open |
Synthesizing vector architecture diagram...
The file is written left to right, but read from the right. A reader first opens the fixed-size Footer at the end of the file, which points to the Meta Index and Index Block. For a lookup, the Filter Block's Bloom filter answers first: if the key is definitely absent, the read stops without touching data. Otherwise the sparse Index Block gives the one 4 KB data block that could hold the key, and only that block is read and decompressed. So a lookup costs about one small block read per file, because the filter and index are usually cached in RAM.
Scenario Execution Matrix: LSM Read, Write, and Compaction Dynamics
| Step # | Event / Input | In-Memory / Distributed State | Evaluation & Transition | Outcome / Latency & IOPS |
|---|---|---|---|---|
| 1 | Point Lookup on Hot KeyGET "user_102" | Key resident in active RAM MemTable (SkipList) | Skip-list search on the concurrent SkipList locates entry; zero disk access required | Direct RAM Hit Returned from memory; zero Bloom filter checks, 0 disk IOPS |
| 2 | Point Lookup on Cold KeyGET "order_4821" | Absent from MemTable; present on disk in Level 2 SSTable | 1. MemTable miss check L0 Bloom filters (negative, skip) 2. L1 non-overlapping index binary search (not in range) 3. L2 Bloom filter returns 1 read sparse index | Single Disk Read Targeted data block read from NVMe; completely bypasses L0 and L1 |
| 3 | Tombstone DeletionDELETE "user_102" | Record exists on disk across L1 and L2 SSTables | Append a tombstone (the key with a delete flag and no value) to WAL and MemTable; subsequent queries return 404 immediately | Instant Tombstone Append (plus the log sync, if on) Stale on-disk versions masked; as compaction carries the tombstone down, it drops the older versions it meets, and the tombstone itself goes once it reaches the bottom level (nothing older can exist below it) |
3. Data Migration, Anti-Entropy & Consistency Protocols
1. Compaction Strategies Compared
Synthesizing vector architecture diagram...
Both panels show one merge step. In the "Size-Tiered Compaction" panel, the engine waits until there are four similar 100 MB SSTables and merges them into one 400 MB file; each byte is rewritten only a few times, but files within a tier can overlap, so a read may have to check several. In the "Leveled Compaction" panel, each level is a set of files with non-overlapping key ranges: a file from L1 covering A-M is merged into the L2 files that overlap it, producing A-D, E-H and I-M. Leveling rewrites data more often (higher write amplification) but lets a read check at most one file per level; choose size-tiered for write-heavy workloads and leveled for read-heavy ones.
| Compaction Strategy | Mechanics | Write Amplification () | Space Amplification () | Read Latency | Best Suited For |
|---|---|---|---|---|---|
| Size-Tiered (STCS) | Accumulates SSTables of similar sizes; merges them into one large SSTable | Low (; about one rewrite per tier) | High (up to temporary disk space in the worst case, a merge of everything) | Higher (Must search multiple files per tier) | Write-heavy workloads that are not strictly time series (Apache Cassandra default); for time series that expire, Cassandra recommends Time-Window (TWCS) |
| Leveled (LCS) | Strictly non-overlapping key ranges per level (); merges into | Higher () | Low ( disk overhead) | Fast (every L0 file, then at most 1 SSTable per level below) | General OLTP workloads (RocksDB, TiKV; CockroachDB's Pebble) |
| FIFO Compaction | Discards oldest SSTables once total disk quota is exceeded | Minimal () | Zero | Every file stays in L0, so a read checks each one (Bloom filters help) | Event logs with a TTL (RocksDB suggests query logs) |
2. WAL Group Commit & Fsync Durability Protocols
Calling fsync() (or fdatasync()) forces the data through the page cache and the drive's volatile write cache to persistent storage. One log does one sync at a time, so with one write per sync it caps at writes a second for a sync latency : at , and "a few hundred" only at . Modern engines implement Group Commit:
- Incoming client threads queue their mutations. The first writer that finds no sync in progress becomes the leader: it writes every queued record to the log and issues one
fdatasync()(RocksDB's write-group leader does exactly this). - Writers that arrive during that sync queue up and form the next batch. No timer is needed: the batch is whoever arrived during the last sync. (A window variant waits on purpose, e.g. up to 2 ms or 64 KB, trading latency for fewer syncs.)
- All client requests in that batch are acknowledged once the sync returns. At writes a second and , about 100 writes share each sync ( syncs a second) and each write waits between and . That is durable against a crash or power loss of this node, if every layer honours the flush; it doesn't survive losing the disk (see the sections What fsync really guarantees, Group commit and Durability modes: what "acknowledged" means of Write-Ahead Log, fsync & Group Commit).
3. Distributed Anti-Entropy, Merkle Trees & Crash Recovery
- Merkle Tree Anti-Entropy Repair: In distributed LSM databases (e.g., Apache Cassandra, ScyllaDB), node replicas detect out-of-sync data without transferring entire datasets. Replicas construct hierarchical Merkle trees over token ranges. During operator-run anti-entropy repairs (
nodetool repair), nodes exchange tree roots () and traverse diverging child branches () to isolate diverging key ranges, streaming only the data in the mismatched ranges. (Cassandra 4.0+ can send a whole SSTable with zero-copy streaming when every partition in that file has to be transferred.) - WAL Checkpointing & Crash Recovery Protocol: On unexpected process termination, the engine reconstructs volatile MemTables deterministically from the on-disk Write-Ahead Log:
- Parse the most recent Checkpoint Marker: retrieves the last flushed Log Sequence Number ().
- Sequential Scan: Replay WAL log records starting strictly from .
- Tail Truncation on CRC Mismatch: Validate a checksum per record (CRC32C in RocksDB and PostgreSQL). If power loss caused a torn write at the file tail, the engine stops at the first bad record and discards the trailing bytes, restoring consistent state; when the engine syncs before acknowledging, those bytes belong only to writes that were never acknowledged. Damage followed by valid records is not a torn tail: what happens then is the engine's recovery mode (RocksDB's
wal_recovery_mode,kPointInTimeRecoveryby default, stops at the first corruption).
Unlock Complete Architecture & Production Runbooks
You have explored the free architectural preview (~47%). Spend 1 Coin to unlock the remaining 7 production deep-dive sections for a full 24 hours.