Course Content
System Design Interview
31 sections · 71 lessons
Key-Value Store: requirements, a single node and storage engines
The prompt: "Design a distributed key-value store."
Stop here and attempt it for 45 minutes on paper first. This is the densest section in the course — half of distributed systems theory arrives in it — and attempting it badly first is what makes the rest stick.
This section is three lessons. This first one scopes the problem and builds a single node properly: the simplest thing that works, broken deliberately, until the index outgrows memory and forces a real on-disk storage engine. The second distributes the data and makes consistency a dial. The third handles failure and walks a read and a write through everything end to end.
What a key-value store is
A store supporting two operations: put(key, value) and get(key). No queries by value, no joins, no schema. That narrow interface is what makes it possible to distribute the data across hundreds of machines, because any key can live anywhere and nothing needs to be co-located.
You have already used several. A cache is a key-value store without durability. A session store is a key-value store with a time-to-live. The URL shortener of Section 10 is a key-value store with two endpoints in front of it.
The five questions
1. How big are keys and values? This changes everything downstream.
- Keys under 100 bytes, values under 10 KB: everything in this section applies directly.
- Values of many megabytes: this is object storage, not a key-value store — the design splits metadata from data and chunks the payload (Section 26).
- Values under 100 bytes at enormous volume: per-key overhead starts to dominate storage, and the design shifts towards packing.
Assume keys under 100 bytes, values under 10 KB unless told otherwise, and say so.
2. Read-heavy or write-heavy? This decides the storage engine (compared later in this lesson) more than any other input. A 100:1 read-heavy workload favours a B-tree; a write-heavy ingestion workload favours a log-structured merge tree. Ask for the ratio; if the interviewer will not give one, state your assumption.
3. Single datacentre or global? Within one datacentre, replication costs 0.5 ms and strong consistency is affordable. Across continents it costs 100–150 ms per round trip, and synchronous cross-region writes stop being viable for anything interactive. This question decides whether the consistency conversation is easy or hard.
4. Which consistency guarantee? Strong, or eventual with tunable quorums? Ask for the user-visible requirement rather than the label: "If two clients write the same key at the same moment, what should the next reader see? And is it acceptable for a read to occasionally return a value that is a second old?" Consistency models, in Section 5, has the vocabulary.
5. Durability requirements. Is losing the last few milliseconds of acknowledged writes acceptable on a node crash? The answer decides whether a write is acknowledged after it is in memory, after it is in a write-ahead log, or after that log has been forced to disk — a difference of roughly three orders of magnitude in write latency.
Two more if there is time
Availability target. Must writes succeed while a node or a whole zone is down? That sets the replication factor and the write quorum.
Scale. Total data size and operations per second. Say something like 10 TB and 100,000 operations per second so the arithmetic has anchors.
The design goals we will build against
Assume the common brief: small keys and values, mixed read/write, multi-datacentre, tunable consistency, high availability, and the ability to grow by adding machines. That is the Dynamo-style shape, and the rest of the section builds it from a hash map upward.
Version 1: a hash map in memory
Build the simplest thing that works, then break it deliberately. Every component later in the section is a repair to a specific failure of this progression.
store = {} # key -> valuedef put(k, v): store[k] = vdef get(k): return store.get(k)This is a real key-value store. Both operations are O(1), and a single machine handles hundreds of thousands of them per second because nothing touches disk.
How far it goes. 1 million keys with 1 KB values is 1 GB — fine on any machine. 10 million keys is 10 GB — fine on a large one. This is not a toy: an in-memory store like this, with a network protocol on top, is what a cache is.
What breaks it: the process restarts and everything is gone.
Version 2: an append-only log with an in-memory index
Do not write the hash map to disk — write every operation to the end of a file as it happens.
put(k, v) -> append (k, v) to data.log ; index[k] = byte offsetget(k) -> seek to index[k] ; read the recordWrites are sequential appends, the fastest thing a disk does — on the order of hundreds of megabytes per second even on spinning disks, and gigabytes per second on NVMe SSDs. Reads are one seek. On restart, replay the log to rebuild the index.
Two immediate problems, both solvable:
- The log grows forever. Overwriting a key appends a new record and leaves the old one. Fix: periodically compact by rewriting the file with only the newest record per key. This idea returns below as compaction in the LSM tree.
- A crash mid-write leaves a partial record. Fix: a checksum per record, and discard a trailing record that fails it.
What breaks it: the index is still entirely in memory.
Version 3: the memory limit that forces a real engine
Do the arithmetic. Each index entry needs the key (say 32 bytes), a byte offset (8 bytes), and hash-table overhead (say 30 bytes) — call it 70 bytes per key.
| Keys | Index memory |
|---|---|
| 10 million | 700 MB — fine |
| 100 million | 7 GB — uncomfortable |
| 1 billion | 70 GB — needs a very large machine for the index alone |
| 10 billion | 700 GB — not possible on one machine |
The values are not the problem; they live on disk. The index is the problem. At a billion keys you are buying an enormous machine to hold a lookup table, and at ten billion you cannot buy one at all.
So the index itself must live on disk, in a structure that can be searched without loading it all. That structure is either a B-tree or a log-structured merge tree, and choosing between them is the rest of this lesson.
And then: one machine is still one machine
Even with an on-disk index, a single node has a bounded disk, bounded throughput, and a 100% chance of eventually failing. At 10 TB of data and 100,000 operations per second, you need many machines — which forces partitioning (Distributing the data), which forces replication, which forces every problem in the two lessons after this one: quorums, conflicts and failure handling.
Storage engines: the B-tree
There are two ways to keep a searchable index on disk. Nearly every database in existence uses one or the other, and the difference is a direct trade between write throughput and read latency.
A B-tree is a balanced tree of fixed-size pages (commonly 4–16 KB), each holding sorted keys and pointers to child pages. It is the classic relational database index, and it updates data in place.
Write path. Find the page containing the key by walking the tree from the root — three or four page reads for a large index. Modify the page in memory. Write the change to a write-ahead log for crash recovery, then eventually flush the modified page to its original location on disk. If the page is full, split it into two, which also modifies the parent.
Read path. Walk the tree: root, internal nodes, leaf. Three or four page reads, of which the upper levels are almost always cached in memory, so in practice one disk read finds any key. Latency is low and — importantly — predictable.
Cost. Every small write turns into at least one full-page write plus the log entry. Writing a 100-byte value can mean writing 8 KB. That ratio is write amplification, and for B-trees it is typically in the region of 10–30×, depending on page size and how random the keys are. The writes are also scattered across the disk rather than sequential.
Strengths: predictable low read latency, efficient range scans (leaves are linked in key order), and mature transaction support.
The log-structured merge (LSM) tree
Never update in place. Buffer writes in memory, flush them out in sorted batches, and merge those batches in the background.
Write path.
- Append the record to a write-ahead log on disk (sequential, for crash recovery).
- Insert it into the memtable — a sorted in-memory structure such as a skip list or balanced tree.
- When the memtable reaches a size threshold (say 64 MB), it becomes immutable, a new one takes over, and the full one is written to disk in one sequential pass as an SSTable (sorted string table): a file of key-value pairs in sorted order, with a small index at the end.
- Compaction runs in the background, merging SSTables together, discarding overwritten and deleted keys.
A write therefore costs an append and a memory insert. There is no random disk I/O on the write path at all.
Read path. Check the memtable. Then check each SSTable, newest first, because a key can exist in several with different versions. That is potentially many file lookups for one get — read amplification. Two mechanisms make it tolerable:
- Bloom filter per SSTable: a compact probabilistic structure that answers "is this key definitely absent?" A negative answer is certain, so the file is skipped without being read. At roughly 10 bits per key it gives about a 1% false-positive rate, so ~99% of unnecessary file reads are avoided.
- Sparse index per SSTable: one index entry per block rather than per key, kept in memory, narrowing a lookup to one block read.
Cost. Compaction rewrites data repeatedly — typical write amplification is in the region of 5–15× depending on the compaction strategy — and it consumes disk and CPU that compete with live traffic. The characteristic production problem is a compaction storm: a burst of writes triggers heavy compaction, which slows reads, which is exactly when you did not want it.
Strengths: very high write throughput, sequential I/O throughout, and good compression (sorted data compresses well). Also cheap deletes — a delete is a tombstone record, a marker that removes the key when compaction eventually processes it.
Choosing an engine
| B-tree | LSM tree | |
|---|---|---|
| Write path | In-place page update + log | Sequential append + memtable |
| Write amplification | ~10–30× | ~5–15× |
| Write throughput | Moderate | High |
| Read latency | Low and predictable | Higher and more variable (several SSTables) |
| Range scans | Excellent | Good |
| Space | Fragmentation from splits | Temporary duplication until compaction |
| Background work | Minimal | Compaction, which competes with traffic |
| Typical users | Relational databases | Wide-column and key-value stores |
For this design: LSM tree, because a distributed key-value store built for scale is absorbing high write volume and single-key reads, and the Bloom filters make the read amplification manageable. State the condition that would flip it: "If this were read-dominated with strict tail-latency requirements and moderate write volume, a B-tree would give more predictable reads and I would prefer it."