Why B-Trees are Slow for Writes
Most relational databases (PostgreSQL, MySQL) use B+ Trees.
- Read: Fast
- Write: Slow. A write requires finding the correct page on disk and modifying it. This is Random I/O.
To handle millions of writes per second (e.g., IoT data, Chat logs), we need a data structure optimized for Sequential I/O. Enter the LSM Tree (Log-Structured Merge Tree).
Components
MemTable (In-Memory)
When a write comes in (PUT key=value), it is written to an in-memory structure called the MemTable.
- Structure: Usually a Red-Black Tree or SkipList.
- Why?: Keeps data sorted in RAM. Writes are fast ( memory op).
- WAL: To prevent data loss if power fails, we write to a sequential Write Ahead Log on disk first.
SSTable (Sorted String Table) (On-Disk)
When the MemTable gets big (e.g., 64MB), it is flushed to disk as an SSTable.
- Immutable: Once written, an SSTable is never modified.
- Sorted: Keys are sorted, allowing binary search scanning.
Bloom Filter
Problem: If we have 100 SSTables on disk, do we search all of them for a key? Solution: Every SSTable has a Bloom Filter.
- It tells us "The key is definitely NOT in this file" or "Maybe in this file".
- Saves 99% of unnecessary disk seeks.
The Write Path
- Append to WAL (for durability).
- Insert into MemTable (RAM).
- Ack
200 OK. - Background: If MemTable > Limit, flush to new SSTable.
The Read Path
- Check MemTable (Is it in RAM?).
- Check Block Cache (Is it in LRU Cache?).
- Check Bloom Filter of SSTable L0 (Most recent).
- Check Bloom Filter of SSTable L... until found
Note: "Read Amplification". Reads are slower than writes because we check multiple levels.
Compaction (The Cleanup)
Over time, you have duplicates.
- SSTable:
Key A = 10 - SSTable 5:
Key A = 20(Newer)
To save space and speed up reads, we run Compaction.
- Read multiple SSTables.
- Merge them into one new larger SSTable (like Merge Sort).
- Discard old values and Deleted keys (Tombstones).
Strategies
- Size-Tiered: Optimized for Write throughput (Cassandra default). Fast writes slower reads.
- Leveled: Optimized for Read throughput (LevelDB/RocksDB). More aggressive merging means fewer files to search.
B-Tree vs LSM Tree
| Feature | B-Tree (MySQL) | LSM Tree (Cassandra) |
|---|---|---|
| Write Speed | Slow (Random I/O) | Fast (Sequential I/O) |
| Read Speed | Fast | Slower (Check multiple levels) |
| Space | Fragmentation | Compact (Compressed) |
| Use Case | Financial, SQL | Logs, Time-Series, Chat |
Production Examples
RocksDB (Meta/Facebook)
Used for social graph storage and MyRocks MySQL storage engine. Handles 100,000+ writes/sec per instance.
Apache Cassandra
Powers Netflix recommendations, Apple iCloud, Instagram messages. 1,000+ node clusters handling millions of writes/sec.
DynamoDB (AWS)
Uses LSM-based storage engine optimized for consistent single-digit millisecond latency. Handles 10+ million requests/sec.
The Three Amplifications: The LSM Trade-off Triangle
Every LSM tuning debate reduces to three quantities, and you cannot minimize all three at once:
- Write amplification: bytes physically written to disk ÷ bytes the user wrote. Every compaction rewrites data that was already written — leveled compaction can amplify a single logical write by 10–30x over its lifetime. This burns SSD endurance and steals I/O bandwidth from foreground work.
- Read amplification: disk reads needed to answer one query. More un-compacted SSTables = more places a key might hide.
- Space amplification: disk consumed ÷ live data size. Old versions and tombstones linger until compaction reclaims them; size-tiered compaction can transiently need 2x the data size during a large merge.
Size-tiered compaction chooses low write amplification (merge rarely, in big batches) at the cost of read and space amplification. Leveled chooses low read and space amplification at the cost of rewriting data more often. This is why Cassandra (write-heavy telemetry) defaults to size-tiered while RocksDB (read-mostly serving) defaults to leveled — same structure, opposite bets.
Tombstones: How Deletion Becomes a Performance Problem
In an immutable-file world you cannot erase a key — you write a tombstone (a "this key is deleted" marker) that shadows older values until compaction physically drops both. Two operational consequences:
- Deletes temporarily increase disk usage. Deleting a million rows writes a million new records. Teams cleaning up disk space by mass-deleting are routinely surprised when usage goes up for days.
- Tombstone scans wreck range reads. A query scanning a range where 99% of rows were deleted must still read and discard every tombstone. Cassandra emits
TombstoneOverwhelmingExceptionwarnings for exactly this; queue-like workloads (insert, consume, delete, repeat) are the classic anti-pattern on LSM databases.
Tombstones can only be dropped after a grace period (gc_grace_seconds in Cassandra, default 10 days) that gives replication time to propagate the delete to all replicas — drop it too early and a lagging replica can "resurrect" the deleted row during repair.
Write Stalls: When the Background Catches Up With You
Compaction is "background" work only while it keeps pace. If sustained write throughput exceeds compaction throughput, L0 files pile up, reads degrade, and eventually the engine applies backpressure: RocksDB first throttles writes (slowdown trigger), then stops them entirely (stop trigger) until compaction catches up. The symptom in production is a saw-tooth latency pattern — smooth for minutes, then a spike when a big compaction runs or a stall hits.
Levers that matter in practice: compaction thread count and I/O rate limits (isolate from foreground I/O), MemTable size (bigger = fewer, larger flushes), and for bulk loads, bypassing the write path entirely by pre-building SSTables and ingesting them directly.
Summary
- MemTable: Buffers writes in RAM (Sorted).
- SSTable: Immutable files on disk (Sorted).
- Bloom Filter: Optimization to avoid disk reads.
- Compaction: Background merge sort to cleanup garbage — and the source of the write/read/space amplification triangle.
- Tombstones: Deletes are writes; physical cleanup is deferred and replication-aware.
- Use Case: Write-heavy workloads like time-series, logs, analytics
Related Concepts
- Redis Internals — In-memory data structures
- Database Sharding — Horizontal partitioning
- CAP Theorem — Availability vs consistency trade-offs
About ScaleWiki
ScaleWiki is an interactive educational platform dedicated to demystifying distributed systems, software architecture, and system design. Our mission is to provide high-quality, technically accurate resources for software engineers preparing for interviews or solving complex scaling challenges in production.
Read more about our Editorial Guidelines & Authorship.
Educational Disclaimer: The architectural patterns and system designs discussed in this article are based on common industry practices, technical whitepapers, and public engineering blogs. Actual implementations in enterprise environments may vary significantly based on specific product requirements, legacy constraints, and evolving technologies.
Related Articles
Redis Internals: Why is it Fast?
Deep dive into Redis architecture: single-threaded event loop, data structures, persistence strategies (RDB/AOF), replication, and cluster mode.
Document Databases
The most popular type of NoSQL database. Storing data in flexible, JSON-like documents with embedded structures and dynamic schemas.
Caching Overview
High-speed data storage to reduce latency. The single most effective way to scale read-heavy systems.