In a planning meeting someone says, "we need a service that stores chat history." You walk to the whiteboard and draw a box for the app, a box for the database, and an arrow between them. That's enough for a first demo.
Then the whiteboard starts offering more. Something to spread requests over several servers would be sensible, and so would something to remember recent answers, something to hold work for later, and a database split across several machines. Nothing in the sentence says which of them you need, because the sentence has no numbers in it. A chat service for fifty people and one for fifty million have the same one-line description, and they are completely different designs: one runs on a single small database server with room to spare, the other collects more data every year than one machine can hold. An engineer sizing a bridge counts the cars first and draws the arches afterwards, for the same reason.
So most of system design is arithmetic done before drawing. You guess how many people use the service and what they do, turn the guesses into requests per second and bytes per year, compare those with what one machine can handle, and add a box only when a number is bigger than that machine's limit. This chapter does it once, for the chat service: from the sentence to numbers, from the numbers to a check against what one machine can do, and from the check to the few components that earn their place. One question runs through it: how do you know which boxes to draw?
01Turn one sentence into numbers
1.1Four guesses
The sentence says what the service does and nothing about how much of it there will be, so the first job is to guess the amounts, out loud, and see what they imply. Four guesses are enough to start.
The first is how many people use it. We'll count daily active users, the number of different people who use the service on an ordinary day, and guess ten million. The second is how many messages each of them sends in a day: say forty. The third is how many bytes one stored message takes, counting its text plus the metadata kept with it, such as the sender, the time and the conversation: say two hundred. Multiplying the first two gives messages per day. Dividing that by the 86,400 seconds in a day gives messages per second, which is the write rate, because storing each new message is one write to the database.
Dividing by the whole day treats every second alike, but people aren't awake evenly, and a system has to survive its busiest moments and not just its average ones. So the fourth guess is the peak factor, how many times busier the busiest period is than the average. For a consumer chat app with an evening rush, three is a reasonable guess. Multiply the average write rate by it and you get the peak write rate, the one we'll compare with capacity.
Storage is the other number we want. Messages per day times bytes per message times 365 gives the bytes added in a year. We'll count a single copy for now. Later we'll add replicas, extra copies of the data on other machines, kept so that one dead disk doesn't lose anyone's messages.
1.2Trying it
The script below does exactly that sum. Each guess is a variable at the top with a comment marking it as an assumption, so anyone reading it can change one and rerun. 10_000_000 is Python's way of writing ten million, with _ marks that only help the eye, the :, inside the f-strings prints thousands separators, and dividing by 1e9, a billion, turns bytes into gigabytes. The qps in two variable names stands for queries per second, here the number of writes per second. Save it as est.py and run python3 est.py.
users_per_day = 10_000_000 # assumption: daily active users
msgs_per_user = 40 # assumption: messages sent per user per day
bytes_per_msg = 200 # assumption: text plus metadata
peak_factor = 3 # assumption: busiest second vs the daily average
msgs_per_day = users_per_day * msgs_per_user
avg_qps = msgs_per_day / 86_400
peak_qps = avg_qps * peak_factor
per_year_gb = msgs_per_day * bytes_per_msg * 365 / 1e9
print(f"messages per day : {msgs_per_day:,}")
print(f"average writes/s : {avg_qps:,.0f}")
print(f"peak writes/s : {peak_qps:,.0f}")
print(f"storage per year : {per_year_gb:,.0f} GB (~{per_year_gb / 1000:.1f} TB, before replicas)")messages per day : 400,000,000
average writes/s : 4,630
peak writes/s : 13,889
storage per year : 29,200 GB (~29.2 TB, before replicas)Read the output from the top. Ten million users sending forty messages each is 400 million messages a day. Spread over the day that's 4,630 writes a second, and at three times the average the peak is 13,889. The history grows by 29,200 GB a year, about 29.2 TB, for one copy.
1.3What the numbers already decide
Those four figures already rule some designs in and out, before anything is drawn. A few thousand writes a second is an ordinary load for one database server, and a few million is not (section 5 puts real numbers on that line), so we already know which side of it we're on. The 29 TB added every year keeps growing, so sooner or later it outgrows any single machine. Neither conclusion needed a diagram. Each is a comparison between a number we computed and a limit, and each says whether one machine is enough or something has to be added.
Hold on to that shape, because the rest of the chapter repeats it: a number, a limit, and a component added only when the number is above the limit.
1.4Say the assumptions out loud
You'll rarely be given every number, which is why we made guesses. A guess is fine as long as it's written down next to the estimate. "Assuming a 3× evening peak" is a design input that someone in a review can correct, while a silent assumption is a bug nobody can find.
Four guesses gave us load and storage. They said nothing about how fast a conversation must open, or whether a message may ever be lost, and those change a design as much as volume does.
02Requirements that carry a number
2.1Two kinds of requirement
How fast a conversation must open, and whether a message may be lost, need numbers too. Requirements come in two kinds. Functional ones say what the system does: send a message, load a conversation, search history. Non-functional ones say how well it does it: how many users, how fast, how durable. Two services with the same functional list and different non-functional answers are different designs, so the second kind is the one to pin down.
2.2The questions that change the design
A few words first, because the questions use them. Latency is how long one request takes, from the moment it's sent to the moment the answer arrives. An average latency hides the slow requests, so targets are usually given as a percentile: p99 is the latency that 99 out of every 100 requests beat, which means the slowest one in a hundred takes longer. A system is durable if a message it has acknowledged is still there after a crash, which in practice means it reached a disk (Chapter 08 explains why that takes a deliberate flush). Consistency, here, asks what a reader sees right after a write, such as whether the sender sees their own message at once. Availability is the share of requests that get a proper answer, so 99.9% means one request in a thousand may fail.
Each of these becomes a question to ask before designing, and each answer changes something:
| Question | Example answer | What it changes |
|---|---|---|
| How many users, and how active? | 20M daily active users | Every other estimate |
| Read-to-write ratio? | 3 reads per write | Whether caching or write throughput dominates |
| How much data per item, and for how long? | 300 B per message, kept forever | Storage growth, whether one node can hold it |
| Latency target, and at which percentile? | p99 under 200 ms to load a conversation | What can sit on the read path |
| Durability? | An acknowledged message is never lost | Copy to a second machine and flush to disk before replying |
| Consistency? | A sender sees their own message at once | Where reads can be served from |
| Availability target? | 99.9% of requests succeed | Redundancy, and what degrades first |
| Peak shape? | Evening peak about 3× the daily average | Capacity you provision |
?Why ask about percentiles and not just "fast"?
Because the average hides the requests that make users leave. A p99 target says what the slowest 1% may take (the slowest requests are called the tail), and that decides whether a synchronous call to a slow dependency is allowed on the path. Chapter 16 shows how quickly the tail grows with load and fan-out.
Now we have the questions. Before we use them, one practical matter: the arithmetic in section 1 needed Python, and in a meeting or an interview nobody hands you a terminal.
03Doing the arithmetic by hand
3.1Conversions worth memorising
A few conversions carry most estimates, rounded so that the multiplication is easy. Request rates are often written QPS, queries per second, as in the script's variable names, so you'll see that name for what we've been calling requests per second.
| Quantity | Exact | Round to |
|---|---|---|
| Seconds in a day | 86,400 | 10⁵ |
| Seconds in a 30-day month | 2,592,000 | 2.5 × 10⁶ |
| Seconds in a year | 31,557,600 | 3 × 10⁷ |
| 1 million requests a day | 11.6 per second | 10/s |
| 1 billion requests a day | 11,574 per second | 10⁴/s |
| 2¹⁰, 2²⁰, 2³⁰, 2⁴⁰ bytes | KiB, MiB, GiB, TiB | 10³, 10⁶, 10⁹, 10¹² |
Rounding a day to 10⁵ seconds makes every estimate about 14% too low, which is smaller than the error in any of our guesses, so it doesn't matter. Two shortcuts follow from the table. A million requests a day is about ten a second, and a billion is about ten thousand.
A product has 10 million daily active users, each making 50 API requests a day. About how many requests per second does the API serve on average?
With the conversions and the questions in hand, we can redo section 1 properly, using a fuller set of requirements that also says how people read.
04A full estimate for the chat service
4.1The requirements
Here are the fuller requirements for the chat service. The numbers are assumptions, stated, so you can change any of them. A conversation is a thread of messages between a set of people, and opening one shows its most recent messages.
- 20 million daily active users.
- Each sends 40 messages a day, averaging 300 bytes with metadata.
- Each opens 100 conversations a day; opening one loads its last 50 messages.
- Messages are kept forever. An acknowledged message must never be lost.
- p99 under 200 ms to open a conversation. Evening peak is 3× the average.
4.2Requests, storage, bandwidth
Reads need one more decision, which is where the messages live. Let's assume the first design stores them in Postgres, a widely used relational database, with one row per message and an index next to the table: a sorted lookup structure that lets the database find the last 50 messages of one conversation without scanning the whole table. In a test table of one million rows, a message's row plus its index entry came to about 340 bytes, a little more than the 300 bytes of the message itself, and we'll use that for storage.
| Messages written per day | 20M × 40 | 800M |
| Average write rate | 800M / 86,400 | 9,300/s |
| Peak write rate | 9,300 × 3 | 28,000/s |
| Conversation loads per day | 20M × 100 | 2B |
| Peak read rate | 2B / 86,400 × 3 | 69,000/s |
| Stored per day, Postgres row + index | 800M × ≈340 B (1M-row test table) | 272 GB |
| Stored per year, one copy | 272 GB × 365 | 99 TB |
| Stored per year, three replicas | 99 TB × 3 | ≈300 TB |
| Peak read bandwidth | 69,000/s × 50 × 300 B | 1.0 GB/s |
| the number that decides the design | ≈100 TB a year | |
Read the table from the top. Twenty million users at forty messages is 800 million writes a day, which becomes 28,000 a second at peak. Each of the 20 million users opens a hundred conversations, so 2 billion loads a day, which becomes 69,000 a second at peak. The storage rows multiply the daily volume by the row size, then by 365, then by three for the replicas.
?Why did bandwidth show up as a big number?
Bandwidth is how many bytes per second move through a link. Every conversation load returns 50 messages, so 69,000 loads a second is 1 GB/s of payload leaving the storage tier, before protocol overhead. That's a reason to cache recent messages close to the app servers, return fewer messages per page, and compress responses. The request count alone would never have told you.
4.3Working set and connections
Two more estimates decide how big a cache is and how many connections the database needs.
The first is the working set, the part of the data that's touched often enough to be worth keeping in fast memory. Most conversation loads are of recent, active conversations. If a cache, a fast copy of recent data kept close to the app servers, holds the last 50 messages of every conversation active today, say 20 million conversations, that's 20M × 50 × 300 B = 300 GB. If a tenth of them take most of the reads, then 30 GB of cache memory covers the hot set. The cache could be Redis, a server that keeps values in memory and looks each one up by a name, its key, such as the conversation ID.
The second is how many database connections we need. Little's Law (Chapter 16) says that the number of requests in flight at once equals the arrival rate times the time each spends in the system. A database connection is busy for the whole query, so requests in flight is also the number of connections in use. At 69,000 reads a second and 0.2 ms per query, about 14 queries are in flight at once. If the database is struggling and each query takes 20 ms, it's 1,400, which is more connections than a Postgres node should hold. A connection pool, a fixed set of open connections that requests borrow in turn, has to be sized and given timeouts for the slow case and not the fast one.
We now have the numbers that matter: 28,000 writes a second, 69,000 reads a second, 1 GB/s of read traffic, 100 TB a year, and a 30 GB hot set. A number is only big or small compared with something, though, so we need to know what one machine can do.
05What one machine can do
5.1Latency: how long one operation takes
The classic reference is the table of approximate timings in Peter Norvig's Teach Yourself Programming in Ten Years, widely circulated as "latency numbers every programmer should know". Its figures are old, so the table below gives the same kinds of operations on current hardware, next to Norvig's.
A few rows need explaining. A loopback round trip sends one byte over TCP (the protocol under most network connections, HTTP included) to a program on the same machine and gets a reply, with no real network involved, so it shows the cost of the software alone. fsync and fdatasync ask the operating system to flush a file's data to the storage device, and F_FULLFSYNC is macOS's stronger request, which also makes the drive empty its own internal cache before answering. p50 is the median, the latency that half the requests beat.
| Operation | Result | Where | Norvig's table |
|---|---|---|---|
| Random memory access, 256 MB working set | 124 ns | Linux VM | 100 ns |
| Read 1 MB sequentially from RAM | 35 µs | Linux VM | 250 µs |
| Loopback TCP round trip, 1 byte | 31 µs (p99 47 µs) | Linux VM | — |
| 4 KB write + fdatasync | 65 µs (p99 up to 600 µs) | Linux VM | — |
| 4 KB write + F_FULLFSYNC (a real flush) | 3.9 ms | laptop SSD | — |
| Postgres: fetch one conversation's messages by index | 0.083 ms | Linux VM, 1 client | — |
| Redis GET, 50 clients | 0.11 ms p50 | Linux VM | — |
| Disk seek (spinning disk) | — | — | 8 ms |
A virtual machine (VM) is a computer simulated in software on a bigger one, and "Linux VM" here means a small one with four CPUs shared with other work. The "laptop SSD" row is a Mac's built-in drive. Each figure is the median of three runs, and all of them will differ on your hardware; what carries over is the size of the gaps between rows.

Look at the flush rows. Even the cheapest one is hundreds of times slower than a memory access, and the two flush rows differ from each other by a factor of 60.
?Why is the VM's flush 60 times faster than the laptop's?
Because they do different amounts of work. On macOS, plain fsync (about 20 µs in the same test) hands the data to the drive without making the drive empty its cache, while F_FULLFSYNC waits for that too, which is where the 3.9 ms goes. The VM's fdatasync lands in a virtual disk that is itself a file on the host machine, so the host may still be holding the data in memory when the VM is told it's done. Treat the 65 µs as a lower bound on what a flush can cost, with no promise that the data is safe. Real server SSDs with power-loss protection sit in between. Chapter 08 explains what fsync promises. For design, the lesson is to time the disk you'll run on before you size anything that commits on every request.
5.2Throughput: how much one node sustains
Latency tells us how long one operation takes. Throughput tells us how many a node completes per second, and that's what we compare with 28,000 writes and 69,000 reads. These are one small node's numbers, meant as orders of magnitude and not as benchmarks.
Here's how to read the rows below.
- A client is one connection that sends a request, waits for the answer, and sends the next.
- Rates are counted in operations per second (ops/s) or, for Postgres, transactions per second (tps).
- Pipelining means a Redis client sends sixteen commands before waiting for any answer, which saves most of the round trips.
pgbenchis Postgres's own benchmarking tool, and-Smakes it run only the simplest read, a lookup of one row by its unique ID, the primary key.- A normal commit makes Postgres record the change in its write-ahead log and flush that log to disk before confirming. Setting
synchronous_commit = offskips the wait for the flush and confirms early, which risks losing the last moments of commits in a crash. - Inserting in batches of 100 rows means one transaction and one commit for every hundred rows.
- "TPC-B-like" is a standard bank-account workload where every transaction updates one account, one teller and one branch row.
| Workload | Throughput | Setup |
|---|---|---|
| Redis GET, no pipelining | 238,000 ops/s | redis-benchmark, 50 clients |
| Redis GET, pipelined 16 deep | 2,380,000 ops/s | redis-benchmark, 50 clients |
| Postgres primary-key SELECT (pgbench -S) | 16,000 tps / 44,000 tps | 1 client / 8 clients |
| Postgres: last 50 messages of a conversation | 12,000 / 37,800 queries/s | 1 client / 8; 1M rows, in RAM |
| Postgres single-row INSERT, commit per row | 35,800 rows/s | 8 clients |
| Same, synchronous_commit = off | 69,900 rows/s | 8 clients |
| Postgres INSERT in batches of 100 rows | 278,800 rows/s | 8 clients |
| Postgres pgbench TPC-B-like (scale 10) | 4,600 tps | 8 clients; contended rows |
All of these come from the same Linux VM, with Postgres left at its default of 128 MB of memory for caching table pages (the shared_buffers setting), each figure the median of three runs of 10 to 15 seconds. "In RAM" means the whole table fits in memory, so no query waits for the disk.
The VM's cheap flush flatters every commit-per-row number. When several clients commit at the same moment, Postgres can cover all of their commits with one flush of the log, which is called group commit. On a disk where each flush takes milliseconds, group commit is what keeps commit-per-row throughput from collapsing, and the numbers above would be lower.
Now lay the chat service's peak writes against the node.
| Peak writes the chat service needs | from section 4 | 28,000/s |
| One node, single-row commits | 28,000 / 35,800 | 78% of the node |
| One node, batches of 100 rows | 28,000 / 278,800 | 10% of the node |
| what grouping writes buys | 78% → 10% | |
Single-row commits would leave one node nearly full at peak, and grouping the writes leaves it nearly empty.
The last row of the table looks wrong at first. A bank-account workload on the same node gets only 4,600 transactions a second, a tenth of what plain inserts reached.
?Why is the TPC-B-like number so much lower than plain inserts?
Because every TPC-B transaction updates one of only ten branches rows at scale 10, so the eight clients spend their time waiting on each other's row locks. That's the difference between a workload's raw cost and its contention, and designs fail on the second one. (The mechanics are in Chapter 13.)
5.3Distance
Some latency no engineering removes. Light in optical fibre travels at about two-thirds of its speed in a vacuum, roughly 200,000 km a second, which is about 1 ms of round trip per 100 km. A cloud provider splits its servers into regions, and each region into several availability zones, separate data centres that fail independently.

| Hop | Distance | Round trip, physics alone |
|---|---|---|
| Same machine, loopback | — | 31 µs from section 5.1 (software, not distance) |
| Between AWS Availability Zones | "within 100 km (60 miles) of each other" (AWS) | up to about 1 ms |
| New York to London | about 5,600 km | at least 56 ms |
| Los Angeles to Tokyo | about 8,800 km | at least 88 ms |
Real paths aren't great circles, so real round trips are longer, which makes these figures a floor. A design that makes three sequential Los Angeles to Tokyo calls on every request has spent at least 264 ms before any code runs, more than a 200 ms target allows.

5.4Which number breaks one machine first?
We can now compare each requirement from section 4 with one node.
Compare these estimates with the single-node numbers in section 5.2. Which requirement is the first to rule out a single Postgres instance?
So the numbers say something specific: writes need grouping, reads need help, and storage needs the data split up. Each of those is a different box, and the next question is which box answers which number.
06Which box to draw
6.1Each component answers a number
Start from the simplest design that could work: stateless app servers (servers that keep nothing between requests, so any of them can answer any request) in front of one database. Then add a component only when an estimate crosses a line.
The table pairs each common component with the situation that justifies it and the price you pay. Some of its terms are defined inside the cells.
| Component | Add it when | What it costs you |
|---|---|---|
| More app servers behind a load balancer, which spreads requests over identical copies | CPU or memory on one server is the limit, or you need to survive losing one | Almost nothing if the servers are stateless |
| Read replicas, extra copies of the database that serve reads | Reads exceed one node, and slightly stale reads are acceptable | Replication lag, since a replica is a little behind; read-your-writes (seeing your own message at once) needs routing |
| Cache (Redis, memcached) | A hot working set is read far more than written, or a read is expensive | Invalidation (knowing when a copy is out of date), stampedes (section 8), a second source of truth |
| Queue (Kafka, SQS), a list that producers add work to and consumers take from later | Work can happen later, bursts exceed steady capacity, or several consumers need the same events | Asynchrony: consumers must be idempotent (handling a message twice does no harm), and you need to watch lag, how far behind they are |
| Partitioning (also called sharding), splitting the data across several databases that each own a slice | Data or write rate exceeds one node | Cross-partition queries, rebalancing, hot keys (one key far busier than the rest) |
| Object storage (S3), which keeps files by name, cheaply | Large blobs, or cold data read rarely | Higher latency per request, no in-place updates |
| CDN, servers near users that keep copies of files | The same bytes are read by many users far from the servers | Cache invalidation at the edge |
| Search index, a structure built for text search | Queries need full text or many filter combinations | A second copy to keep in sync |
?Why not add them all up front, to be safe?
Because each one is a new way to fail. A cache can serve stale data or fall over and send its entire load to the database at once. A queue can back up for hours without anyone noticing. Partitioning makes every query that crosses partitions slower and every migration harder. Complexity is a cost you pay on every deploy and every incident, so you pay it only for a number.
6.2How far one box goes
Simple designs go further than you'd expect. Nick Craver's Stack Overflow: The Architecture – 2016 Edition reported, for one day in February 2016, 209,420,973 HTTP requests to the load balancers and 504,816,843 SQL queries, served by 11 web servers, 4 SQL Servers in two clusters and 2 Redis servers. That's about 2,400 requests and 5,800 queries a second on average. Craver adds that they're "down to needing only 1 web server", with the caveat "I'm not saying it's a good idea."
6.3How the chat design grows
Now we can build the chat design one forced component at a time. Two terms first. The primary is the copy of a database that accepts writes, and a synchronous standby is a second copy that receives every change before the commit is confirmed, so it can take over if the primary dies. In the picture, each card at the database is a load from the estimate, and its colour says whether one node can carry it: amber means close to the limit, red means over it, and green means handled.
Each step has a number next to it, and each number came from section 4. A design review can challenge any step by challenging its number.
Step five split the data by conversation ID, and that choice was more than a detail: a different key would have made the most common query much more expensive.
07Access patterns decide storage
7.1Choosing the partition key
Picking a key means asking which query runs most often, and we know it: opening a conversation, "the last 50 messages in this conversation", at 69,000 a second. Consider two ways to split conversation 7, which has three members, Ana, Ben and Cy, and four new messages. Option A picks the partition from the conversation ID. Option B picks it from the sender's user ID.
c7 · #101 means conversation 7, message 101. Each must be stored on one partition. Option A: choose the partition from the conversation ID.?Why partition by conversation and not by user?
Because the dominant query is "the last 50 messages in this conversation". Partitioning by conversation makes it a single-partition index lookup. By user, a group conversation's messages are scattered across the senders' partitions, and every load has to ask every partition. The partition key should be whatever the hottest query filters on. Chapter 29 covers how keys map to partitions and how data moves when you add one.
That settles the key for one query. The API has more than one, though, and the others don't filter on conversation.
7.2Write the queries next to the API
The API is the handful of calls clients can make. For each one, write down the query it runs and how often. For the chat system:
| API call | Query | Rate at peak | Needs |
|---|---|---|---|
POST /conversations/:id/messages | Insert one row | 28,000/s | Durable commit, IDs that only go up within a conversation |
GET /conversations/:id/messages?before=… | Last 50 by (conversation_id, id) | 69,000/s | An index sorted by conversation, then ID; a cache |
GET /users/:id/conversations | Conversations by last activity | Lower | A separate, per-user index |
GET /search?q=… | Full-text over a user's messages | Low | A search index, updated asynchronously |
?Why does the third row need a separate index?
Because it's keyed by user, and the messages are partitioned by conversation. Any query that doesn't filter on the partition key has to ask every partition, or be served from a second structure keyed the way it reads. That's the general pattern: one copy of the data per access pattern that matters, kept in sync by the write path or by a stream of changes.
The third and fourth rows are small. A feed shows the same problem on a much larger scale, because every reader's feed is a query over many authors' content.
7.3Fan-out on write, or on read
The biggest access-pattern decision in feed-like systems is when to do the work of assembling what a user sees. A timeline is the list of recent posts from the accounts a user follows, and fan-out is one post being delivered to many readers' lists. Twitter's home timeline is the well-documented example. From Raffi Krikorian's Timelines at Scale talk (2013), as summarised by High Scalability: "300K QPS are spent reading timelines and only 6000 requests per second are spent on writes." With reads 50 times more frequent than writes, Twitter did the work at write time:
The summary gives the details. Each home timeline is capped at 800 entries in Redis, and fan-out aims to finish "under 5 seconds, but it doesn't always work, especially when celebrities tweet". The same page says it could take up to 5 minutes for a tweet to reach Lady Gaga's 31 million followers. For accounts like Taylor Swift's, the advice was to skip fan-out and "merge in her timeline at read time".
| Strategy | Work happens | Read cost | Write cost | Breaks when |
|---|---|---|---|---|
| Fan-out on write | When content is created | One lookup | One insert per follower | An author has millions of followers |
| Fan-out on read | When a feed is loaded | One query per followed source, then merge | One insert | A reader follows thousands of sources, or reads dominate |
| Hybrid | Write for most authors, read for the largest | Lookup plus a small merge | Bounded | You need to maintain both paths |
The boxes are drawn and the key is chosen. Two checks remain before the design is finished: whether the latency target is reachable, and what happens when each box misbehaves.
08Latency budget and failures
8.1Spend the latency budget
A p99 target is a latency budget, and every hop on the request path spends some of it. We'll walk one conversation load from a phone to the database and back. Two terms appear on the way. TLS is the handshake that encrypts a new connection, and it costs extra round trips. The edge is the first server, near the user, that terminates that connection.
?Why count sequential calls and not total calls?
Because parallel calls cost you the slowest one, while sequential calls cost you the sum. If three sequential calls each have a p99 of 20 ms, budgeting 60 ms for them is the safe assumption, since all three can be slow together. Three parallel calls cost about the slowest of the three, which is still a tail, but only one of them. Chapter 16 covers what fan-out does to the tail.
The cache sits on that path and answers most loads, so we should also ask what happens when it stops answering.
8.2Walk each box's failure
For each component, ask three questions: what if it's slow, what if it's down, and what if it's overloaded? Start with the cache, because it's the box whose failure surprises people most. Take 100 conversation loads a second (a round number, so the percentages are easy to read) and a cache that holds what 95 of them ask for. Watch the database behind it.
What the picture shows has a name, a cache stampede. A cache that absorbs most reads means the database behind it was sized for the few that get through. When the cache empties, after a restart or a failover, the database suddenly gets the full read rate. The defences in the last frame have names too. Coalescing misses means letting one request fetch a key while the rest wait for its answer. Load shedding means refusing some requests quickly so that the rest still finish. The third defence, sizing the database for a cold cache, is capacity planning. Timeouts, retry budgets and load shedding for moments like this are covered in Chapter 40.
?Why does the cache row matter so much?
Because a cache is a load reducer until it fails, and then it's a load multiplier. With a 95% hit rate the database sees 5% of the loads, and with an empty cache it sees all of them, twenty times as many. The same three questions apply to every other box. For the chat design:
| Component | If it's down | The design's answer |
|---|---|---|
| App server | Its requests fail | Stateless; the load balancer routes around it |
| Cache | Every read goes to the database: a stampede | Database sized for a partial cold cache; request coalescing; cache replicas |
| Database primary for one partition | Writes to those conversations fail | Synchronous standby, automatic failover; other partitions unaffected |
| Fan-out or search indexer | New messages aren't searchable yet | Asynchronous by design; alert on lag, not on errors |
| Whole region | Everything fails | A decision: accept it within the availability target, or pay for a second region |
Most real designs are finished in this walk, and not in the earlier steps. What remains is to write the design down so that someone else can challenge it.
09Writing it down
9.1The method you just followed
A design that lives only in a whiteboard photo can't be reviewed. Before the document, here are the steps you took, in order, since each one answered the question the previous one left open.
Read top to bottom, the steps are a series of questions, each answered with a number, and each answer narrows what the next step can be. That's why every box in the chat design arrived with a number from section 4 attached to it.
9.2The design document
The document doesn't need to be long. It needs to show its reasoning:
- Requirements, with every number and every assumption.
- Estimates: peak QPS, storage per year, bandwidth, working set.
- API and data model, with the query each call runs.
- The design, with the "this is here because" sentence for each component.
- Alternatives considered, and the number that ruled each one out.
- Failure modes from section 8.2.
- What you'll measure to know it's working: the SLIs (service level indicators, measurements such as the share of conversation loads under 200 ms), and the capacity metrics that tell you when the next step in section 6.3 is due. Chapter 39 covers how to collect them.
9.3Rules that hold up
- Estimate before drawing, and write every assumption next to its number.
- Estimate more than requests per second: storage per year, bytes per second, working set and requests in flight.
- Multiply by a stated peak factor, and size pools and timeouts for the slow case.
- Start with one service and one database, and add a component only when you can write "this is here because number exceeds limit".
- Choose the partition key from the hottest query, and keep a second copy for each other important access pattern.
- Put expensive work on the rarer side of the ratio.
- Size the database for a cold cache, and coalesce misses.
- Run your own query shape on the hardware you'll deploy before the design depends on someone else's throughput numbers.
9.4What you trade for what
| You add | You get | You pay | When the bill arrives |
|---|---|---|---|
| Batching writes | Up to 7.8× more rows a second from one node | Each write waits for its batch to flush | As added latency on every write |
| A cache | Most reads never reach the database | Stale copies, invalidation, a second source of truth | When it empties, as a stampede |
| Partitioning | Storage and write rate beyond one node | Cross-partition queries, rebalancing, hot keys | On any query that doesn't filter on the key |
| Fan-out on write | One lookup per read | One insert per follower | When an author has millions of followers |
| Any extra component | An answer to one number | A new way to fail | On every deploy and every incident |
9.5Symptom, cause, fix in design reviews
| Symptom in the design | Likely cause | Fix |
|---|---|---|
| Five services and a queue for a system with 50 requests a second | Boxes drawn before numbers | Estimate; start from one service and one database |
| "We'll shard (partition) later" with 100 TB a year in the estimate | Storage never estimated | Choose the partition key now, from the hottest query |
| p99 target unreachable on paper | Sequential calls on the request path | Parallelise, precompute, or move the call off the path |
| Cache hit rate assumed at 99% with no working-set estimate | Working set never computed | Estimate the hot set; size the database for a cold cache |
| Queue added "for scale" | No burst or asynchrony requirement | Remove it, or state the burst it absorbs |
| Estimates quoted as averages | No peak factor | Multiply by a stated peak before comparing with limits |
| Benchmarks borrowed from a blog | Nobody ran the test on the target hardware | Run one node on the hottest query and the write path |
10Summary
- Numbers decide the boxes. Estimate before drawing, and give every component a reason of the form "this number exceeds that limit".
- Four guesses turn a sentence into load. Users, messages per user, bytes per message and a peak factor gave 13,889 writes a second and 29 TB a year for the first draft.
- Non-functional requirements shape the architecture. Users, ratio, size, latency percentile, durability and peak factor are the inputs.
- State every assumption. A written assumption can be corrected; a silent one is a bug.
- 86,400 seconds a day, about 10⁵. One million requests a day is about 12 a second; one billion is about 12,000.
- Know one node's limits. On one small VM: Redis 238k GETs a second, Postgres 36k single-row commits, 279k rows a second in batches.
- Batching is the biggest write lever. It was 7.8× on the same node.
- QPS rarely decides the design. In the worked example storage did, at about 100 TB a year; bandwidth and working set came next.
- Access patterns choose the partition key. Partition by what the hottest query filters on, and keep a second copy per other important pattern.
- Put expensive work on the rarer side of the ratio. Fan-out on write when reads dominate, on read for the accounts where writes explode.
- Finish with the latency budget and the failure walk. Sequential calls add their tails, and a cache that fails multiplies load.
11Build this
An estimate you can check.
- Pick a system you run. Write its requirements as in section 2.2 and estimate peak QPS, storage per year, bandwidth and working set.
- Test one node: run
pgbenchwith a custom script containing your hottest read and your write, andredis-benchmarkwith your value size. Compare with the estimate. - Find the first number that crosses a single-node limit. Is it the one your current architecture was built around?
- Write the one-page design document from section 9.2 for the next 10× of growth, with a number next to every component you'd add.
12Interview questions
beginnerHow do you turn daily active users into requests per second?›
Multiply users by requests per user per day, then divide by 86,400 (or 10⁵ for a quick estimate). 10 million users making 50 requests each is 500 million a day, about 5,800 a second on average. Then multiply by a stated peak factor, say 3×, before comparing with any capacity limit.
beginnerWhen would you add a cache?›
When a hot working set is read much more often than it's written, or when a read is expensive to compute, and the estimate shows the database can't serve the read rate or the latency target alone. Estimate the working set to size the cache, and size the database for the moment the cache is cold.
intermediateWhat usually forces a system off a single database?›
Often storage, not queries. Reads can be scaled with replicas and caches, and writes with batching, but data that grows by tens of terabytes a year doesn't fit on one node. In the chat example, 800 million messages a day at about 340 bytes each is roughly 100 TB a year, which forces partitioning, while the read and write rates alone didn't.
intermediateHow do you choose a partition key?›
From the hottest query. If the dominant read is "the last 50 messages of a conversation", partition by conversation ID so it's a single-partition index lookup. Then check the key spreads load evenly (no giant conversation owns a partition) and list the queries that don't filter on it, since each needs a second structure or a fan-out.
intermediateFan-out on write or fan-out on read for a news feed?›
Fan-out on write when reads greatly outnumber writes and most authors have modest audiences: precompute each reader's feed so a read is one lookup. Fan-out on read for authors whose audience is so large that writing to every follower is too slow. Twitter reported 300K QPS of timeline reads against 6,000 writes a second, fanned out on write, and merged the largest accounts in at read time.
deepYour design meets its average latency but not its p99. Where do you look?›
At sequential dependencies on the request path, because each adds its own tail to yours; at fan-out, where the slowest of N calls decides the response; and at shared resources near saturation, where queueing grows the tail first. Fixes are parallelising calls, precomputing, hedging slow calls (sending a duplicate request if the first is slow), moving work off the synchronous path, and adding headroom to the hottest resource.
deepWhy is a cache a risk as well as an optimisation?›
Because the database behind it is sized for the misses. If the cache serves 95% of reads, the database sees 5%. When the cache is emptied by a restart, failover or flood of evictions, the database suddenly sees all of it, twenty times its usual load, and can fall over, which keeps the cache from refilling. The design has to size for a partially cold cache, coalesce concurrent misses for the same key, and shed load while it warms.
13Go deeper
One billion requests a day: about how many per second on average?›
About 11,600. Divide by 86,400, or estimate with 10⁹ / 10⁵ = 10⁴.
800 million 340-byte records a day: how much storage per year for one copy?›
About 99 TB: 272 GB a day × 365.
Why did batching 100 rows per insert raise throughput 7.8×?›
The per-commit costs (round trip, transaction bookkeeping, write-ahead log flush) are paid once per batch instead of once per row.
What does 1 ms of round trip buy you in fibre?›
About 100 km of distance, since light in fibre covers roughly 200,000 km a second, there and back.
Real request and query counts, real server counts, and a long argument for running on a few big machines. nickcraver.com
The origin of the latency table everyone copies; compare it with your own timings. norvig.com
The book-length treatment of replication, partitioning, and derived data kept in sync: the components section 6 adds.
The two tools used for section 5.2. Custom pgbench scripts (-f) let you test your own query shape. pgbench docs, redis-benchmark
14Related chapters
The utilisation curve, Little's Law and fan-out: the arithmetic behind headroom and latency budgets. Chapter 16.
Service level objectives (SLOs), error budgets and what each box should do when a dependency fails. Chapter 40.
How to measure the SLIs and capacity signals a design document promises. Chapter 39.
What the cache box in section 6 is doing inside, and where it runs out. Chapter 22.
How keys map to partitions and how data moves when you add one. Chapter 29.
What fsync promises, and why a flush costs what it does. Chapter 08.