Interview Setup
Interview Prompt
Design the storage engine of a single node relational database (PostgreSQL style) handling 50K TPS on a 1 TB dataset with 256 GB buffer pool, guaranteeing ACID semantics and sub 5ms commit latency on NVMe.
Clarifying Questions (ask before designing)
| Question | Why it matters |
|---|---|
| Are we designing a single node engine or a distributed database? | A single node keeps the discussion on page management, WAL, MVCC, indexing, and local concurrency. A distributed design adds replication and consensus as separate concerns. |
| What is the primary workload? | OLTP favors row oriented storage and B trees for point and range access. OLAP workloads often favor columnar storage and different execution strategies. |
| Do we need Serializable isolation, or a stronger real time consistency guarantee? | This determines whether Read Committed, Repeatable Read, or Serializable isolation is sufficient. PostgreSQL Serializable uses Serializable Snapshot Isolation, which is not identical to strict serializability. |
Scope
In scope
- Storage engine and page layouts
- Write-Ahead Logging (WAL) for durability
- Concurrency control (MVCC)
- Index structures (B tree)
- Vacuum, tuple visibility, and transaction ID aging
- Heap maintenance structures: free space map, visibility map, HOT, and TOAST
Out of scope (state explicitly)
- Detailed SQL grammar and advanced query optimizer implementation
- Implementing the PostgreSQL wire protocol itself beyond a representative message flow
- Implementing distributed consensus itself, since the primary design is single node
Functional Requirements
Start by asking your interviewer whether you are designing a single node storage engine or a distributed database, because this problem focuses on the former. Confirm whether the workload is OLTP or OLAP and whether the consistency requirement is Read Committed, Repeatable Read, Serializable, or something stronger. In PostgreSQL, Serializable is implemented with Serializable Snapshot Isolation rather than strict serializability.
- Relational Model: Store data in tables with rows and typed columns.
- SQL Support: Support standard SQL queries (SELECT, INSERT, UPDATE, DELETE, JOINs, aggregations).
- ACID Transactions: Atomicity, Consistency, Isolation, and Durability must be guaranteed.
- Indexing: Support B tree indexes to speed up point queries and range scans.
- Concurrency Control: Allow multiple clients to read and write simultaneously without locking the entire database (MVCC).
Non-Functional Requirements
Your interviewer will stress test durability through WAL before data page flushes and will probe whether MVCC lets readers proceed while writers update other versions. They will also ask what happens when dead tuples accumulate, where VACUUM, autovacuum, visibility, and transaction ID wraparound become important staff level follow ups.
- Durability (Zero Data Loss): Committed transactions must survive power failures and crashes.
- High Performance: Must leverage in memory caching through the Buffer Pool and sequential WAL I/O to maximize throughput.
- Scalability: Vertical scaling is the primary model, though read replicas should be supported for scaling read traffic. Replica reads are potentially stale and read after write paths that require freshness should route to the primary or use LSN aware routing.
- Crash Recovery: The system must recover to a consistent state rapidly after an unexpected shutdown.
Capacity Estimations
Run page count and buffer pool math before sizing WAL throughput. At 50K TPS, group commit and checkpoint spreading determine whether the commit p99 can stay within the 5 ms target.
| Metric | Calculation | Value |
|---|---|---|
| Database size | Given (production instance) | 1 TB |
| 8 KB pages | 1 TB ÷ 8 KB | ~125M pages |
| OLTP read/write mix | Given | 70% read / 30% write |
| Peak transactions/sec | Given | 50K TPS |
| Buffer pool (shared_buffers) | Given | 256 GB (~25% of data hot) |
| WAL write rate (peak) | 50K TPS x ~500 B WAL/transaction | ~25 MB/s |
| Client connections (via PgBouncer) | 10K microservices x 50 application connections = 500K client side connections, while PgBouncer caps active PostgreSQL backends at 500 | 500 pooled backends |
Using the stated planning approximation of 125M pages, a 256 GB buffer pool holds roughly 25% of the 1 TB dataset at 8 KB per page. The page count is a planning approximation because PostgreSQL uses an 8 KiB page while storage capacity labels may use decimal or binary units. A point lookup through a B tree with depth 3 to 4 touches 3 to 4 pages. A cache hit requires roughly 0.1ms, while a miss takes 1 to 5ms on NVMe storage. At 50K TPS, if group commit batches ~10 to 50 transactions per WAL fsync, the design can target commit p99 under 5ms.
Architecture Diagram
PostgreSQL style storage combines a SQL execution layer with a transactional storage engine that provides ACID semantics through WAL and MVCC. At 50K TPS on a 1 TB dataset with only 25% of the data fitting in the 256 GB buffer pool, the design relies on sequential WAL writes, buffered page access, and snapshot based concurrency rather than read locks on every row.
The diagram traces a query from the PostgreSQL wire protocol through query execution, the buffer pool, B tree indexes, and the WAL based durability path used for crash recovery.
In the room
Clarify single node engine vs distributed database architecture such as Citus or CockroachDB in the first two minutes. Then walk through WAL append and fsync on COMMIT before the corresponding dirty page reaches durable storage. That durability story separates a real database from a key value store with SQL on top.
Component Deep Dives
1. Query Execution Pipeline
Interview depth on PostgreSQL internals centers on four mechanisms: how queries become execution plans, how the buffer pool amortizes disk I/O, how WAL makes commits durable without flushing every dirty page synchronously, and how MVCC lets readers and writers coexist without blocking each other.
Every query passes through parsing, semantic analysis, planning, and execution. EXPLAIN ANALYZE is how you verify whether the chosen plan actually uses the expected index or falls back to a sequential scan.
When a SQL query arrives, the execution pipeline works through the following stages:
- Parser: Verifies syntax and generates a Parse Tree.
- Analyzer: Checks semantics (verifying whether tables and columns exist and validating permissions).
- Optimizer (Cost-Based): Generates multiple execution plans (such as Index Scan vs Sequential Scan, or Hash Join vs Nested Loop Join) and estimates the I/O and CPU cost of each using statistical distributions. It picks the cheapest plan.
- Executor: Processes the plan node-by-node (Volcano iterator model), pulling tuples from the storage engine.
EXPLAIN ANALYZE SELECT * FROM users WHERE age > 30 ORDER BY name LIMIT 10;
Limit (cost=0.42..1.53 rows=10 width=45) (actual time=0.045..0.055 rows=10 loops=1)
-> Index Scan using users_name_idx on users (cost=0.42..111.45 rows=1000 width=45)
Filter: (age > 30)
Rows Removed by Filter: 15
Planning Time: 0.150 ms
Execution Time: 0.080 ms2. Buffer Pool (Shared Memory)
Disk I/O is the primary bottleneck. The buffer pool caches 8 KB database pages in shared memory, with a hit ratio above 95% as the operational target. PostgreSQL uses a clock sweep style buffer replacement strategy rather than a classic LRU implementation.
Disk I/O is slow. PostgreSQL allocates a large chunk of RAM called the shared_buffers. The database reads data from storage in fixed size blocks, usually 8 KB, called pages.
- When a query requests a row, the storage engine checks whether the page containing that row is already in the Buffer Pool.
- If it is (Cache Hit), it's returned immediately.
- If not, the page is loaded from storage into the Buffer Pool, and an older page may be evicted according to PostgreSQL's clock sweep based replacement policy.
- Modifications are made in memory first. The page becomes "dirty".
3. Write Ahead Log (WAL) and Durability
WAL is the durability mechanism. PostgreSQL flushes the required WAL records at commit, while dirty data pages are written later by checkpoint and background write activity. The key rule is WAL before the corresponding data page. For downstream CDC ingestion, see the Change Data Capture (CDC) Pipeline.
If the database crashes before dirty pages are flushed to disk, those in memory data pages can be lost, but committed changes remain recoverable from durable WAL. Flushing 8 KB data pages for every transaction would create excessive random I/O. The solution is the WAL, which lets the commit path persist a much smaller sequential log record while data pages are flushed later.
- Before modifying a page in the Buffer Pool, a small log entry describing the change is appended to the WAL in memory.
- On
COMMIT, the required WAL records are flushed to durable storage before the commit is acknowledged. Group commit can let multiple transactions share one flush. - Dirty data pages are written during checkpoint processing, while the
bgwriteralso writes pages proactively to smooth future write pressure. - Rule: The required WAL records must reach durable storage before the corresponding dirty data page is written to durable storage.
- Checkpoint: A checkpoint flushes dirty data pages and records a redo starting point. A checkpoint is not a transaction commit and does not itself determine whether an individual transaction is committed.
- Full page writes: With full page writes enabled, the first modification to a data page after a checkpoint can log the full page image in WAL to protect against torn page writes. This increases WAL volume after checkpoints.
LSN 0/1A2B3C: Transaction ID: 5092 Resource Manager: Heap Action: INSERT Relation: users (filenode: 16384) Block: 42 Offset: 12 Tuple Data: (id=5, name='Bob', age=35)
4. Multi Version Concurrency Control (MVCC)
MVCC lets PostgreSQL readers use transaction snapshots while updates create new tuple versions. Ordinary reads therefore do not need row read locks that conflict with writers, although write write conflicts can still block or abort transactions.
MVCC is designed so ordinary reads can proceed without acquiring row locks that conflict with concurrent writers. An UPDATE creates a new tuple version, allowing readers to continue using an appropriate transaction snapshot.
- Each heap tuple carries MVCC metadata including
xminandxmax, which identify the inserting and deleting or updating transaction IDs. - Tuple visibility is evaluated against the transaction snapshot, including transaction IDs that were committed, aborted, or still in progress. It is more precise than a simple
xmin/xmaxcomparison. - Transaction status and hint bits: PostgreSQL uses transaction status data to determine whether relevant XIDs committed or aborted, and hint bits can cache visibility information in tuple headers so repeated checks do not always require another status lookup.
- Repeatable Read provides a stable snapshot for the transaction, while Serializable adds PostgreSQL's Serializable Snapshot Isolation checks to detect dangerous dependency patterns.
- Vacuum Horizon: VACUUM cannot remove tuple versions that might still be visible to an active transaction. Long running transactions can therefore hold back cleanup and increase table bloat.
- Autovacuum: PostgreSQL normally schedules background vacuum and analyze work based on table activity. The goal is to reclaim reusable space and keep planner statistics current before bloat or transaction ID aging becomes dangerous.
5. Indexes (B tree and Hash)
B tree indexes reduce point lookup search from O(N) scanning to O(log N) tree navigation. A tree depth of 3 to 4 at billion row scale means only a few index levels must be traversed, although some pages may already be cached and an index scan may also need a heap page fetch.
Indexes reduce point lookup search from O(N) scanning to O(log N) tree navigation. PostgreSQL's B tree access method is the general purpose ordered index structure.
- Leaf entries identify heap tuples through tuple identifiers unless the index can satisfy the query as an index only scan.
- Because the tree is balanced and highly branched, even very large tables can have only a few index levels. The exact number of storage reads depends on cache residency and whether the query can use an index only scan.
- PostgreSQL also supports Hash indexes, GiST for extensible search such as geospatial workloads, and GIN for inverted indexes such as full text and JSONB containment.
- Index Only Scans: A B tree query can avoid heap fetches when all required columns are available from the index and the visibility map confirms that the relevant heap pages are all visible.
6. Connection Pooling (PgBouncer)
A dedicated PostgreSQL backend process per client connection becomes expensive when connection counts grow into the thousands. PgBouncer reduces the number of active PostgreSQL backends by multiplexing many client connections through a smaller server side pool.
PostgreSQL normally uses one backend process per client connection. The ~10MB per connection figure is a planning estimate, not a protocol guarantee, and the actual memory cost varies with workload and session state. At scale, thousands of client connections can exhaust memory and process capacity.PgBouncer is used as a lightweight connection pooler sitting in front of PostgreSQL, multiplexing thousands of client connections onto a small pool of actual database connections.
7. High Availability and Replication
Single node PostgreSQL scales primarily through CPU, RAM, and storage improvements. Streaming replication provides standby copies, while a failover manager such as Patroni coordinates promotion without changing the SQL model.
To survive a catastrophic node failure, PostgreSQL uses Replication.
- Physical Replication (Streaming): The primary streams WAL records to a standby, and the standby replays that WAL to maintain a physically equivalent database state.
- Logical Replication: Decodes WAL into logical row changes for selected publications and subscriptions. It is useful for selective replication and many migration paths, including major version upgrades with minimal downtime.
- Automated Failover: Tools such as Patroni use a distributed configuration store such as etcd, Consul, or ZooKeeper for leader election and cluster state. If the primary fails, the selected standby can be promoted and the routing layer can be updated.
- Replica Lag: Asynchronous standbys can lag behind the primary. Applications that require read after write consistency should route those reads to the primary or wait for a replica to replay the required WAL position.
8. Heap Maintenance: FSM, Visibility Map, HOT, and TOAST
These structures explain several behaviors that appear in production once update rates and row sizes increase. They are the bridge between the logical MVCC model and the physical heap layout.
- Free Space Map (FSM): Tracks approximate free space available on relation pages so INSERT and update placement can find candidate pages efficiently.
- Visibility Map: Maintains all visible and all frozen information for heap pages. The all visible bit can let an index only scan avoid heap visibility checks, while VACUUM uses the map to skip work that is already known to be clean and to track pages suitable for freezing related maintenance.
- HOT Update: When indexed columns are unchanged and the page has enough free space, PostgreSQL can place the new tuple version on the same heap page without adding new index entries. This reduces index maintenance and write amplification but depends on page space and update patterns.
- TOAST: Oversized variable length values can be compressed and stored out of line in a related TOAST table as chunked values, with the main heap row retaining the reference to the external data.
API Design
Message Flow (Query)
PostgreSQL does not expose its database engine through REST. Clients normally use the stateful binary TCP based PostgreSQL Wire Protocol, which carries query, row description, data row, completion, and transaction state messages.
// Frontend (Client) -> Backend (PostgreSQL)
// Query Message ('Q')
Q: "SELECT id, name FROM users WHERE age > 30;"
// Backend -> Frontend
// Row Description ('T')
T: [Field1: "id" (Int4), Field2: "name" (VarChar)]
// Data Row ('D')
D: [5, "Bob"]
// Data Row ('D')
D: [12, "Alice"]
// Command Complete ('C')
C: "SELECT 2"
// Ready For Query ('Z')
Z: (Transaction Status: Idle)This stateful connection model is why PgBouncer is useful at scale. In session pooling mode it keeps backend sessions attached for the client session, while transaction or statement pooling can multiplex more aggressively when the application does not depend on session state.
Data Model
PostgreSQL stores tables and indexes as fixed size database pages, usually 8 KB in a standard PostgreSQL build. The database page size is an engine format choice and should not be confused with the underlying filesystem or device block size.
// Simplified 8 KB page layout in PostgreSQL
struct PageHeaderData {
uint64 pd_lsn; // LSN associated with the latest WAL change to the page
uint16 pd_checksum; // Detects corruption when page checksums are enabled
uint16 pd_flags; // Flag bits
uint16 pd_lower; // Offset to start of free space
uint16 pd_upper; // Offset to end of free space
uint16 pd_special; // Offset to special space used by index access methods
uint16 pd_pagesize_version;
uint32 pd_prune_xid; // Hint for the oldest XMAX that may need pruning
};
struct ItemIdData { // Line pointer metadata
unsigned lp_off:15; // Offset to tuple data
unsigned lp_flags:2; // State of tuple
unsigned lp_len:15; // Tuple length
};
// Heap tuple data grows backward from the free space region.
struct HeapTupleHeaderData {
TransactionId t_xmin; // Insert XID
TransactionId t_xmax; // Delete or update XID
union {
CommandId t_cid; // Insert / delete command ID
TransactionId t_xvac; // XID used when VACUUM moves a row version
} t_field3;
ItemPointerData t_ctid; // Current TID of this or a newer tuple version
uint16 t_infomask2; // Attribute count and flag bits
uint16 t_infomask; // Tuple status and visibility flags
uint8 t_hoff; // Offset to user data
// A null bitmap and user attributes follow the fixed header as needed.
};- Line Pointers (ItemIds): Grow forward from the page header and point to tuple locations through offset and length metadata.
- Tuple Data: Grows backward from the free space region. The space between item identifiers and tuple data is available for new allocations.
- MVCC Coherence: MVCC metadata is stored in the heap tuple header, so PostgreSQL can determine tuple visibility from tuple metadata plus the transaction snapshot. PostgreSQL traditionally uses heap tuple versions rather than a separate general purpose undo log for ordinary table updates.
- Free Space and Visibility Maps: PostgreSQL maintains separate relation level maps for approximate free space and page visibility. The free space map helps find pages with room for new tuples, while the visibility map records pages whose tuples are known to be visible to all transactions and can also help index only scans.
- HOT and TOAST: HOT updates can avoid extra index entries when updates stay on the same page and indexed columns do not change. TOAST moves oversized variable length values out of the main heap tuple into a related storage structure when necessary.
Fault Tolerance
Recovery Process (WAL based, ARIES inspired)
If the server crashes, in memory dirty pages are lost. That is safe for committed transactions because the required WAL records were already flushed before commit was acknowledged.
- Upon startup, PostgreSQL locates the latest valid checkpoint. At checkpoint time, required dirty data pages have been flushed and a checkpoint record identifies the redo starting point. Older WAL can become recyclable when no archiving, replication, or other retention requirement still needs it.
- It begins reading the WAL from the checkpoint redo point forward and replays the required records.
- REDO Phase: WAL records are replayed from the required redo point so data files reach a consistent state after the crash. This is roll forward recovery, not a replay of every SQL statement.
- Transactions that were in progress and had not committed are not treated as committed during visibility checks. Their tuple versions remain unreachable to normal transactions and are reclaimed later by VACUUM.
Routine maintenance is handled by autovacuum, which vacuums and analyzes tables based on activity. Normal VACUUM makes dead tuple space reusable inside the relation, while VACUUM FULL rewrites the table and requires an ACCESS EXCLUSIVE lock. Continuous WAL archiving can additionally support point in time recovery. That differs from standby failover because recovery can restore the database to a chosen historical target.
Additional Considerations
Interview Walkthrough
- 25 minute cut
Skip deep architectural variants unless targeting staff level.
- Single node engine vs distributed database architecture (5 min)
- Slotted page layout and buffer pool hit or miss mechanics (6 min)
- Durability lifecycle: WAL append and fsync on commit (5 min)
- MVCC snapshot isolation and VACUUM tuple reclamation (5 min)
- B tree indexing with depth 3 to 4 at billion row scale (4 min)
- Clarify single node engine vs distributed database architectures such as Citus or CockroachDB, confirming this problem focuses on the single node engine.
- Walk 8 KB page layout: line pointers forward, tuples backward, free space in the middle.
- Durability story: append WAL records, flush the required WAL on COMMIT before the corresponding dirty page reaches durable storage, then replay WAL from the checkpoint redo point after a crash.
- MVCC: UPDATE creates a new tuple version, readers evaluate visibility against a snapshot, and VACUUM reclaims dead versions once no active transaction can see them.
- Index: B tree search is O(log N) with a few levels at billion row scale, and EXPLAIN ANALYZE verifies whether the planner actually chooses the index.
- Scale client connections with PgBouncer, because the one backend process per client connection model does not survive 10,000 microservices.
- Common pitfall: long running transactions hold back the MVCC horizon, prevent dead tuples from being reclaimed, and can drive table bloat and worsening scan performance.
Engineering Trade-offs
PostgreSQL MVCC vs. MySQL InnoDB MVCC
PostgreSQL keeps old row versions in the heap, so frequent UPDATE and DELETE activity creates dead tuples that autovacuum must reclaim. Aborting a transaction does not rewrite every affected row to an earlier physical state. The tuple versions remain invisible according to MVCC and are cleaned up later. MySQL InnoDB keeps older row versions in undo records, giving it a different storage and history management model. Long running transactions can keep both systems from reclaiming old versions promptly.
B tree vs. LSM Tree (LevelDB/RocksDB)
B trees provide efficient point and range reads, but random insert workloads can cause page splits and random writes. LSM trees, used by systems such as Cassandra and RocksDB, reduce random write pressure by buffering writes and flushing immutable sorted files, but compaction and searches across levels can increase read and write amplification.
Synchronous vs. Asynchronous Replication
Async Replication: The primary can acknowledge a commit before the standby has received the WAL. This gives lower latency but a failover can lose transactions that were committed only on the primary.
Sync Replication: The commit waits for the configured synchronous standby to confirm the required WAL has reached durable storage. This can provide RPO=0 for the configured failure model, but adds at least a network round trip and can reduce availability if the synchronous standby is unavailable.
Vertical Scaling vs. Sharding
Relational databases can scale vertically very well, but eventually CPU, RAM, storage bandwidth, or single node write capacity becomes the limiting factor. Sharding splits data across PostgreSQL nodes, for example with Citus. Sharding can increase aggregate capacity, but cross shard joins, global constraints, and distributed transactions become more complex. See Distributed Transactions: 2PC vs Saga.
Review
How helpful was this walkthrough?
Click a star to rate. We actively use this feedback to refine and update our system design content.
Discussion
Share your thoughts, ask questions, or help others.