Why estimate at all
Numbers decide architecture. A service with 50 writes per second fits on one well-tuned relational database; one with 500,000 writes per second needs partitioning, batching and probably a log-structured store. Without a quick estimate you cannot tell which world you are in, so you either over-engineer a simple system or under-build one that will fall over. Interviewers use estimation to see whether you can reason quantitatively under uncertainty.
The output should be a handful of numbers that drive decisions: average and peak QPS, read/write ratio, bytes stored per day and over the retention period, network bandwidth, and how many servers or how much memory that implies. Every number should trace back to a stated assumption so it can be revised when the interviewer changes a parameter.
Powers of two and data sizes
Storage and memory are counted in powers of two, and the useful trick is that each step of 2^10 is roughly a factor of a thousand. That lets you move between bits, kilobytes and petabytes without a calculator. Remember also the sizes of common things: an ASCII character is one byte, a 64-bit integer or timestamp is 8 bytes, a UUID is 16 bytes, a short text post is a few hundred bytes, a compressed photo is a few hundred kilobytes, and a minute of HD video is tens of megabytes.
- 2^10 ≈ 1 thousand = 1 KB
- 2^20 ≈ 1 million = 1 MB
- 2^30 ≈ 1 billion = 1 GB
- 2^40 ≈ 1 trillion = 1 TB
- 2^50 ≈ 1 quadrillion = 1 PB
Latency numbers worth memorizing
A classic table popularized by Jeff Dean lists the cost of common operations. The exact values have drifted as hardware improved, but the ratios still teach the important lessons. Memory is fast and disk seeks are slow; avoid random disk access where possible. Simple compression is cheap relative to sending bytes across a wide-area network, so compress before shipping data between regions. Round trips between continents cost on the order of 150 ms, which is why multi-region systems avoid synchronous cross-region calls on the hot path.
The practical use of this table is to sanity-check designs. If a page requires ten sequential cross-region calls, you already know it cannot load in under a second. If a lookup needs a random disk read per request at 10,000 QPS, you know a cache or SSDs are needed.
- L1 cache reference: ~0.5 ns; main memory reference: ~100 ns
- Compress 1 KB with a fast codec: ~10 µs (microseconds)
- Random 4 KB read from SSD: ~150 µs; read 1 MB sequentially from memory: ~250 µs
- Round trip within a data center: ~0.5 ms
- HDD disk seek: ~10 ms; read 1 MB sequentially from HDD: ~30 ms
- Packet California to Netherlands and back: ~150 ms
Availability and SLAs
Availability is the fraction of time a service is operational, usually quoted in nines. Cloud providers publish SLAs (service level agreements) that commit to a level such as 99.9% or 99.99% and give credits when they miss. Each additional nine cuts allowed downtime by a factor of ten, and the engineering cost to get there rises steeply because you must eliminate ever rarer failure modes.
Two composition rules matter. Components in series multiply: if a request touches three services each at 99.9%, the whole path is about 0.999^3 ≈ 99.7%. Components in parallel (redundant replicas) combine as one minus the product of failure probabilities: two independent 99% replicas give 1 - 0.01 x 0.01 = 99.99%. This is the quantitative argument for redundancy.
- 99% → about 3.65 days of downtime per year
- 99.9% → about 8.77 hours per year (~1.44 minutes per day)
- 99.99% → about 52.6 minutes per year (~8.6 seconds per day)
- 99.999% → about 5.26 minutes per year
Worked example: a photo-sharing app
Assume 100 million DAU. Each user uploads 0.5 photos per day, so uploads are 50 million per day. Dividing by 86,400 gives about 580 write QPS; with a 3x peak factor, roughly 1,750. Each user views 50 photos per day, giving 5 billion views per day, or about 58,000 read QPS on average. The read:write ratio is 100:1, which immediately suggests heavy caching and a CDN for image bytes.
For storage, assume an average photo of 500 KB after compression. Daily new storage is 50 million x 0.5 MB = 25 TB per day. Over a year that is 25 x 365 ≈ 9,125 TB, about 9 PB, and with three-way replication roughly 27 PB per year. Egress bandwidth for views is 5 billion x 0.5 MB = 2.5 PB per day, about 29 GB/s on average, a number that only a CDN can absorb economically. Metadata, by contrast, is tiny: 50 million rows x 200 bytes is 10 GB per day.
Technique and presentation tips
Round aggressively. Treat 86,400 seconds as 100,000 and 99,987 / 9.1 as 100,000 / 10. The point is to keep arithmetic simple enough to do in your head while talking, and the rounding error is far smaller than the uncertainty in your assumptions anyway.
Write assumptions down where the interviewer can see them, and label every unit. Writing 5 next to storage is meaningless; 5 TB per day is not. Separate averages from peaks, and state the peak factor you chose. Common things to estimate are QPS, peak QPS, storage, cache size, bandwidth and number of servers; practice each until it takes under a minute.
- Convert per-day counts to per-second by dividing by ~10^5.
- 1 million requests per day ≈ 12 per second.
- Multiply storage by the replication factor and retention period.
- Sanity-check results against known systems you have seen.
From numbers to design decisions
Estimates are only valuable if they change the design. A write QPS under a few thousand usually means a single primary database is fine. Hundreds of terabytes per year means object storage rather than database blobs. A working set that fits in tens of gigabytes means one cache cluster can serve it from RAM. Bandwidth in gigabytes per second means a CDN is mandatory, not optional.
Close the loop out loud: state the number, then state the consequence. For example, 58,000 read QPS with a 100:1 ratio means we cache feeds and serve images from a CDN; 9 PB a year means blob storage with lifecycle policies to cold tiers. This is what interviewers mean when they say a candidate reasons with data.
Key numbers
Key terms
- QPS
- Queries (requests) per second handled by a service, quoted as average and peak.
- Peak factor
- The multiplier from average to peak load, commonly assumed between 2x and 10x.
- SLA
- A contractual commitment to an availability or performance level, with penalties for missing it.
- Nines
- Shorthand for availability percentages such as 99.9% (three nines).
- Read:write ratio
- How many reads occur per write, which drives caching and replication choices.
- Replication factor
- The number of copies of each piece of data, which multiplies raw storage.
- Working set
- The subset of data actively accessed in a time window, which sizes the cache.
Common mistakes
- Chasing precise arithmetic instead of rounding and moving on.
- Forgetting to state assumptions, so numbers cannot be checked or adjusted.
- Dropping units or mixing bits and bytes (network links are quoted in bits per second).
- Computing only averages and ignoring peak load.
- Forgetting replication and retention when sizing storage.
- Producing numbers and never connecting them to an architectural decision.
Further study
- Latency Numbers Every Programmer Should Know (Jeff Dean / Peter Norvig)
- Designs, Lessons and Advice from Building Large Distributed Systems (Jeff Dean, LADIS 2009)
- Google SRE book (chapter on embracing risk and error budgets)