Course Content
System Design Interview
31 sections · 71 lessons
Building blocks: CAP, consistency models and choosing a database
CAP is the most-quoted and least-understood idea in system design interviews. Quoting it badly costs marks; stating it precisely earns them.
This section is a reference: every later case study links back here instead of re-explaining quorums or retries. This first lesson covers the guarantees a distributed store can and cannot make, and how to pick the store. It starts with CAP and its more useful successor, PACELC. It then walks through the consistency models between "eventual" and "strong", learnt by the symptom each one prevents. It ends with the seven database families and how the access pattern chooses between them.
What CAP actually says
Three properties, defined narrowly:
- Consistency (C): every read returns the most recent completed write, as if there were one copy of the data. This is linearizability. It is not the C in ACID, which means "the database's own invariants hold" — an entirely different idea that happens to share a letter.
- Availability (A): every request to a non-failing node returns a non-error response. Not "the service mostly works" — every node, always, no errors.
- Partition tolerance (P): the system keeps operating when the network drops or delays messages between groups of nodes.
The theorem: when a network partition occurs, a distributed system must give up either consistency or availability. During a partition, a node that cannot reach its peers either answers with data it cannot verify is current (giving up C) or refuses to answer (giving up A). There is no third option.
What CAP does not say
It is not "pick two of three." Partitions happen whether you want them to or not. You do not choose P; the network chooses it for you. The real choice is a binary one that only matters during a partition.
It says nothing about normal operation. Almost all the time your network is fine, and CAP is silent about what you should do then. Since partitions are rare, this means CAP describes a tiny fraction of your system's life.
"CP or AP" is a caricature. Real systems are configurable per operation. A quorum-based store (Partitioning and replication) can be strongly consistent for one read and eventually consistent for the next, decided by a parameter on the call. A relational database with synchronous replication behaves one way; the same database with asynchronous replication behaves another.
PACELC: the more useful frame
PACELC extends CAP with the part that actually matters day to day:
If there is a Partition, choose between Availability and Consistency; Else, choose between Latency and Consistency.
The "else" branch is where your system lives 99.99% of the time. Strong consistency during normal operation is not free: it requires a round trip to a leader or a quorum before a read can be answered. Inside one datacentre that is ~0.5 ms and hardly matters. Across regions it is 100–150 ms per read, which is the entire latency budget of most consumer products.
Two concrete classifications:
- A relational database with a single leader and synchronous replication is PC/EC: it refuses writes during a partition, and it pays latency for consistency normally.
- A leaderless quorum store configured with one-node reads and writes is PA/EL: it stays up during a partition and answers fast, at the cost of possibly stale data.
What to say about CAP in the interview
Do not announce "this is an AP system". Say what happens during a partition, concretely:
"If the replica in the other region becomes unreachable, I'd keep serving reads from the local replica even though they might be a few seconds stale, because the requirement was availability over freshness for the feed. For the payment balance I'd do the opposite — refuse the read rather than serve a number that might be wrong."
That answer contains CAP without needing the word, and it shows the decision is per-operation rather than per-system.
Consistency models, weakest to strongest
CAP only names the two ends. A consistency model is a promise about what a read can return when several copies of the data exist, and there are several useful points in between. The useful way to learn them is by the symptom a user sees when you pick each one.
Eventual consistency. If writes stop, all replicas converge to the same value, in some unspecified amount of time. Nothing else is promised.
Symptom: a user posts a comment and refreshes — the comment is gone. They refresh again and it is back. A like counter reads 42, then 41, then 43, because three reads hit three replicas at different stages of catching up.
Read-your-writes consistency. A user always sees their own writes. Other users may see stale data.
Symptom: fixed for the author, unchanged for everyone else. This is the guarantee that stops the "I uploaded it and it vanished" bug from Database replication in Section 2. Implemented by routing a user's reads to the leader for a short window after they write, or by having the client carry the position of its last write.
Monotonic reads. Once a user has seen a value, later reads never show an older one — time does not run backwards for them.
Symptom without it: the user sees a new message, refreshes, and it disappears, because the second read landed on a further-behind replica. Implemented by pinning a user's reads to one replica (by hashing their identifier) rather than spreading them.
Consistent prefix / causal consistency. If write A happened before write B, nobody sees B without A. Causally related operations are ordered for everyone; unrelated ones may be seen in different orders.
Symptom without it: a chat shows a reply before the message it replies to. This is a genuinely bad user experience and it is why causal consistency is worth the extra machinery in messaging systems.
Strong consistency (linearizability). Every read returns the value of the most recently completed write, as though there were a single copy.
Symptom: the system behaves the way people naively expect. The cost is a coordination round trip on the read or write path, and unavailability during a partition.
The comparison
| Model | User-visible symptom of its absence | Typical cost | Use when |
|---|---|---|---|
| Eventual | Values flicker; own writes vanish briefly | Cheapest; local reads | Counters, feeds, view counts |
| Read-your-writes | "Where did my post go?" | Leader reads for a short window | Anything a user creates and immediately views |
| Monotonic reads | Data appears then disappears | Replica pinning per user | Timelines, message history |
| Causal | Reply before message | Version tracking / dependency metadata | Chat, comments, collaborative editing |
| Strong | Two people see different truths | Quorum or leader round trip; unavailable in a partition | Balances, inventory, unique constraints |
The practical rule
Pick the weakest model whose symptom your users would tolerate, and then say the symptom out loud. "Eventually consistent" is a phrase; "a follower may not see the post for up to two seconds, which is fine for a feed and would not be fine for a bank balance" is an argument.
Note that different fields in the same system can have different models. The news feed of Section 13 is eventually consistent while the same product's account settings are strongly consistent, and that mixture is normal rather than sloppy.
Databases: the seven families
The consistency you need is one input to the next decision: which database. There are seven families worth knowing. Choosing between them is one of the two or three most common places an interviewer probes, and the answer is decided by the access pattern.
Relational (SQL). Rows in tables with a fixed schema, related by keys, queried with Structured Query Language. Supports joins and ACID transactions — Atomicity, Consistency, Isolation, Durability. Strength: relationships and ad-hoc queries you have not thought of yet. Limit: writes concentrate on one leader, and sharding is manual work you take on yourself. The "distributed SQL" family (sometimes called NewSQL) keeps the interface and shards underneath, at the cost of cross-shard transaction latency.
Key-value. get(key) and put(key, value). No querying by value. Strength: a single lookup at very high volume, and trivial horizontal partitioning. Limit: if you need to ask any question other than "what is at this key", it cannot help you. Used for sessions, caches, feature flags, and short-link resolution (Section 10).
Document. Stores self-describing records, usually JSON, with secondary indexes and a flexible schema. Strength: an entity fetched and written as a whole — a product listing, a user profile. Limit: relationships across documents are the application's problem, and denormalised copies drift.
Wide-column. Rows addressed by a partition key plus a clustering key, physically sorted within each partition, built for very high write throughput. Strength: time-ordered data retrieved as "the last N items for this key" — messages, feeds, event history. Limit: you must know your queries before you design the table, because the table is the query.
Graph. Nodes and edges, traversed directly. Strength: multi-hop questions — "friends of friends who live in this city" — which in SQL become a self-join per hop and collapse under depth. Limit: a niche shape; most social products store the graph in a wide-column or relational store and traverse in application code.
Time-series. Append-only, timestamp-primary, heavily compressed, with built-in downsampling. Strength: metrics at enormous volume (Section 22). Limit: it does one thing.
Search. An inverted index mapping terms to documents, with ranking and fuzzy matching. Strength: full-text queries and relevance. Limit: it is a secondary index over a source of truth, not a source of truth — it is usually eventually consistent with the database it indexes.
The decision table, keyed on access pattern
| Access pattern | Choose | Because |
|---|---|---|
| Entities with relationships, unpredictable queries, transactions | Relational | Joins and ACID |
| One key in, one value out, at very high volume | Key-value | O(1) lookup, trivially partitioned |
| Whole documents, flexible fields, few relationships | Document | Read and write the aggregate as one unit |
| Huge write volume, always read as "last N for this key" | Wide-column | Sorted partitions, write-optimised engine |
| Multi-hop relationship traversal | Graph | Traversal without join explosion |
| Timestamped numeric series, aggregated over ranges | Time-series | Compression and downsampling built in |
| Free-text search with ranking | Search index | Inverted index and relevance scoring |
The sentence to say
Never say "I'd use MongoDB" and stop. Say the pattern, the choice, the alternative, and the condition:
"The access pattern is 'give me the last fifty messages in this conversation', always by conversation and always in time order, at about 40,000 writes per second. That's a wide-column store — partition by conversation, cluster by message identifier. I'd consider a relational store instead if the volume were an order of magnitude lower, because then I'd rather have transactions and joins than write throughput."
Three seconds longer, and it demonstrates the axis Section 1 said decides level.