System Design Methodology and Trade-off Analysis Questions
The end-to-end approach to an open-ended design problem and the judgment that resolves it: clarifying scope and constraints, gathering functional and non-functional requirements, capacity and back-of-envelope estimation, and mapping requirements to a high-level architecture, then reasoning explicitly about competing options on cost, complexity, latency, and reliability to defend a choice. Covers driving a design interview from ambiguity to a proposal, trade-off frameworks, decision-making under uncertainty and incomplete information, reversible-versus-irreversible decisions, and defending choices under scrutiny. The process-and-judgment skill underneath every system-design case study.
Walk through the process you'd use to produce a quick capacity and cost estimate for a new system when you only have a handful of customer-provided numbers (like average request rate and daily data volume). What do you ask for, and how do you sanity-check the result?
Sample Answer
Direct answer
With only a couple of customer-provided numbers you cannot produce a precise estimate, but you can produce a defensible range: convert the given numbers into a small set of derived quantities using clearly labeled assumed multipliers, present the result as a low/likely/high band with every assumption visible, and immediately ask for the handful of additional numbers that would narrow the range the most, peak-to-average ratio, payload size, retention period, and read/write mix.
Structured elaboration
What to ask for beyond the customer's two numbers:
- Peak-to-average ratio (how bursty is the traffic relative to the average given).
- Typical request and response payload size.
- Retention period for any stored data (drives storage growth over time, not just a snapshot).
- Read/write ratio and replication or durability requirements.
- Regions served (affects network egress, data leaving the cloud provider's network to the internet or another region, which providers typically bill for separately from compute and storage, unlike incoming/ingress traffic, and multi-region cost multipliers).
Sanity-check method, more useful than checking the numbers in isolation: confirm the derived figures scale consistently with the two customer-given numbers, doubling the stated average request rate should roughly double the compute line and leave the storage line untouched, since storage tracks data volume, not request rate. If a change to one input moves every output line by the same factor, an assumption has been applied incorrectly.
Worked example
Customer gives two numbers: average request rate = 200 requests per second (RPS), daily data volume ingested = 50 GB/day.
Assumed, clearly labeled as illustrative since the customer didn't provide them: peak-to-average ratio = 3x (typical for diurnal web traffic), average response payload = 5 KB, retention = 90 days, replication factor = 2.
Compute: peak RPS =200×3=600 RPS. Assuming, illustratively, that one core sustainably handles 100 RPS at acceptable latency: cores needed at peak =600/100=6, provisioned with headroom to 8 cores.
Storage: 50 GB/day×90 days=4,500 GB=4.5 TB raw. With replication factor 2: 4.5×2=9 TB provisioned.
Network egress (using a decimal GB convention throughout, 1 GB = 1,000,000 KB, since egress is what vendors bill on and vendors bill decimal): 200 RPS×86,400 s/day×5 KB=86,400,000 KB/day=86.4 GB/day≈2.6 TB/month (86.4×30=2,592 GB=2.592 TB).
Cost banding, using illustrative unit prices purely to demonstrate the method, not tied to any specific vendor's current published rate: compute at $0.05/core-hour, storage at $30/TB-month, egress at $80/TB.
Compute=8×24×30×$0.05=$288/month
Storage=9×$30=$270/month
Egress=2.592×$80≈$207/month
Likely total=$288+$270+$207=$765/month
Applying a discovery-stage uncertainty band of ±40%: Low =$765×0.6≈$459, High =$765×1.4≈$1,071.
Trade-offs & pitfalls
- Presenting a single point number instead of a range reads as false precision when 2 of the 5 inputs used were assumed, not given, always show the band and label which numbers came from the customer versus which were assumed.
- Applying the same peak-to-average ratio to every workload type without asking is a common shortcut that silently mis-sizes bursty workloads (batch/ETL) versus steady ones (background jobs).
- Forgetting network egress is a frequent gap, and for read-heavy services it is often the largest line item, not a rounding error.
- Keeping the assumptions explicit and separate from the customer's real inputs means the estimate can be corrected later by swapping one assumption for a measured value, instead of redoing the whole model from scratch.
You need to cut the latency of a key product flow from 200ms to 50ms. How would you go about identifying the likely bottleneck, network, serialization, database, or algorithmic, before you start optimizing?
Sample Answer
Direct answer
Don't optimize the layer that looks slow, instrument the request path end to end first. Get a latency budget broken into per-hop numbers (network, serialization, database, business logic) that actually sum to the 200 ms observed, then attack the hop with the best ratio of milliseconds saved to effort required, re-measuring after every change rather than assuming which layer is guilty before the data says so.
Structured elaboration
Method, in order:
- Baseline with distributed tracing across the full request path, capturing per-hop timing, not just a total.
- Form one hypothesis per layer (network/TLS overhead, serialization cost, database query time, business logic compute) and check it against the trace data rather than intuition.
- Rank candidate fixes by (milliseconds likely saved) divided by (implementation effort and risk), not by which one is technically most interesting.
- Ship the highest-ranked fix, re-measure the full trace, and repeat, because fixing the biggest hop changes which hop is now biggest.
Isolation checks, when tracing alone doesn't localize it: compare with keep-alive/connection pooling on versus off to isolate network/TLS overhead, compare payload size before and after trimming to isolate serialization cost, and compare with and without a query cache or added index to isolate the database's contribution.
Worked example
Assume tracing on the current 200 ms path yields this breakdown (illustrative numbers, chosen to sum to the measured total):
| Hop | Current (ms) | Fix | Target (ms) | Savings (ms) |
|---|---|---|---|---|
| Network / TLS | 40 | keep-alive + connection pooling + regional colocation | 10 | 30 |
| Serialization | 15 | compact binary format, trim payload | 5 | 10 |
| Database query | 100 | targeted index + cache hot reads | 25 | 75 |
| Business logic | 45 | remove redundant recomputation | 10 | 35 |
| Total | 200 | 50 | 150 |
Reproducing the arithmetic: current total 40+15+100+45=200ms, matching the measured baseline. Target total 10+5+25+10=50ms, matching the 50 ms goal, and the sum of savings 30+10+75+35=150ms accounts for exactly the gap (200−50=150). The database hop is the largest single lever (75 ms, half the total savings) and gets prioritized first for that reason, not because it's assumed to be the culprit before measuring.
Trade-offs & pitfalls
- Jumping straight to rewriting business logic when tracing shows the database is half the budget is solving the wrong problem first, always rank by measured contribution, not by which layer is the most familiar to fix.
- Not re-measuring after each change stacks unverified assumptions, a fix that looked good in isolation can interact badly with the next one.
- Chasing 90% of the theoretical win on the hardest 10% of the effort (a protocol rewrite) before taking the cheap 30 ms keep-alive win first wastes the easiest gains.
- Caching for latency introduces a correctness trade-off (staleness) that needs an explicit owner and time-to-live (TTL), "just add a cache" without that ownership is a common wrong turn.
- Reserve architectural changes (removing a network hop entirely, changing the protocol) for after the low-risk, high-yield fixes are exhausted, they carry more deployment and compatibility risk and should be justified by the remaining gap, not reached for first.
You're building a stateful, write-heavy service that needs to sustain 10,000 writes per second with low latency. How does that write-heavy profile change your datastore and architecture choices compared to a read-heavy service?
Sample Answer
Direct answer
A sustained 10,000 writes-per-second, low-latency, stateful workload pushes you away from a design tuned for reads (a single write primary, heavy indexing, read replicas) and toward one built for write scaling: a storage engine optimized for sequential writes, a partitioning scheme that spreads writes across many nodes, and a replication model with an explicit, tunable durability-versus-latency trade-off rather than a single write bottleneck.
Structured elaboration
Why a read-optimized design breaks down here. Traditional B-tree storage engines perform random-access writes and update every index on every insert, each additional index roughly adds another write per record. A single-writer relational primary caps total write throughput at whatever one node's disk and CPU can sustain, and read replicas do nothing for write capacity, they only copy the primary's write stream.
What changes for write-heavy:
- Storage engine: log-structured merge (LSM) tree engines (used by databases like Cassandra, HBase, and the storage layer behind DynamoDB-style stores) append writes sequentially and merge them in the background, trading some read amplification (a single logical read may have to check several separate on-disk files before it can answer, since recent and older writes land in different segments) for much higher sustained write throughput than a B-tree.
- Partitioning: writes are sharded across many nodes by a partition key. The key must be chosen for even cardinality, a monotonically increasing key (like a timestamp or auto-increment ID) concentrates all new writes on one shard regardless of how many nodes exist.
- Replication and durability: instead of one primary with no built-in fan-out, use a quorum-based replication scheme, writes are acknowledged once a majority of replicas confirm, giving a tunable point between "acknowledge on one node" (fast, risks data loss) and "acknowledge on all nodes" (safest, slowest).
- Indexing discipline: keep secondary indexes to the minimum the write path can afford, every index is a write, this is the opposite instinct from a read-heavy design where more indexes are usually free wins.
Worked example
Assume, illustratively, that a single write-optimized node sustains 2,000 writes per second at the target latency.
Nodes needed for raw throughput: 10,000/2,000=5 shards.
For durability, replicate each shard three ways (tolerate one node failure without data loss): 5×3=15 total storage nodes.
A quorum write with N=3 replicas and a write quorum of W=2 means the client waits only for the second-fastest replica to acknowledge, not the slowest, bounding tail write latency while still guaranteeing the write survives a single node failure.
Cost contrast, provisioned versus per-operation pricing. At an illustrative $0.00001 per write operation under a consumption-priced managed service:
ops/day=10,000×86,400=864,000,000 writes/day
daily cost=864,000,000×$0.00001=$8,640/day≈$259,200/month
Against 15 provisioned nodes at an illustrative $400/node/month: 15×$400=$6,000/month. At this sustained write rate the per-operation model costs roughly 40 times more, which is why sustained high-volume writes usually favor provisioned or self-managed clusters, and why consumption pricing fits bursty, low-average workloads instead.
Trade-offs & pitfalls
- Carrying over every index from a read-heavy design roughly multiplies write cost by the number of indexes, audit which indexes the write path can actually afford.
- A low-cardinality or monotonically increasing partition key creates a hot shard that caps total throughput no matter how many nodes you add, this is the single most common write-scaling mistake.
- Waiting for all replicas (W=N) is the safest durability setting but the slowest; a majority quorum balances safety and latency, the exact quorum size is itself a trade-off decision, not a default.
- High write concurrency needs connection pooling and write batching, naive one-connection-per-request patterns hit connection limits long before they hit the storage engine's real capacity.
flowchart LR
Client --> Router[Write Router]
Router --> ShardA[Shard A Leader]
Router --> ShardB[Shard B Leader]
Router --> ShardC[Shard C Leader]
ShardA --> ShardARep[Shard A Replicas x2]
ShardB --> ShardBRep[Shard B Replicas x2]
ShardC --> ShardCRep[Shard C Replicas x2]
When designing a relational schema, how do you decide whether to normalize a table or denormalize it? Walk through the reasoning you would use, including what you gain and what you give up with each choice.
Sample Answer
Direct answer
Normalize when write correctness and storage efficiency matter most: each fact lives in exactly one place, so an update touches one row and there is no duplicate copy to drift out of sync. Denormalize when read speed matters most: copying a value into the table that needs it removes a join at read time, at the cost of extra storage and extra write work to keep every copy consistent. The decision is really about where you are willing to pay a cost: on the write path (normalized) or on the read path (denormalized).
Structured elaboration
What normalization buys you
- A single source of truth for each fact (a customer's name lives in one row in the customers table). Rename a customer once, and every order referencing that customer's ID sees the new name immediately, because nothing else stored a copy.
- No update anomalies: you cannot end up with two rows disagreeing about the same customer's e-mail address, because there is only one row.
- Smaller row sizes and less redundant storage, since each attribute is stored once.
What it costs
- Reads that need a full picture (an order plus the customer's name and the product's title) require joining across multiple tables. As the number of tables in the join grows, so does read latency and database load per request.
What denormalization buys you
- Fast reads: a single table scan or index lookup returns everything the page needs, no join required. This matters most for read-heavy, latency-sensitive paths (a product listing page, an order-history feed).
- Fewer round-trips and less join computation on the database, which matters at high read volume.
What it costs
- Duplicated data: the same fact (a product's name, a customer's e-mail) now lives in more than one row.
- Write amplification and staleness risk: change the source fact once, and every duplicate copy must also be updated, or the duplicates drift and become wrong. If you skip updating one copy, you now have silently inconsistent data.
- More total storage, since the same bytes are stored multiple times.
How to actually decide
- Estimate the read:write ratio on the specific table or field in question, not the system as a whole. A field read a thousand times for every write is a strong denormalization candidate; a field written as often as it is read is not.
- Ask how often the would-be-duplicated value actually changes. A product's category ID rarely changes; a live inventory count changes constantly. Denormalizing something that changes constantly multiplies your write cost and your staleness risk.
- Ask how expensive staleness is if a duplicate briefly lags. A denormalized display name that is a few seconds stale is usually fine; a denormalized account balance is usually not.
- Consider partial solutions before going fully one way: a materialized view or a cached read model gives you denormalized-shaped reads without hand-maintaining duplicate columns in the source tables, at the cost of a refresh lag you must define and tolerate.
Worked example
Take an orders schema. Normalized (third normal form): an orders table (order ID, customer ID, timestamp), an order_items table (order ID, product ID, quantity, unit price), a customers table, and a products table. Rendering an order-detail page means joining order_items to products (for the product name and image) and joining orders to customers (for the customer's name), a three- to four-way join.
Suppose the system processes 1,000,000 orders a month, averaging 3 line items per order, so 3,000,000 order_items rows are written per month. A normalized order_items row (order ID, product ID, quantity, unit price as fixed-width fields) is roughly 28 bytes. A denormalized version that also copies in the product name (about 24 bytes), product category (about 12 bytes), customer name (about 20 bytes), and customer e-mail (about 24 bytes) adds about 80 bytes per row:
At 3,000,000 rows a month, that is:
3,000,000×80 bytes=240,000,000 bytes≈240 MBof pure duplicate data added every month, before counting index overhead or replication. That is the storage side of the cost. The write side shows up when a product gets renamed: if that product already appears in 50,000 historical order_items rows, a normalized schema needs a single row updated in products; a denormalized schema that copied the product name into order_items needs all 50,000 rows updated (or accepts that historical order rows show the old name, which is a legitimate choice for orders specifically, since an order should arguably show the name as it was at purchase time, not the current name).
That last point is the real lesson: denormalizing an order line item's product name is often correct, not just a performance hack, because an order is a historical record and should not silently change when a product is renamed later. Denormalizing a customer's current e-mail address into the same row would be the wrong call, because you want that field to always reflect the customer's latest value, and a copy will drift.
Trade-offs & pitfalls
- Over-normalizing a read-heavy path (a product catalog page hit thousands of times a second) forces the database to redo the same multi-table join on every request, which is real, measurable load that a single denormalized read model would remove.
- Over-denormalizing a field that changes often multiplies write cost for a marginal read benefit, and creates a data-integrity bug class (stale duplicates) that is easy to miss in testing and expensive to debug in production.
- A common pitfall is denormalizing before measuring the actual read:write ratio, based on an assumption that reads are always dominant. Analytics and reporting schemas intentionally denormalize heavily (star-schema fact and dimension tables in an online analytical processing, OLAP, warehouse), because they are overwhelmingly read-heavy and batch-loaded; the live transactional path behind an online transaction processing (OLTP) system usually should not copy that pattern wholesale.
- The strongest senior answer treats this as a per-field decision, not a whole-schema philosophy: a single table can normalize some columns and denormalize others based on how each specific column is actually read and written.
Explain the difference between latency and throughput, and how the two relate to each other.
Sample Answer
Direct answer
Latency is how long a single request takes from request to response; throughput is how many requests the system completes per unit time. They are related through concurrency: at a fixed level of concurrency, throughput is roughly concurrency divided by latency, so throughput can be raised either by lowering per-request latency or by running more requests concurrently, at least until the system runs out of capacity, at which point requests start queueing and both latency and its variance rise sharply.
Structured elaboration
Little's Law is the bridge between the two
L=λW
where L is the average number of requests in the system (concurrency), lambda is the arrival rate (throughput), and W is the average time each request spends in the system (latency). This one relationship connects the two metrics completely.
Two regimes
- Below capacity: adding concurrency raises throughput roughly linearly without raising latency much, since the extra work overlaps with idle capacity.
- Near or above capacity: requests start queueing behind each other, and latency rises non-linearly. A small increase in load causes a disproportionate jump in the tail, a pattern basic queueing models (for example M/M/1: the standard textbook queueing model for one server with random arrivals and random service times, the model that produces the classic curve where wait time explodes as utilization approaches 100%, named here, not derived) predict and production systems reliably show.
Why percentiles, not the average, matter once load is near capacity
The mean can look fine while a growing minority of requests wait behind a queue; the 95th and 99th percentile (P95/P99) surface exactly what the average hides.
The two are not always aligned
Batching, processing many items in one call to raise throughput, typically raises the latency of any individual item in the batch. A system tuned to maximize throughput at all costs (large batches, very high concurrency) can make its own P99 latency worse, which is why the metric worth optimizing depends on the workload: a public, user-facing API should optimize for tail latency at a given throughput target; an overnight batch job should optimize for total throughput and mostly ignore any single item's latency.
Worked example
Suppose a service's actual processing time per request (its service time) is 10 ms, and the target is 500 requests/sec sustained. Little's Law says the average number of requests being served concurrently at that point is:
L=λ×Wservice=500 req/s×0.01 s=5 concurrent requests
If the worker pool has exactly 5 workers, the system is running at 100% utilization, and queueing theory's core warning applies: at or near full utilization, queue length and wait time become highly unstable, since there is no slack to absorb any variance in arrival timing or request duration. Sizing to a target utilization of about 75% instead:
capacity=ρtargetL=0.755≈6.7⇒7 workers
gives the system headroom to absorb bursts without its tail latency exploding, at the cost of running roughly 40% more capacity than the bare-minimum number, capacity that sits partly idle most of the time. That headroom is not waste; it is the price of a stable P99.
Beyond the mechanics, defending a capacity decision to non-engineering stakeholders usually means presenting this same relationship visually: a P50/P95/P99 latency trend next to a throughput trend over the same time window, so a viewer can see the point where rising throughput starts dragging tail latency up, rather than being told about it in the abstract.
Trade-offs & pitfalls
- Quoting only an average latency number, which hides that the system may already be close to its queueing knee for a meaningful fraction of requests.
- Treating "increase throughput" and "decrease latency" as the same goal; batching and running near full utilization both raise throughput while making individual-request latency worse.
- Sizing capacity to exactly the average expected load instead of leaving headroom, which looks efficient on a spreadsheet and causes a tail-latency incident on the first genuinely busy day.
- What a senior answer adds: naming which of the two metrics the workload actually cares about, rather than reciting the definitions of latency and throughput and stopping there.
Unlock Full Question Bank
Get access to all System Design Methodology and Trade-off Analysis interview questions and detailed answers.
Sign in to ContinueJoin thousands of developers preparing for their dream job.