Numbers to know
Latency, throughput, and capacity numbers for sizing designs, not vibes.
Capacity estimates turn a whiteboard design from guesses into trade-offs. This page gives you the ranges, formulas and review triggers to size a system without pretending the numbers are laws.
Read this if your last attempt…
- You proposed sharding at 1M rows (a single Postgres table can often reach hundreds of millions or billions of rows with the right indexes and partitioning, depending on row size, write rate and maintenance)
- You can't estimate QPS from DAU in your head
- You don't know roughly how much a single Postgres or Redis node can handle, or what that depends on
- Your capacity math is "a lot" instead of a number
The concept
Modern hardware is more capable than many candidates assume. On a current cloud instance, one well-tuned Postgres primary can often serve tens of thousands of simple indexed reads per second (hot data in memory, pooled connections, cheap queries), and one Redis shard is commonly planned at ~100K simple operations per second (small values, no pipelining). If your system does 1K QPS and you're proposing sharding, you're probably adding complexity for no reason.
But single-node numbers move a lot with the workload: query shape, row and value size, index count, durability settings, instance class, network, and the p99 you're promising. So learn them as ranges with assumptions, say the assumption out loud, and treat thresholds as review triggers ("this is where I'd re-check the design"), not hard limits.
Six orders of magnitude separate a memory reference from a cross-continent round trip. Each rung is roughly an order (or more) slower than the one before — which is exactly why an in-memory cache hit beats a database call, and why keeping data in-region beats reaching across the planet.
Rough single-node envelopes on modern hardware — planning ranges with assumptions, not guarantees. Benchmark your own workload.
| System | Rough envelope (assumptions) | Review trigger |
|---|---|---|
| Redis | ~100K simple ops/s per shard (small values, no pipelining); higher with pipelining, lower with big values or Lua. No fixed memory cap: shard size is set by failover/resync time, persistence headroom and managed node sizes | Memory, p99, failover time or one hot key exceeds your budget |
| Postgres / MySQL | ~10K-50K simple indexed reads/s with hot data in memory; low thousands to ~10K+ durable writes/s (row size, indexes, fsync and sync replicas matter) | Sustained writes near the primary's headroom, or data size where restore, index builds and vacuum hurt (often several TB) |
| Cassandra | Roughly 10K+ small writes/s per node; historically ~1-3 TB of data per node for comfortable repair and streaming | Add nodes; watch compaction, repair time and per-node density |
| Kafka | Plan in MB/s: tens to a few hundred MB/s of producer traffic per broker with acks=all and RF=3 (disks, network, batching and compression matter) | Broker disk or network saturation, slow rebuild of a failed broker, or long retention (consider tiered storage) |
| Elasticsearch / OpenSearch | Hundreds to thousands of searches/s per node depending on query cost; Elastic suggests shards of roughly 10-50 GB | Search latency over SLO, heap pressure, or shards outside the size guidance |
| App server | Hundreds to ~10K req/s per instance, entirely dependent on the work per request | CPU above ~70% or latency over SLO |
| Nginx / Envoy | Tens of thousands of simple req/s per instance, scaling with cores | CPU or open connections near limits |
| S3 | At least 3,500 PUT/COPY/POST/DELETE and 5,500 GET/HEAD requests/s per partitioned prefix (AWS docs) | Spread keys across more prefixes; expect some 503 Slow Down responses while S3 scales up |
- These ranges are for back-of-envelope sizing. Real limits depend on schema, query mix, payload size, durability settings, instance class and the p99 you promise.
- Thresholds are review triggers: the point where you re-check the design, not the point where it breaks.
How interviewers grade this
- You estimate QPS from DAU fluently — a number, not "a lot".
- You know rough single-node envelopes (Redis ~100K simple ops/s per shard, Postgres tens of thousands of simple reads/s) and the assumptions behind them.
- You do capacity math when making scaling decisions, not by vibes.
- You know latency numbers and use them in your budget.
- You push back on premature sharding by naming a review trigger, not a magic threshold.
Variants
Latency numbers
The time each hop takes — from nanoseconds to hundreds of milliseconds.
These are the numbers that define your latency budget. Every serial hop in your design adds roughly this much (typical ranges, not guarantees):
- L1 cache reference: ~1 ns
- Main memory reference: ~100 ns
- SSD random read (4KB): ~100 microseconds (fast NVMe can be lower)
- Redis / Memcached GET (same AZ): ~0.5-1 ms
- Simple indexed DB query (same AZ): ~1-10 ms
- Cross-AZ round trip (same region): ~1-2 ms
- Cross-region (same continent, e.g. US-East to US-West): ~30-70 ms
- Cross-continent: ~70-90 ms US-East to Europe, ~200+ ms US-East to Singapore
- CDN edge hit from a nearby PoP: ~5-20 ms
- DNS cold lookup: tens to ~200 ms
- New TCP + TLS 1.3 connection: 2 round trips — a few ms in-region, roughly 2× the RTT across regions
The critical insight: the gap between in-memory (nanoseconds) and cross-continent (~100 ms) is about six orders of magnitude. That's why caching works so well — you're replacing a ~10 ms DB call with a ~0.5 ms Redis call, roughly a 20× improvement. And it's why global deployments with regional data matter for low-latency apps — cross-continent latency is physics, not engineering.
Choose this variant when
- Every design — these numbers should be muscle memory
Quick estimation formulas
The five formulas that size every system.
- 1QPS from DAU: DAU x actions_per_user / 86,400. Peak is roughly 2-3x average for steady business traffic, 3-5x for consumer apps with a daily cycle, and 10x or more for flash sales, scheduled events or viral spikes.
- 1Storage: rows_per_day x avg_row_size x 365 x retention_years. "1 GB/day" sounds small until you multiply by 5 years = 1.8 TB. Multiply by the replication factor for raw disk.
- 1Bandwidth: QPS x avg_response_size. This tells you your CDN/egress bill and whether you need compression.
- 1Cache size: working_set_fraction x total_items x avg_object_size gives the payload. A common starting assumption is a ~20% working set (the 80/20 rule); real access skew varies, so validate with a measured hit-rate curve. Provision roughly 2x the payload per copy for key overhead, fragmentation and headroom, plus replicas.
- 1Server count: peak_QPS / QPS_per_server + ~30% headroom. Size for peak, not average.
Compute these in the first few minutes of a design round, then use them to drive the infrastructure decisions that follow. "We need about 3 Redis shards because the working set is ~60 GB of payload, ~120 GB provisioned" beats "we'll use Redis".
Choose this variant when
- Capacity estimation section of every design round
- Before proposing any infrastructure component
Common over-engineering mistakes
Things candidates propose too early because they don't know the numbers.
Sharding at 1M rows: row count alone is rarely the problem. With the right indexes and time-based partitioning, a single Postgres table often reaches hundreds of millions or billions of rows; what limits it is row size, write rate, index maintenance, vacuum and restore time for your workload. Reach for sharding when a named trigger fires — sustained writes near what one primary absorbs with headroom, or data size where backup/restore, index builds and vacuum get painful (often several TB) — not because the row count sounds big.
Cache for 100 req/sec: a single Postgres instance handles this easily. You cache when read QPS approaches DB capacity or latency exceeds your SLO, not as a default.
Kafka for 100 events/sec: a database table with a polling consumer, or a managed queue, handles this. Kafka earns its operational cost when you need high sustained throughput, multiple independent consumer groups, or replay.
Microservices for a new product: a single well-structured service handles most early-stage products. Split when you have team-scale or isolation problems, not just load problems.
The rule: if you can't justify the complexity with a number (QPS, storage, latency), you're probably over-engineering. Knowing when NOT to add infrastructure is as important as knowing when to add it.
Choose this variant when
- When you hear yourself saying "let's add X just in case"
Worked example
Scenario: size the infrastructure for a URL shortener with 100M MAU. (Every number below is an assumption you say out loud.)
Step 1 — QPS:
- DAU = 100M / 3 = ~33M (consumer app DAU/MAU ratio ~1/3).
- Write: 1 new URL per user per day = 33M / 86,400 = ~380 writes/sec avg. Peak 3x = ~1,150.
- Read: 100:1 read:write ratio = ~38,000 reads/sec avg. Peak 3x = ~115,000.
Step 2 — Storage (5 years):
- 33M URLs/day x 365 x 5 = ~60B URLs. x 500 bytes = ~30 TB (before replication).
- ~1,150 writes/sec is fine for one primary, but ~30 TB is well past where most teams are comfortable with a single primary (restore time, index builds). Storage, not QPS, is the trigger: shard by short_code hash (say 4 shards of ~7.5 TB) or use a managed KV store.
Step 3 — Cache:
- Most redirects hit recently created or trending links. Assume links active in the last 30 days are hot: 33M x 30 = ~1B URLs x ~200 bytes (code + target) = ~200 GB of payload.
- Provisioned: ~2x for key overhead, fragmentation and headroom = ~400 GB per copy, e.g. 8 shards of ~50 GB, each with a replica.
- QPS per shard: 115K / 8 = ~15K/sec — comfortable for Redis.
- Target a 95% hit rate and treat it as an assumption to validate with real traffic. Don't derive a hit rate by stacking the 80/20 rule on itself.
Step 4 — Servers:
- 115K reads/sec peak. Assume a lightweight redirect handler does ~5K req/sec per instance (measure it).
- 115K / 5K = 23 servers. Add 30% = ~30 servers.
Step 5 — Origin check:
- At a 95% hit rate, origin sees ~5.8K reads/sec at peak, spread over 4 shards and their replicas. Easy. At 80% it's ~23K/sec — still workable, and that's the number to watch.
Result: ~30 app servers, ~8 Redis shards with replicas, Postgres in ~4 shards (or a managed KV store). Every number is backed by a stated assumption.
Good vs bad answer
Interviewer probe
“How many servers do you need?”
Weak answer
"A lot — we'll have millions of users. Let's use auto-scaling and it'll figure it out."
Strong answer
"33M DAU at 100:1 read/write gives ~115K reads/sec at peak. Assuming ~5K req/sec per app server, that's ~23 servers plus 30% headroom — call it 30. Cache: links active in the last 30 days are ~200 GB of payload, ~400 GB provisioned across ~8 Redis shards with replicas. Storage: ~30 TB over 5 years, which is past what I'd keep on one primary, so 4 Postgres shards by short_code hash. If the cache hits 95% — an assumption I'd validate — origin sees ~6K reads/sec at peak, which 4 shards handle easily."
Why it wins: Every infrastructure decision has a number and a stated assumption behind it. The interviewer can challenge any number and get a defended answer.
When it comes up
- During capacity estimation — the first 3–5 minutes of most design rounds
- When proposing infrastructure — every component needs a number
- When the interviewer challenges "do you actually need this?"
- When sizing cache, DB instances, or server fleet
- When pushing back on over-engineering (sharding, microservices, Kafka)
Order of reveal
- 1State the base unit. "1 day = 86,400 seconds. Base number for every QPS estimate."
- 2Compute QPS from DAU. "DAU × actions/user / 86,400 = avg QPS. Peak is roughly 2-3× for steady traffic, 3-5× for consumer apps, 10× or more for flash sales."
- 3State storage with retention. "Rows/day × size × retention_years. Without retention, the number is meaningless."
- 4Size each component from rough single-node envelopes. "Redis: ~100K simple ops/s per shard as a planning figure. Postgres: tens of thousands of simple reads/s, low thousands to ~10K durable writes/s. App servers: depends on the work per request, so I'll assume ~5K/s and say so. Divide peak by per-instance to size the fleet."
- 5Add 30% headroom. "Always size for peak, not average, and add headroom for growth and failure margin."
- 6Call out cache hit rate as the lever. "At X% hit rate, origin sees only (1-X)× the QPS. This is where the biggest scaling wins live, so I'll treat the hit rate as an assumption to measure."
- 7Push back on premature scaling. "At Y QPS / Z TB, we do NOT need sharding/Kafka/microservices yet. Here's the trigger at which we would."
Signature phrases
- “1 day = 86,400 seconds” — The single most-used number in capacity estimation.
- “Shard on a named trigger, not a feeling” — Pushes back on premature complexity without inventing a universal threshold.
- “Name the peak factor: 2-5× daily, 10×+ for events” — Prevents under-sizing for the worst hour of the day.
- “Assume ~20% hot, then measure” — A starting heuristic for cache sizing, honestly labelled.
- “Row count alone rarely forces sharding; write rate and data size do” — Counters the reflexive "Postgres doesn't scale" narrative without overclaiming.
- “Every number has an assumption” — Frames your sizing as defendable, not vibes.
Likely follow-ups
?“Your system is at 1k write QPS. Do you need to shard the database?”Reveal
Probably not. A single well-tuned Postgres primary can usually absorb low thousands to ~10K+ durable writes/sec, depending on row size, indexes and durability settings. At 1k QPS you likely have plenty of headroom — confirm it with a benchmark or your metrics.
Review triggers for sharding (re-check the design when any of these approach):
- 1Sustained writes near the primary's headroom on a hot table, after batching and index tuning
- 2Data size where operations hurt — backup/restore, index builds and vacuum take too long for your RTO (often several TB)
- 3Access-pattern divergence — different tables want different placement strategies
At 1k QPS, the right optimisations are:
- Indexes tailored to access patterns (the biggest lever)
- Read replicas for read scaling
- Connection pooling (PgBouncer)
- Partitioning by time (for large but low-QPS tables)
Sharding at 1k QPS adds real operational complexity (cross-shard queries, rebalancing, ops surface) for little benefit. Pushing back on it, and naming the trigger that would change your mind, shows judgment.
?“How do you estimate cache size for a system?”Reveal
Working set × object size, then overhead, then a hit-rate target you validate.
Step 1 — estimate the working set. A common starting assumption is the 80/20 rule: about 20% of items get about 80% of traffic. More skewed workloads (viral posts, celebrity content) concentrate traffic in an even smaller set; near-uniform access (random-key lookups) means the working set is close to all data. If you have access logs, a measured hit-rate curve beats any rule of thumb.
Step 2 — compute the cache payload.
working_set_fraction × total_items × avg_object_sizeExample: 1B URLs × 20% working set × 500 bytes = 100 GB of payload.
Step 3 — provision memory and shards. Add roughly 2× for per-key overhead, allocator fragmentation and headroom (~200 GB per copy), then add replicas. Redis has no fixed per-instance memory limit, but very large shards resync and fail over slowly and persistence forks need spare memory, so many teams keep shards in the tens of GB. For example, 4 shards of ~50 GB, each with a replica.
Step 4 — validate the hit-rate target. Typically design for a 90–95% hit rate. At 95%, origin sees 5% of traffic. If origin is comfortable at that load, the cache is sized right. Then measure.
Step 5 — plan for cold start. On a full cache flush, origin must handle 100% of traffic temporarily. Either (a) warm the cache with a replay job, or (b) size origin (or its rate limits) for the full load.
?“A cross-region call takes 60 ms. What does that mean for your design budget?”Reveal
Every serial cross-region hop costs you ~60 ms off the latency budget.
Implications:
- 1Parallelise what you can. If a request needs data from 3 regions, do all 3 calls in parallel — the cost is max(60, 60, 60) = 60 ms, not 180 ms.
- 2Cache aggressively in the caller region — eliminate the round trip entirely.
- 3Co-locate request path data. If users in EU always need EU data, serve them from EU; don't ping US-East on every request.
- 4Use async for non-critical cross-region work. Replication, backup, analytics — send cross-region through a queue, don't block the request.
Rough round trips (they vary by provider and route):
- US-East ↔ US-West: ~60-70 ms
- US-East ↔ EU-West: ~70-90 ms
- US-East ↔ Singapore: ~200+ ms
- Speed of light in fiber: ~200,000 km/s — the floor no optimisation breaks
At a 100 ms p95 budget, you have little or no room for a cross-region hop on the request path. Regional deployments with local data are the architecture that works.
Common mistakes
Modern instances are far more capable than many candidates assume. Don't propose sharding for a few thousand simple QPS. Do the math with current instance sizes, and say which assumptions (query shape, payload, durability) your numbers depend on.
"1 KB per row x 1M rows/day = 1 GB/day." For how long? 1 year = 365 GB. 5 years = 1.8 TB. Retention determines whether you need sharding. Always multiply by time horizon.
Peak is roughly 2-3x average for steady business traffic, 3-5x for consumer apps with a daily cycle, and 10x or more for event-driven spikes. Size infrastructure for peak. Your SLA doesn't say "available on average".
"Redis maxes out at X GB" or "Postgres can't do more than Y TPS" are not laws. Real limits depend on the workload and on what you're optimising for (p99, failover time, cost). Quote a range, state the assumption, and name the trigger that would make you re-check.
1 day = 86,400 seconds. This is the single most-used number in capacity estimation. Memorize it. 100M events/day / 86,400 = ~1,157 events/sec.
Practice drills
A social app has 300M MAU, 50M DAU, and users read about 20 posts per session. What's the read QPS?Reveal
50M DAU x 20 reads / 86,400 = ~11,574 QPS average. Peak (3x) = ~35K QPS. If each read were a simple lookup by primary key, one Postgres with read replicas and a cache could serve this. A real home feed is not that: it involves the follow graph, ranking and hot celebrity accounts, which is why feeds are usually precomputed per user or assembled from caches. So this is a feed-design and caching problem before it is a sharding problem.
Each post is 1 KB and the app writes 500M posts/day. How much storage per year?Reveal
500M x 1 KB = 500 GB/day. x 365 = ~180 TB/year before replication. That is far beyond what you would keep on one primary, so storage (not QPS) forces sharding — for example by author_id — plus tiering old posts to cheaper storage to keep the hot dataset small.
You need a Redis cache for 1B short URLs at 100K reads/sec. How many Redis shards?Reveal
QPS: 100K reads/sec could fit one shard at the usual planning figure. Memory: 1B x ~200 bytes = ~200 GB of payload, so provision ~400 GB per copy for overhead and headroom. There is no hard per-instance limit, but you would still shard: huge single instances resync and fail over slowly, persistence forks need spare memory, and managed node sizes cap you. For example, 8 shards of ~50 GB, each with a replica. Memory and failover time drive the shard count here, not throughput. Use Redis Cluster (or a Valkey equivalent) sharded by short URL key.
Cheat sheet
- •1 day = 86,400 seconds. Memorize it.
- •QPS from DAU: DAU x actions/user / 86,400. Peak ≈ 2-3x (steady), 3-5x (consumer daily cycle), 10x+ (events).
- •Redis: ~100K simple ops/s per shard (planning figure). Postgres: ~10K-50K simple reads/s; low thousands to ~10K+ durable writes/s.
- •Cross-AZ: ~1-2 ms. Cross-region: ~30-70 ms. Cross-continent: ~70-250 ms.
- •Shard on a named trigger: writes near one primary's headroom, or data size where restore and maintenance hurt.
- •Don't cache until read QPS approaches DB capacity or latency exceeds the SLO.
- •Working set for cache: start with ~20% of data, then measure. Provision ~2x payload, plus replicas.
- •Kafka: size in MB/s. Storage = write MB/s x retention x replication factor.
- •1 KB x 1M rows = 1 GB. The base unit for storage math.
- •Server count: peak_QPS / QPS_per_server + 30% headroom.
Practice this skill
No problem is tagged directly to Numbers to know yet. These published problems still exercise the same interview category.
Read this if