System Design Interview

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

What the store promises, and what it will notIt promises• get and put on an opaque key• Single-digit millisecond p99• Tunable consistency per call• Survives a node, and a rackIt does not promise• Joins, or any query by value• Range scans over the key space• One global ordering of writes• A transaction across two keys
Every question in this problem is really asking which of these promises you are willing to buy.

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.

Three versions, each broken on purposev1: hashmap in memoryv2: appendlog plus indexv3:SSTables, compactionStill one machinetopbottomRead bottom to top; each layer repairs the failure below it.
Every distributed component later in the lesson is a repair to a failure this progression makes visible first.
Text
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.

Text
put(k, v)  ->  append (k, v) to data.log ; index[k] = byte offsetget(k)     ->  seek to index[k] ; read the record

Writes 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.

KeysIndex memory
10 million700 MB — fine
100 million7 GB — uncomfortable
1 billion70 GB — needs a very large machine for the index alone
10 billion700 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.

  1. Append the record to a write-ahead log on disk (sequential, for crash recovery).
  2. Insert it into the memtable — a sorted in-memory structure such as a skip list or balanced tree.
  3. 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.
  4. 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-treeLSM tree
Write pathIn-place page update + logSequential append + memtable
Write amplification~10–30×~5–15×
Write throughputModerateHigh
Read latencyLow and predictableHigher and more variable (several SSTables)
Range scansExcellentGood
SpaceFragmentation from splitsTemporary duplication until compaction
Background workMinimalCompaction, which competes with traffic
Typical usersRelational databasesWide-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."

LSM tree — optimised for writesMemtablesorted, in RAMWALappend-onlyL0 SSTables4 filesL1 SSTables40 filesL2 SSTables400 filesflush when fullcompactcompactWRITE: append to WAL, insert into memtable.No disk seek, no read.READ: check memtable, then every level,newest first. A miss touches several files —bloom filters make most of those checks free.B-tree — optimised for readsroot pageinternalinternalleafleafleafREAD: one page read per level. Depth 3–4 covers billions ofrows, so a point lookup is 3–4 seeks — and that is the worstcase, not the average.WRITE: find the page, modify it in place, and split it if it is full.A random write is a random seek.LSM turns random writes into sequential appends and pays for it on reads. B-trees do the opposite. Pick by which side of your workload is heavier.
Neither is faster overall — they move the cost between the write path and the read path.