Database Internals and Storage Engines Questions
How a single database engine works internally, at the mechanism level: storage-engine architectures such as B-tree versus LSM-tree, on-disk page and buffer-pool management, write-ahead logging (WAL) and checkpointing, and MVCC's internal mechanics (version chains, vacuum and bloat, transaction-ID wraparound). Also covers compaction and its write/read amplification trade-offs, including engine-specific behavior such as Bloom filters and tombstone/TTL expiry in LSM-based stores like Cassandra, and hands-on design of simplified storage components: on-disk data layouts and indexes, WAL recovery routines, compaction strategies, and crash-safe log or queue formats. Tests depth beyond usage: why the engine behaves as it does, not how to operate, tune, recover, or scale it.
Explain what a write-ahead log (WAL) is and why databases use it. Describe the typical order of operations for a durable commit using WAL and when an fsync is required to ensure durability. Mention the role of checkpointing in recovery time.
Sample Answer
Direct answer
A write-ahead log, WAL, is an append-only file where a database records a description of every change before it changes the actual data page on disk. The rule that gives it its name and its durability guarantee is that the log record describing a change must reach non-volatile storage before the transaction that made the change is allowed to report success to the client; if the process or machine crashes right after that, replaying the log on restart reproduces exactly the changes that had already been acknowledged, and nothing else.
Structured elaboration
- Order of operations for a durable commit: (1) the transaction's changes are applied to in-memory copies of the affected pages, in the buffer pool (the shared-memory cache of recently used disk pages); (2) a WAL record describing the change is appended to an in-memory WAL buffer; (3) at commit, the WAL buffer (everything up through this transaction's commit record) is flushed to disk and the storage device is told to actually persist it, an operation commonly called an fsync, which forces the operating system to stop buffering the write in memory and confirm it has reached durable storage; (4) only after that fsync returns does the database tell the client the commit succeeded; (5) the actual data page on disk can be written out later, lazily, because the WAL record alone is enough to reconstruct the change if the page write has not happened yet when a crash occurs.
- When fsync is required: any time a transaction commits and the client is being told it succeeded, unless the administrator has explicitly opted into weaker durability. Postgres exposes this as
synchronous_commit; its default,on, fsyncs the local WAL before acknowledging the commit. Setting it tooffskips that wait (a commit can be acknowledged and then lost if the server crashes within roughly the next WAL-flush interval), which trades durability for latency and has to be a deliberate decision, not a default, because it changes what "the database told me it committed" actually guarantees. - Group commit: because the fsync is the expensive part (a physical storage operation, not just a memory copy), most engines batch multiple transactions' WAL records into one fsync when they commit close together in time. Postgres additionally exposes
commit_delay(default 0) andcommit_siblings(default 5) to deliberately widen that batching window undersynchronous_commit=on, trading a little added latency per transaction for fewer, larger fsyncs under high concurrent commit load. - Checkpointing and recovery time: a checkpoint is a point the engine guarantees all data pages have actually been written to disk up through some WAL position; after a crash, recovery only needs to replay WAL starting from the most recent checkpoint, not from the beginning of time, because everything before the checkpoint is already durably reflected in the data files. That is the entire reason checkpoints exist: they bound recovery time. Postgres triggers a checkpoint on a timer (
checkpoint_timeout, default 5 minutes) or when WAL has grown past a size budget (max_wal_size, default 1 GB), and spreads the checkpoint's I/O out over most of that interval (checkpoint_completion_target, default 0.9) rather than doing it in one burst, to avoid a periodic I/O spike. Shortening the checkpoint interval reduces crash-recovery time at the cost of more frequent background I/O; lengthening it does the reverse.
Worked example
On a live Postgres 16 instance, pg_current_wal_lsn() (the current write position in the WAL, expressed as a log sequence number, LSN, a monotonically increasing byte offset into the WAL stream) read 0/16119AC8. A single INSERT of one small row was committed. pg_current_wal_lsn() afterward read 0/1611B940, an advance of 7,800 bytes for a row that itself is only a few dozen bytes wide. That gap is the concrete shape of "changes are described before they're applied": the table's page had not been modified since the last checkpoint, so Postgres's full-page-write protection (on by default) logged the entire roughly 8 KB page image the first time it changed after that checkpoint, not just the new row, precisely so that a torn (partially written) page during a future crash can be reconstructed from the WAL instead of from a corrupted on-disk page.
Trade-offs and pitfalls
The most common misunderstanding is treating "the WAL was written" and "the data file was written" as the same event: they are not, and the data file can lag the WAL by a long time, which is exactly what makes WAL cheap (sequential appends) instead of expensive (random page writes on every change). The second common mistake is tuning checkpoints purely for background I/O smoothness and forgetting the other side of the trade: a longer checkpoint_timeout or bigger max_wal_size means more WAL has to be replayed after a crash, directly increasing recovery time, on a system where recovery time is often the metric that actually matters.
What is a B-tree index and why is it well-suited to range queries? Describe briefly how B-tree insert/delete operations maintain balance, why page splits occur, and how page splits can impact write amplification and concurrency.
Sample Answer
Direct answer
A B-tree index (in almost every production database this really means a B+-tree: a self-balancing tree where every row pointer lives at the leaf level, and leaf pages are linked to their neighbors) keeps keys sorted at every level. Because the leaves are sorted and linked, finding a range's start takes one root-to-leaf descent (height, typically 2 to 4 levels even for hundreds of millions of rows), and after that the engine walks sideways along the leaf chain instead of re-descending for each row, which is why range queries are cheap. Inserts and deletes stay balanced by splitting a page that overflows and allowing a page to run under capacity rather than aggressively merging on every delete, so the tree's height only grows when the root itself splits, keeping lookups at roughly log(n) page reads.
How B-tree balance, splits, and concurrency interact
- Structure: every page holds a sorted array of (key, pointer) entries. Internal pages point to child pages; leaf pages point to the actual row (or, for a clustered index like InnoDB's primary key, contain the row itself). Fanout, how many entries fit in one page, determines height: with fanout F, a tree over N keys needs roughly ceil(log base F of N) levels.
- Insert: the engine descends to the correct leaf and inserts in sorted order. If the leaf has room, this is a single page write. If it is full, the leaf splits: half its entries move to a new page, and a separator key is inserted into the parent, pointing at the new sibling. If the parent is also full, the split cascades upward; the tree only grows a level when the root itself splits, which is what keeps height balanced rather than letting the tree degrade into an unbalanced structure.
- Delete: the entry is removed from its leaf. Well-behaved implementations let a leaf run under its target fill percentage rather than merging on every delete (merging on every delete would thrash under mixed insert/delete workloads); many engines instead reclaim mostly-empty pages lazily.
- Why splits cost more than the one write they look like: a split has to write the original page, the new sibling page, and the parent page that now holds one more separator key. Under write-ahead logging (WAL, the mechanism where every page change is described in a log record and forced to disk before the change is considered durable) each of those pages generates its own log record, so one logical insert can produce a page's worth of extra durable writes: that is the "write amplification" indexes add on top of the base table's own write. Concurrency-wise, a split has to briefly hold an exclusive lock (a short, in-memory latch, not a transactional row lock) on the splitting page and its parent while it moves entries and updates the parent's pointer; concurrent readers that land mid-split are not blocked in most modern engines (Postgres's B-tree implementation, for example, lets a reader who arrives at a just-split page follow a temporary "right link" to the new sibling instead of waiting), but concurrent writers targeting the same page do queue briefly.
Worked example
I built this on a live Postgres 16 instance: a 3,000,000-row orders table with a plain B-tree index on order_date, then ran EXPLAIN (ANALYZE, BUFFERS) (Postgres's command to show the actual, executed query plan, including real row counts and I/O, not just the planner's estimate) for a 7-day range predicate (order_date >= '2021-06-01' AND order_date < '2021-06-08'), matching 11,669 of the 3,000,000 rows.
Without the index (parallel sequential scan): Buffers: shared hit=16056 read=6003 dirtied=7080 written=5907, filtering out 996,110 rows per worker, 70.977 ms.
With the index (a bitmap index scan, which uses the index to find matching row locations instead of reading the whole table, then a bitmap heap scan to fetch just those rows): Buffers: shared hit=1252 read=526 written=464 for the base table plus shared read=13 for the index itself, roughly 1,791 buffer touches total (1,252 + 526 + 13), against about 22,059 without the index, 7.037 ms.
pageinspect's bt_metap('idx_orders_date') reports the index root is at level = 2, meaning the tree has 3 levels (root, one internal level, leaf level) over those 3,000,000 rows: 2,615 leaf-level pages at 8 KB each (20 MB total), and the root page held only 13 downlink entries (checked with bt_page_stats), because at the root you only need enough entries to address every second-level page, not every leaf. That 13-entry root points directly to those same 13 second-level pages, and each second-level page in turn addresses roughly 201 leaf pages (2,615 leaf pages / 13 second-level pages is about 201): together that is the concrete shape of "height stays logarithmic": adding a million more rows would not add a fourth level until fanout at the second level was exhausted too.
Trade-offs and pitfalls
Every additional B-tree index is a second (or third, or tenth) sorted copy of a subset of the table's columns: each insert into the base table now also does an insert into every index on it, each of those inserts can itself trigger a split, and each split is extra write-ahead-log volume and extra page writes on top of the row's own write. On a workload that mixes a frequent-write pipeline (an extract-transform-load, ETL, job loading rows continuously) with frequent ad hoc dashboard reads, this is a real trade-off: the indexes that make the dashboard's queries fast are the same indexes slowing down the ETL job's inserts and growing on-disk size roughly in proportion to how many columns are indexed. The common mistake is treating "add an index" as free because one SELECT got faster; the honest comparison is against the write path's added cost, not just the read path's benefit. A second pitfall is assuming insert cost is independent of key shape. Monotonically increasing keys (an auto-increment ID, a timestamp) mostly append at the rightmost leaf, which splits far less often than a workload inserting uniformly random keys across the whole key space, because random inserts hit a different, potentially already-full leaf every time.
Explain the difference between a database's logical storage (tables, schemas, indexes, views) and its physical storage (pages, files, and how those pages are actually organized on disk or in object storage). Why does an engine maintain this separation, and how does the physical storage choice affect query latency and cost for a read-heavy reporting workload?
Sample Answer
Direct answer
A database's logical storage is the model applications see: tables, columns, indexes, views, the schema queries are written against. Its physical storage is how that logical model is actually laid out as bytes: fixed-size pages (commonly 4 to 16 KB) grouped into files, and where those files physically live, a local block device, a network-attached disk, or an object store. The engine maintains a strict separation between the two so it can change how data is physically organized, or even what kind of storage backs it, to fix performance or reduce cost, without any application code changing at all.
Structured elaboration
- Why the separation exists: the logical model is a contract (an application's queries do not say "read block 4,391 of file segment 7"); the physical layout is an implementation detail the engine is free to change. This lets the engine reorganize physical layout for performance, defragmenting pages, changing compression, moving cold data to cheaper storage, as a purely internal operation. It also lets the same logical schema be served by very different physical backends, block-oriented local storage versus an object store, without the application noticing, which matters a lot for the cost and latency trade-off below.
- What physical storage actually consists of: pages are the unit of I/O, the engine reads and writes whole pages, not individual rows, even when a query only touches one row, because the underlying storage medium is efficient at page-sized (or larger) transfers and inefficient at single-byte ones. A buffer pool (an in-memory cache of recently used pages) sits between the logical query engine and physical storage so that "hot" pages do not have to be re-read from disk on every access.
- The consequence of ignoring physical layout: a schema that is logically fine (correct columns, correct types) can still perform badly if its physical access pattern is bad, many small random page reads instead of a few large sequential ones, which is exactly the axis storage-engine design (B-tree versus log-structured merge-tree layouts) optimizes for.
Worked example
Keep this to the one comparison the question actually asks about, block-based versus object storage, for a read-heavy reporting workload: a dashboard running many ad hoc aggregate queries over a large historical table. Local block storage (an attached solid-state drive) serves a random 8 KB page read in roughly sub-millisecond time, commonly cited in the tens to low hundreds of microseconds for solid-state media. An object store (the kind of storage behind Amazon S3-style APIs, addressed by whole-object network requests rather than in-place random reads) typically serves one such request in the tens of milliseconds, one to two orders of magnitude slower per request, because each request is a full network round trip rather than a local I/O operation. A reporting query that would need a few thousand scattered small reads against block storage is a workload object storage handles badly unless the engine restructures it, reading large, sequential, compressed batches instead of many small random reads, which is exactly why systems built on top of object storage are designed around big sequential scans rather than random point access. These are illustrative, commonly cited orders of magnitude for the two storage classes, not a benchmark run for this answer, and actual numbers vary by hardware and provider.
Trade-offs and pitfalls
The pitfall is assuming the logical and physical split means physical choices are free to ignore: they are not, they just do not require changing your schema or your queries to fix. A common mistake is choosing a physical backend for cost (object storage is typically far cheaper per gigabyte than provisioned block storage) without accounting for the very different latency and access-pattern assumptions it requires, then being surprised the same reporting queries that were fast on block storage are slow after a supposedly transparent migration.
What is MVCC (multi-version concurrency control)? Describe how it enables readers to avoid blocking writers (and vice versa), how versions are tracked, and provide a simple scenario where two concurrent transactions can each read a valid snapshot yet still produce an anomaly together.
Sample Answer
Direct answer
Multi-version concurrency control, MVCC, is a way a database gives every transaction a consistent view of the data without making readers and writers wait on each other: instead of locking a row so only one transaction can touch it at a time, the engine keeps multiple versions of each row around and hands each transaction the version that was committed as of the moment its own transaction started (its "snapshot"). A reader never blocks a writer because it simply does not look at the writer's uncommitted or newer version, and a writer never blocks on a reader for the same reason: they are working with different versions of the same logical row.
Structured elaboration
- Version tracking: every row carries markers for which transaction created it and, if it has been superseded, which transaction superseded it. Postgres calls these
xmin/xmaxand stores the old and new row physically side by side until cleanup removes the old one; MySQL's InnoDB instead keeps only the current row in place and reconstructs older versions on demand by following a pointer into an undo log. Different storage, same idea: a version is visible to a given transaction if its creator is in that transaction's already-committed set and its (if any) superseder is not. - Snapshot construction: when a transaction running under a snapshot-based isolation level (a rule governing what one transaction can see of another transaction's concurrent work; Postgres's REPEATABLE READ and SERIALIZABLE, and the default REPEATABLE READ in MySQL's InnoDB, are both snapshot-based) starts, the engine computes the set of transaction ids already committed at that instant. Every row version it reads for the rest of that transaction is checked against that fixed set, which is exactly why the same query run twice inside one transaction returns the same answer even if other transactions commit changes in between.
- Why this makes reads and writes non-blocking of each other: a writer creating a new version does not need to wait for any reader holding an older snapshot, because that reader is never going to look at the new version anyway; a reader does not need to wait for a writer's in-progress change to finish, because it is reading a version that already existed and was committed before its snapshot was taken.
Worked example
I ran this on a live Postgres 16 instance.
Non-blocking proof: table acct(id, balance), one row with balance=100, then updated to balance=150. A reader transaction opens at REPEATABLE READ and reads balance=150 (its snapshot's value). While that reader transaction is still open, a second, completely separate transaction updates balance=999 and commits immediately, without waiting on the open reader. The first transaction then reads again, inside the same still-open transaction, after the writer's commit, and still gets 150: the writer was never blocked by the reader, and the reader's answer never changed mid-transaction.
Anomaly proof (an anomaly two individually valid snapshot reads can still produce together): table oncall(doctor, shift_date, is_oncall) with two rows, alice and bob, both is_oncall = true. The rule the application is trying to enforce is "at least one doctor stays on call"; the database itself does not enforce this, the application is expected to check before letting a doctor go off call. Transaction A (REPEATABLE READ) checks count(*) where is_oncall and sees 2, so it proceeds to set alice off call. Transaction B, started while A is still open (its own REPEATABLE READ snapshot taken before A's update), independently checks the same count, also sees 2 on its own snapshot, and proceeds to set bob off call. Both commit. Final state: both alice and bob are off call, an outcome neither transaction would have permitted if it could see what the other was doing. This is write skew: each transaction's individual read was a valid, correct snapshot, and neither write conflicted with the other at the row level (A touched only alice's row, B touched only bob's row), so ordinary snapshot isolation lets both through. Only true serializability (which additionally detects that the two transactions' reads and writes interleave in a way no one-at-a-time ordering could produce) would catch this.
Trade-offs and pitfalls
The most common mistake is assuming "my database uses MVCC, so I do not need to think about isolation levels" for invariants that span more than one row. MVCC solves single-row, single-statement anomalies cleanly (a reader never sees a half-committed write) but does nothing by itself to stop write skew, because write skew is a cross-row, cross-transaction invariant that neither transaction violates by itself. A second pitfall: because MVCC keeps old versions around instead of overwriting in place, something has to reclaim them, or storage and scan cost grow over time. Postgres does this with a background process (autovacuum) that removes versions no transaction can still see; forgetting this exists, or misconfiguring it, is a separate, very real operational cost of choosing an in-place MVCC design.
Explain the difference between B-tree and LSM-tree indexing structures, for example as used in PostgreSQL versus RocksDB. For a high-ingest logging service, which index type would you typically choose, and why?
Sample Answer
Direct answer
A B-tree index mutates data in place: it walks a balanced, sorted tree structure straight to the on-disk page holding a key and updates that page directly, so writes are scattered, random-access page updates. An LSM-tree (Log-Structured Merge-tree) index never mutates data in place: writes go into an in-memory sorted buffer first, then get flushed as new, immutable sorted files that a background process periodically merges. PostgreSQL's default index type is the B-tree; RocksDB is built entirely around an LSM-tree. For a high-ingest logging service, the LSM-tree is the right default, because the workload is almost pure appends and write throughput is the thing under pressure, not point-lookup latency.
Structured elaboration
How each one actually writes a record
B-tree write path: the engine traverses the tree from the root to find the one leaf page where a key belongs, loads that page into the buffer pool if it is not already resident, and mutates it directly (in place). The change becomes durable once the write-ahead log (WAL, the append-only file a database writes an intent-to-change record to before touching the real data) for that change is flushed; the actual data page itself is written back to disk later, in its own time. Every logical write touches one specific, essentially random location in the index file.
LSM-tree write path: the engine appends the change to an in-memory sorted structure (the memtable, commonly a skip list), plus its own WAL entry for durability, and returns. Nothing on disk changes yet. When the memtable fills, its entire contents are written out, unchanged, as one new immutable sorted file (an SSTable, short for "sorted string table"). Later, a background process called compaction merges multiple SSTables together, both to bound how many files a read has to check and to physically drop keys that were since overwritten or deleted.
Why that difference matters for a high-ingest workload
| Property | B-tree (e.g. PostgreSQL default) | LSM-tree (e.g. RocksDB) |
|---|---|---|
| Write pattern | Random-access page mutation | Sequential append (memtable, then sequential SSTable flush) |
| Point read | Direct tree traversal, one predictable path | May check the memtable plus several SSTables (newest first) |
| Read cost control | Not usually needed; tree depth is the cost | A Bloom filter (a compact probabilistic structure that quickly says "definitely not present" for most non-matching keys) is used to skip files that cannot contain the key |
| Background work competing with foreground I/O | Vacuum (a background process that reclaims space left behind by updated or deleted rows) / page write-back, generally lighter | Compaction, which rewrites data multiple times over its lifetime and can compete with live writes for disk bandwidth |
| Best fit | Workloads with frequent updates to existing rows and a need for tight, predictable read latency | Workloads dominated by new writes, tolerant of read paths that touch a few files instead of one |
The recommendation, and what would flip it
For the stated workload (a high-ingest logging service), choose an LSM-tree-backed engine. Log lines are almost never updated after being written, the dominant operation is "append a new record," and the write path never has to do a random page mutation for that, only a sequential append to memory and, eventually, to a new file. A B-tree engine can certainly ingest logs too, but every insert is still a random-access page write under the hood, which is the more expensive pattern for this specific access shape.
This recommendation would flip if the service also needed frequent point lookups or short range scans on recent data under a tight p99 (99th percentile) latency target as a first-class requirement, not just occasional ad hoc queries, because a B-tree's read path is more uniformly fast and does not depend on compaction having kept up. It would also flip if log records needed frequent in-place correction after being written, which is unusual for logs but common for, say, a mutable event status field bolted onto the same table.
Worked example
At 50,000 log lines/sec averaging 300 bytes each, the ingest rate is 50,000 x 300 = 15,000,000 bytes/sec, about 15 MB/s. On the LSM-tree path, that is 15 MB/s of purely sequential writes (memtable, then flush), the pattern disks and SSDs are best at. On the B-tree path, if each insert dirties one index page independent of the others (a reasonable worst case when new rows do not cluster on the same page), that is roughly 50,000 random page writes/sec of demand placed on the storage layer, before accounting for any write-back batching the buffer pool manages to do. The B-tree engine is not incapable of this, modern SSDs and a well-tuned buffer pool absorb a lot of it, but the LSM-tree design converts the same logical write rate into the I/O pattern (sequential) that is cheapest to sustain, which is exactly why it is the standard choice for logging and telemetry systems (RocksDB itself began as a store for exactly this kind of write-heavy workload).
Trade-offs & pitfalls
- "LSM is faster" is only true for writes. A point read on an LSM-tree can be more expensive than on a B-tree, especially for older data buried under several compaction levels; a logging service that also needs fast recent-data dashboards should plan for that, typically by keeping a well-tuned Bloom filter and bounding how many SSTable levels a hot read has to check.
- Compaction is not free background maintenance, it is a second write of the same data. Under sustained heavy ingest, compaction can fall behind, which shows up as growing read latency (more files to check) and, if it falls far enough behind, as write stalls the engine imposes on purpose to let compaction catch up.
- B-tree engines are not one-size-fails-all for logging. If update volume is modest and the team already runs PostgreSQL everywhere, a B-tree-backed table with good partitioning can still work; the LSM-tree recommendation is about the shape of this specific workload, not a blanket rule.
- Several real engines are not purely one or the other. MySQL's InnoDB is B-tree based. WiredTiger, the storage engine MongoDB uses, has an LSM-tree table type in addition to its default B-tree tables, though that LSM mode is legacy and effectively unmaintained, so production MongoDB deployments run on WiredTiger's B-tree path in practice. Naming "PostgreSQL versus RocksDB" is a clean illustration of the two families, not a claim that every engine falls cleanly into one bucket.
That is every published Database Internals and Storage Engines question for Full-Stack Developer so far. Browse the other topics in this category, or practice this one interactively.