4.9 Technology deep dives
9.1 LSM-tree engines (RocksDB, Cassandra, ScyllaDB, HBase, LevelDB)
Problem it solves. Sustain very high write throughput on disks that strongly prefer sequential writes, while still supporting sorted key access and range scans.
Why wasn't a B-tree enough? B-trees convert scattered key writes into random page overwrites, plus they write every value at least twice (WAL + page), sometimes writing a whole page for a few changed bytes. On write-heavy workloads that caps throughput and burns SSD endurance.
Why wasn't the hash-index-over-a-log enough? The hash table must fit in memory, restarts require a full log scan, and range queries are impossible.
How it works internally. §2 in full: memtable → immutable sorted segments → background merge/compaction; WAL for memtable durability; tombstones for deletes; Bloom filter + sparse index per segment; size-tiered or leveled compaction.
Deployment. Embedded (RocksDB inside another system — Kafka Streams, TiKV, CockroachDB, MyRocks all embed it) or as a distributed database (Cassandra/Scylla: peer-to-peer ring, replication factor + tunable consistency; HBase: on HDFS with a master).
Monitoring.
- Compaction debt / pending compaction bytes — the single most predictive metric. If it climbs monotonically, you are losing and will eventually stall.
- Write stalls / backpressure events — RocksDB explicitly suspends writes when the memtable can't flush fast enough; count these.
- Number of SSTables per level and L0 file count (L0 overload is the classic read-amplification cliff)
- Read amplification — SSTables touched per lookup; Bloom filter false-positive rate
- Space amplification — bytes on disk ÷ logical bytes; especially during size-tiered compaction
- Tombstone counts and droppable-tombstone ratio (Cassandra) — this is where "deleted data still returns" and "queries suddenly time out" come from
Scaling. Horizontal via partitioning (Ch 7). Vertically: more memory for memtables + block cache, faster NVMe, and more compaction threads — but compaction competes with foreground I/O, so this is a genuine dial, not a free win.
Backup. Trivially good: snapshot = flush the memtable, record which segment files exist, and hard-link them. Files are immutable, so nothing needs copying until compaction wants to delete them. This is why RocksDB checkpoints and Cassandra snapshots are near-instant.
What actually breaks in production.
- Compaction can't keep up → L0 files pile up → read amplification explodes → the engine applies backpressure → write latency goes from 1 ms to 30 s with no code change. Almost always caused by write rate exceeding sustained compaction throughput, or too few compaction threads, or a saturated disk.
- Tombstone accumulation. Cassandra's canonical failure: a queue-like table where rows are written and deleted. A read scanning past 100,000 tombstones times out. Deletes are writes, and they're not free until they've propagated through every compaction level.
- Range queries far slower than point queries, because Bloom filters don't help and every segment must be scanned. People benchmark point lookups, deploy, then discover their scan workload.
- Space amplification during size-tiered compaction — merging four large tables needs temporary disk equal to their combined size. Running at 70% disk usage is how you get a full disk during a routine compaction.
- Short benchmarks lie. An empty LSM does no compaction, so early numbers are 3–5× the sustained figure. Benchmark past the point where the working set exceeds memory and compaction is in steady state.
- "Deleted" data still on disk for compliance purposes, until tombstones reach the oldest level.
9.2 B-tree engines (PostgreSQL, MySQL/InnoDB, SQL Server, Oracle)
Problem it solves. Predictable, low-latency point lookups and range scans by key, with in-place updates and mature transactional semantics.
Why wasn't an LSM enough? Reads must consult multiple segments; range queries can't use Bloom filters; compaction introduces background load and latency spikes; and the mature transaction/isolation machinery of relational databases grew up around page-based storage.
How it works internally. §3 in full: fixed-size pages, page numbers as on-disk pointers, branching factor in the hundreds, depth 3–4, splits propagating to the root, WAL for crash recovery, buffered dirty pages flushed later.
Deployment. Primary + replicas via WAL shipping (Ch 6). Buffer pool sized to a large fraction of RAM (Postgres conventionally 25% shared_buffers plus reliance on the OS page cache; InnoDB conventionally 70–80% in innodb_buffer_pool_size because it bypasses the OS cache).
Monitoring.
- Buffer pool / cache hit ratio — the cliff-edge metric; when the working set stops fitting, latency changes by orders of magnitude
- WAL generation rate and checkpoint frequency — a checkpoint storm is a classic latency spike
- Index bloat and table bloat; vacuum progress (Postgres) — §4.4's fragmentation problem, made operational
- Page splits per second (InnoDB) — high rates mean a bad insert pattern
fsynclatency — durability is bounded by it; a degraded disk shows up here first- Lock waits, deadlocks, long-running transactions
Scaling. Read replicas; partitioning; bigger buffer pool. The dominant lever is keeping the working set in memory.
Backup. Base backup + WAL archiving → point-in-time recovery. The WAL is what makes both replication and PITR possible — same mechanism, two uses.
What actually breaks.
- Random-insert index hotspots vs. monotonic keys. A UUIDv4 primary key on InnoDB scatters inserts across the whole clustered index → random writes, constant page splits, poor fill factor. A monotonic key (UUIDv7, auto-increment) keeps inserts at the right edge — but then that page becomes a contention hotspot under high concurrency. Both directions have a failure mode.
- Index bloat after mass deletes: pages remain allocated in the middle of the file and can't be returned to the OS without a rebuild.
- Torn pages on hardware without atomic page writes — which is why Postgres has
full_page_writes(and why it inflates WAL volume so much after each checkpoint). - Checkpoint spikes — buffered dirty pages all flushed at once, saturating the disk and stalling foreground queries.
- Too many indexes. Every index multiplies write cost; the "just add an index" reflex silently halves write throughput.
- Long-running transaction blocks vacuum → bloat grows unbounded → eventually the table is mostly dead tuples and even indexed queries slow down.
9.3 Columnar formats and engines (Parquet/ORC + Snowflake, BigQuery, ClickHouse, DuckDB)
Problem it solves. Aggregate over billions of rows while reading only the few columns the query touches, and spend as little CPU per row as possible.
Why wasn't a row store enough? A 100+-column fact table forces a row store to read ~100× more bytes than the query needs, then parse and discard them.
Why weren't indexes enough? Indexes help you find rows; they don't help when you must touch a billion rows regardless. Analytical queries are scan-bound, not lookup-bound.
How it works internally. §7 in full: column chunks with per-column codecs; bitmap + RLE encodings; sorting to create long runs; block/row-group granularity with min/max statistics for skipping; LSM-style batched writes; JIT compilation or vectorized batch operators; SIMD; operating directly on compressed data.
Parquet file anatomy (worth knowing concretely):
Deployment. Files in object storage, a table format (Iceberg/Delta) for atomicity and time travel, a catalog for discovery, a query engine on top — the four-layer split of §7.1.
Monitoring. Bytes scanned per query (the cost metric); row groups pruned vs read (if pruning is 0%, your sort/partition key is wrong); file size distribution; spill-to-disk volume; per-operator time in the query profile.
Scaling. Partition + sort so that predicates prune; size row groups around 128 MB; scale compute independently of storage.
Backup. Files are immutable and versioned by the table format; the catalog is the thing that must be backed up.
What actually breaks.
- The small-file problem. A streaming writer producing a file per minute yields millions of tiny Parquet files. Query planning time exceeds query time, because the engine must read a footer per file. Requires scheduled compaction.
- Wrong sort key, chosen once and effectively permanent, so no query ever prunes and every query is a full scan.
SELECT *in a columnar store — defeats the entire design and is often 50× slower than naming columns.- High-cardinality
GROUP BYOOM — the hash table doesn't fit; the engine spills or dies. - Row-at-a-time inserts into ClickHouse/Druid → "too many parts", the columnar analogue of L0 overload.
- Schema evolution mismatches — a column added in 2025 doesn't exist in 2023 files; engines differ in whether that's a null or an error.
- Timezone/type coercion between writer and reader (INT96 timestamps, anyone) producing silently wrong dates.
9.4 Full-text search engines (Lucene / Elasticsearch / OpenSearch)
Problem it solves. Find documents by keywords appearing anywhere in the text, ranked by relevance, with tolerance for typos and grammatical variation.
Why wasn't LIKE '%term%' enough? It's a full scan with no index and no ranking. Why wasn't a B-tree on the text column enough? It indexes whole values, not the terms inside them.
How it works internally. Analysis chain (tokenize → lowercase → stopwords → stemming → synonyms) produces terms; inverted index maps term → postings list; postings stored in SSTable-like sorted files merged in the background — Lucene is an LSM engine, and its "segments" and "merges" are exactly §2's. Ranking via BM25. Fuzzy matching via a Levenshtein automaton over a trie of terms. Numeric/geo fields use BKD-trees.
Deployment. Index split into shards (each an independent Lucene index) with replicas; a coordinating node scatters the query to shards and gathers/merges results.
Monitoring. Segment count and merge queue; JVM heap pressure and GC pauses (the classic Elasticsearch failure); indexing rate vs refresh interval; search latency by phase (query vs fetch); field data / fielddata cache size; shard count per node.
Scaling. Shard count is fixed at index creation and is the single most consequential decision — too few and you can't scale; too many and per-shard overhead dominates (each shard is a full Lucene index with its own segments and merges). Reindex to change it.
Backup. Snapshot API to object storage; incremental because segments are immutable — the same property as §4.4.
What actually breaks.
- Oversharding. Thousands of tiny shards, each with segments and merge threads, consuming heap for metadata until the cluster falls over.
- Mapping explosion — dynamic mapping on user-supplied JSON keys creates thousands of fields, and the cluster state grows until updates stall.
- Deep pagination (
from: 100000) — every shard must produce and sort 100,000+ hits. Usesearch_after. - Refresh interval vs indexing throughput — a 1 s refresh during a bulk load creates enormous segment churn.
- Split brain / red cluster after a network blip with badly configured quorum settings (Ch 9–10 territory).
- Using it as a system of record. It's a derived data system; it has no transactions worth the name, and its answer to corruption is "reindex from the source."
9.5 Vector databases (pgvector, Faiss, and dedicated stores)
Problem it solves. Retrieve semantically similar items — the retrieval half of RAG — where lexical matching fails because the query and the document share no words.
Why wasn't full-text search enough? "How to close my account" and "canceling your subscription" have zero terms in common. Synonym lists don't generalize.
Why wasn't an R-tree enough? R-trees don't work well for vectors with many dimensions — above roughly 10–20 dimensions, space-partitioning indexes degenerate to full scans (the curse of dimensionality). Embeddings have 384–4096 dimensions.
How it works internally. §8.3: flat (exact, slow), IVF (partition by centroid, nprobe controls the accuracy/speed dial), HNSW (multi-layer proximity graph, greedy descent). Plus quantization — PQ/SQ compress vectors to a fraction of their size, trading recall for memory, which is what makes billion-scale indexes affordable.
Deployment. Either an extension on your existing database (pgvector — strongly preferred when the vector count is in the millions, because you keep transactions and joins with metadata) or a dedicated store when you need billions of vectors or heavy filtered search.
Monitoring. Recall@k measured against a flat/exact baseline — this is the metric everyone forgets, and without it you cannot tell a broken index from a bad embedding model. Also: index build time and memory, query latency vs nprobe/ef_search, and index staleness after writes.
Scaling. HNSW memory is roughly (dim × 4 bytes + M × 8 bytes) × N — a million 1536-dim vectors is ~6 GB before graph overhead. Quantize, or shard.
Backup. The embeddings are derived data — regenerable from source documents, but regeneration costs real money in embedding-API calls, so back up the vectors. Critically: back up which model version produced them.
What actually breaks.
- Silent recall collapse.
ef_search/nprobeset too low, so the index returns plausible-but-wrong neighbours. Nothing errors. Search quality quietly degrades and users just stop trusting it. - Mixing embeddings from two model versions in one index. Vectors from different models are not comparable, and the results are noise. This happens every time someone upgrades an embedding model without a full reindex.
- Filtered search performance. "Nearest neighbours where tenant_id = X" is the hard case: pre-filtering breaks the graph's connectivity, post-filtering may return nothing. This is the #1 real-world vector-search problem and most benchmarks ignore it.
- HNSW deletes — most implementations only tombstone, so a heavily updated index degrades until rebuilt.
- Index build time measured in hours, blocking a deploy nobody planned for.
- Normalization mismatch — cosine similarity on unnormalized vectors gives wrong rankings.