System design math

A back-of-the-envelope calculator. Set daily active users and watch it cascade into requests/sec, servers, storage, and a monthly bill.

12 min read

Every system design conversation starts with the same move: estimate the load before drawing any boxes. Not to the dollar, and not to the server, but well enough to know whether the answer is three machines or three thousand, because those are completely different architectures. It is called napkin math, and it is mostly multiplication from a single number: daily active users.

Back-of-the-envelope calculations are estimates you create using a combination of thought experiments and common performance numbers to get a good feel for which designs will meet your requirements. (Jeff Dean)

One number, and everything below it is arithmetic

Start with daily active users and how many requests each makes. Multiply for requests a day, divide by 86,400 seconds for the average request rate. Then do the step most estimates forget: multiply again for the peak. Users are not spread evenly across the day, and real traffic peaks at two to five times the mean.

Try
50M users × 20 requests = 1B a day ÷ 86,400 s = 12K/s average × 3 = 35K/s at peak
50M
20
3x
Fifty million users making twenty requests a day is about 11,600 requests a second on average, and three times that at the peak.

A fleet sized for the average is a fleet that falls over every evening. Every figure below reads from this one, so change the users or the peak here and watch the servers, cache, database and storage follow.

How many boxes, and how much the cache saves

The application fleet is peak traffic divided by what one server handles. But not at 100%: running flat out means any deploy, garbage-collection pause or failed node tips you over, and latency has already climbed a cliff well before utilisation reaches 1. Size for about 70%, and keep a spare.

Peak traffic
35K/s
Servers needed
50
Run, with a spare
51
1K/s
Peak divided by capacity, at 70% so a bad node and a spike can happen at once, plus one spare.

Then split the traffic. Most products are read-heavy, and a cache in front of the database is the single biggest lever on everything downstream. Thanks to the 80/20 rule you only need to hold the hottest fifth of the data to absorb most of the reads.

  • Cache hits25K/s · 72%
  • Database reads6.2K/s · 18%
  • Database writes3.5K/s · 10%
  • Cache memory for the hottest 20% of objects: 343 GB
90%
80%
2 KB
At an 80% hit rate the cache absorbs most of the reads. Turn it off and every read falls through to the database.

Shards are for writes. Replicas are for reads

The database tier needs two numbers, computed from different terms, and confusing them is the most common mistake in a capacity estimate. A replica is a full copy: it serves reads, but it also has to apply every write. Adding replicas multiplies read capacity and leaves write capacity exactly where it was. Only sharding, splitting the keyspace so each node owns a slice of the writes, moves that number.

Shards (writes ÷ 2K)
2
Copies per shard (reads ÷ 5K)
1
Database nodes
2
Each column is a shard: the orange copy takes its writes, teal copies add read capacity.
5K/s
2K/s
Shards come from writes, replicas from reads. Raise the cache hit rate above and the replica count falls; it never touches the shards.

That is why “are we read-bound or write-bound?” decides the architecture. Replicas are cheap to add later. Re-sharding a live database is one of the most invasive changes you can make. Storage follows the same arithmetic: writes, times object size, times replicas, every day you keep the data.

New data a day
572 GB
After 3 years
612 TB
Peak egress
512 Mbps
204 TB
year 1
408 TB
year 2
612 TB
year 3
3x
3 years
Writes times object size times replicas, every day, for as long as you keep it.

The handful of figures that make this doable in your head

None of the arithmetic is hard. What makes it fast is knowing about a dozen constants well enough that you never reach for a calculator, and knowing which orders of magnitude are simply implausible.

A day86,400 s, call it 100,000: drop five zeros from requests a day to get requests a second
App serverabout 1,000 requests a second for plain CRUD
Postgres nodeabout 5,000 reads and 2,000 writes a second
Redisabout 100,000 operations a second
Cache20% of objects serve 80% of reads
Utilisationsize the fleet for 70% so spikes have room
Peak2 to 3 times the daily average, up to 10 times for spikes
Replication3 copies is the durable default
Data unitsKB, MB, GB, TB, PB: each about 1,000 (really 1,024) times the last

The other half is time: how long operations actually take. The most useful comparison on the list is this one. Reading a megabyte from an SSD takes about a millisecond. A round trip from California to the Netherlands takes about 150, set by the speed of light in fibre and not improvable by any engineering. That is why a chatty protocol across regions is unfixable while the same protocol inside a rack is merely untidy.

CPU and cacheMemoryStorageNetwork
  • L1 cache reference1 ns
  • Branch mispredict3 ns
  • L2 cache reference4 ns
  • Mutex lock and unlock17 ns
  • Main memory reference100 ns
  • Compress 1 KB (Snappy)2 µs
  • Send 1 KB over 1 Gbps10 µs
  • SSD random read16 µs
  • Read 1 MB from memory250 µs
  • Round trip in a datacenter500 µs
  • Read 1 MB from SSD1 ms
  • Disk seek10 ms
  • Read 1 MB from disk20 ms
  • Round trip California to Netherlands150 ms
Log scale: each step right is about ten times slower.

Finally, availability is quoted in nines, and each extra nine cuts the allowed downtime by ten. The catch is that availability multiplies along a request path: five services at 99.9% give you about 99.5% end to end. That is why redundancy has to exist at every tier, not just one.

Target availability
Down per day
1.4 min
Down per month
43.8 min
Down per year
8.76 hours
A common baseline SLA for cloud services. Chain 5 services at 99.9% and the request as a whole is up 99.501% of the time, 1.82 days down a year.
5
Each nine cuts downtime tenfold. Chaining services multiplies their availabilities, so the whole is always worse than its parts.

The short version

  • Daily users times requests, divided by 86,400, then multiplied by the peak factor.
  • Size fleets for about 70% utilisation; caches absorb most reads with 20% of the data.
  • Shards scale writes, replicas scale reads, and only one of them is easy to change later.
  • Memorise a dozen numbers, especially that a cross-ocean round trip is about 150 ms.