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.
For a social feed serving 200M monthly active users and 10k writes/sec, would you fan out a new post to followers' feeds on write, or compute the feed on read? What does each choice cost you, and how would a celebrity account with millions of followers change your answer?
Sample Answer
Direct answer
For 200 million monthly active users (MAU) and 10,000 writes/sec, a pure fan-out-on-write pushes every new post into every follower's feed at write time, buying very low read latency at the cost of massive write amplification and storage. Pure fan-out-on-read defers that work to feed-view time, keeping writes cheap but making every read do more work. A celebrity account with millions of followers breaks the pure push model outright, which is why the practical answer is a hybrid: push for ordinary accounts, pull (or a separate merge step) for very high-fan-out accounts.
Structured elaboration
| Dimension | Fan-out-on-write (push) | Fan-out-on-read (pull) |
|---|---|---|
| Storage | High: one copy of the post lands in every follower's inbox | Low: one canonical copy per post |
| Read latency | Very low: a feed read is a simple lookup | Higher and more variable: must merge recent posts from every followee at read time |
| Write amplification | O(followers) per post; scales with fan-out size | O(1) per post; writes stay cheap regardless of follower count |
| Rebuild after failure | Complex: losing the inbox store means replaying historical writes | Simple: the canonical post store is the source of truth, caches are just recomputed |
| Best fit | Accounts with small-to-medium follower counts | Accounts with very large follower counts (celebrities) |
The decision criterion is the read:write ratio implied by a given account's follower count, not a single global choice: an account followed by 200 people generates trivial fan-out and huge read-latency benefit from push; an account followed by millions generates enormous fan-out for a benefit (marginally faster reads for those followers) that pull-at-read can approximate at read time instead.
Worked example
At 10,000 writes/sec, assume (illustrative, pinned input) an average of 300 followers per post for non-celebrity accounts:
fan-out ops/s=10,000 writes/s×300 avg followers=3,000,000 inbox writes/s
That is the write-amplification cost a pure push model pays continuously just for ordinary accounts.
Now take one celebrity post going to 5 million followers, and assume (illustrative) a fan-out cluster capable of sustaining 500,000 inbox writes/s:
time to fan out one celebrity post=500,000 writes/s cluster capacity5,000,000 followers=10 s
A single celebrity post would take roughly 10 seconds to fully propagate through push fan-out, and that's before accounting for every other post competing for the same fan-out capacity at the same time. This is the concrete reason celebrity accounts change the answer: pushing their posts synchronously into millions of inboxes is not just expensive, it measurably delays delivery to everyone else sharing that fan-out capacity.
Trade-offs & pitfalls
- Treating fan-out-on-write and fan-out-on-read as a single global choice, rather than a per-account decision keyed on follower count, is the most common shallow answer.
- A hybrid design still needs a merge step at read time for celebrity posts, so pull-style merge logic doesn't disappear; it just gets scoped to a small fraction of accounts instead of all of them.
- Async, idempotent fan-out pipelines are required regardless of strategy, because retries and partial failures are certain at this scale; a synchronous fan-out-on-write implementation is a reliability risk independent of the storage trade-off.
- Caching the celebrity's own recent posts aggressively (rather than fanning them out) reduces the read-time merge cost without reintroducing full push fan-out.
You're designing a user profile service with global, low-latency reads. Fields like email, password, and account status need strong consistency. Fields like display name and profile picture can tolerate eventual consistency. How would you decide, field by field, which guarantee each needs, and how would you defend keeping the split instead of making everything strongly consistent?
Sample Answer
Direct answer
Decide per field with a simple test: what does a user or the business lose if this field is read stale for a few seconds, and does that loss involve authorization, money, or identity? Email, password, and account status gate who can act as whom, so they get a linearizable (single, globally agreed order) read/write path even at a latency cost. Display name and avatar are cosmetic: a stale value for a few seconds costs nothing but a visual blip, so they get eventual, region-local, low-latency writes and reads. Defending the split means showing what making everything strong actually costs on the read path, not just asserting that it is safer.
Structured elaboration
Per-field decision table
| Field | Guarantee | Why | Cost of getting it wrong |
|---|---|---|---|
| Password / auth credentials | Strong (linearizable) | A stale read could let an old, revoked credential keep working | Account takeover window |
| Account status (banned/suspended) | Strong | A stale read lets a banned account keep acting | Abuse, trust and safety failure |
| Email (used for login/recovery) | Strong | Same identity-resolution risk as password | Locked-out or hijacked account |
| Display name | Eventual | Cosmetic; a few seconds of staleness is invisible risk | Momentary visual mismatch only |
| Profile picture | Eventual | Same as display name; also a large binary, cheap to serve from cache or object storage | Momentary visual mismatch only |
| Billing / payment state (extension) | Correctness-critical but not necessarily linearizable | Money is at stake, but the fix is compensating transactions, not blocking global writes | Double charge or missed charge, needing a refund/reversal workflow |
Mechanism
This paragraph is implementation detail, useful to know by name but not required to follow the field-by-field argument made above it. Two logical stores per user: a small, strongly-consistent store (consensus-replicated, for example a Raft-based database, where Raft is an algorithm that gets a cluster of replicas to agree on the same order of writes, or a globally-consistent database) for the identity-critical fields, and a multi-region, eventually-consistent store (Dynamo-style or similar) for everything else. Reads compose a single user object from both stores, so only the strong-store portion pays the cross-region latency cost. Read-after-write for the strong fields comes from routing that specific read to the writer's region or the current leader; monotonic reads (once a client has seen a value, a later read never shows it an older one) for the weak fields come from a session token, not from the strong store.
Extending the framework: billing correctness without going fully strong
Billing state is the case that tempts people into "just make everything strong." Resist it: instead of a synchronous global commit for every billing event, use compensating transactions, an idempotent charge (safe to run the same charge request twice, say after a retry, without actually billing the customer twice) plus a defined reversal or refund path if a downstream step (fraud check, inventory hold) fails after the charge already happened. This gets you correctness (the ledger is right once reconciliation finishes) without paying the linearizable-everything latency tax on a field written far less often than it is read.
Defending the split with a number, not an opinion
The strongest defense against "why not just make it all strong" is quantifying what "all strong" costs on the read path, since profile reads vastly outnumber profile writes.
Worked example
Assume a region-local cache read costs 5 ms, and a linearizable read from the strong store (contacting a majority of replicas across 3 regions, with an illustrative one-way inter-region round-trip time (RTT) of 100 ms) costs roughly two one-way trips:
strong-store read latency≈2×100 ms=200 ms latency multiplier if every read used the strong path=5 ms200 ms=40×If, say, 95% of profile reads only ever touch display-name or avatar fields (illustrative traffic mix, would come from real access logs), forcing all of them through the strong store means 95% of read traffic pays a 40x latency tax for a guarantee only the remaining 5% of fields ever needed. That is the number to put in front of someone asking why you didn't make everything strongly consistent.
Now the revenue-risk quantification (the second absorbed angle): the case for still investing in correctness on the billing fields, even though they don't get the fully linearizable treatment either.
assumed error rate on a race-prone billing path=0.1%=0.001 assumed volume=200,000 billing transactions/day at average value $50 expected daily exposure=200,000×0.001×50=$10,000/dayTen thousand dollars a day of exposure (illustrative; in practice pulled from real incident and error-rate data) is what justifies spending engineering time on compensating transactions for billing.
Trade-offs & pitfalls
- The strong store becomes a small, high-value target: shard it narrowly (identity fields only) so its lower throughput ceiling never becomes the bottleneck.
- Session tokens that carry the last-seen strong-store commit are what give read-your-own-writes on the critical fields without every read hitting the leader; skipping this is a common miss that reintroduces stale-password bugs.
- Pitfall: treating "eventual consistency" as a synonym for "no correctness work needed." The weak store still needs a conflict-resolution rule (last-writer-wins or a merge function), or two concurrent display-name edits silently lose one.
- Pitfall: treating billing as either fully strong or fully eventual instead of reaching for the third option, compensating transactions, which is usually the right cost and correctness balance for money-adjacent but not identity-adjacent fields.
How do you decide the right granularity when splitting a system into services? Walk through how coupling versus cohesion, data ownership, and team boundaries change your answer.
Sample Answer
Direct answer
Split along business capability and data ownership, not by technical layer, and treat coupling and cohesion as the actual test: a service boundary is right when it groups things that change together and separates things that don't, and when one team can own its full lifecycle (build, deploy, operate) without waiting on another team to also deploy. Team size and deployment cadence usually decide the timing more than the theory does: a well-modularized monolith can run comfortably until the coordination cost of shared deploys and shared blast radius starts to exceed the operational cost of running the same code as separate services.
Structured elaboration
The criteria, applied together
- Bounded context or business capability: one service per coherent business concept (Orders, Inventory, Billing), not per database table.
- Data ownership: the service that owns a piece of data is its only writer; everyone else goes through its API or its events, never a shared schema.
- Deployment independence: if two "services" cannot be deployed on separate schedules without breaking each other, they are one service wearing two names, a distributed monolith.
- Team boundaries (Conway's Law: a system's structure tends to mirror the structure of the team that builds it): align a service to a team that can own it end to end, so ownership and org chart don't fight each other.
- Transaction boundary: keep operations that need a real ACID (atomicity, consistency, isolation, durability) transaction inside one service; cross-service consistency should default to eventual consistency plus an explicit compensating action, not a distributed transaction.
- Chattiness: if two components exchange many synchronous calls per user request, the network hop between them is pure overhead with no ownership benefit; merge them.
The team-size-driven worked example
Consider an org at 200 people, organized as roughly 20 teams, running a well-modularized monolith with clear internal module boundaries (a modular monolith). Model the shared deploy pipeline as a single server processing one deploy at a time, 30 minutes each, across a 16-hour working day (960 minutes):
deploy capacity/day=30960=32 deploys demand at 20 teams (1 deploy/day each)=20 deploys/day utilization=3220=62.5%At 62.5% utilization there is queueing delay, but the pipeline is stable. Now grow to 500 people, roughly 50 teams, same one-deploy-at-a-time pipeline:
demand at 50 teams=50 deploys/day>32 deploys/day capacityDemand exceeding capacity on a single-server queue means the queue is unstable: it does not just get slower, it grows without bound. That crossing point, not a stylistic preference for microservices, is the concrete signal to start extracting services along the module boundaries the modular monolith already has, so teams stop sharing one serialized deploy pipeline and one shared blast radius.
Anti-patterns that signal you split wrong (or didn't split at all)
- Shared database schema across "separate" services: the clearest sign of a distributed monolith with extra network hops.
- Splitting by technical layer (a UI service, an API service, a database-access service) instead of by capability: nothing can deploy alone, because every user-facing change touches all three.
- A "god" service or shared library that every team depends on for routine changes: it recreates the same coordination bottleneck a monolith had, with worse debugging.
- Over-splitting a capability that still needs real ACID guarantees just because a diagram looks tidier with more boxes.
Trade-offs & pitfalls
- Splitting too early, before the coordination cost above actually bites, buys distributed-systems complexity (network calls, partial failure, eventual consistency) for a coordination problem you didn't have yet.
- Splitting too late means the deploy-pipeline math above turns into a real, measured queue of waiting teams, not a hypothetical.
- The bounded-context choice is the expensive one to get wrong: correcting a wrong service boundary later means a data migration, not just a configuration change.
- Watch for teams treating microservices as a goal instead of a response to a specific coupling problem; the checklist above should produce the boundary, not the other way around.
A model needs to serve 10,000 queries per second at p95 latency under 50ms. Sketch the capacity plan: how many replicas would you provision, and what CPU and memory would you budget per replica?
Sample Answer
Direct answer
Size the fleet from a single relationship, Little's Law, rather than guessing a replica count directly: pick a per-request service time and a concurrency budget per replica, derive that replica's sustainable throughput, divide the target queries per second (QPS) by it, then add headroom so the fleet runs below saturation, since running near 100% utilization is exactly what blows up the 95th-percentile (P95) tail this plan is trying to protect.
Structured elaboration
The relationship. Little's Law states L=λW: the number of requests in flight (L) equals throughput (λ) times average time in the system (W). Inverting it per replica: if a replica holds C requests concurrently and each takes W seconds, its sustainable throughput is C/W.
Decision criteria for the assumptions:
- Service time budget: must leave room under the 50 ms P95 target for queueing and network overhead, not consume the whole budget on compute alone (recall from queueing math that latency blows up as utilization nears 100%, so some of the 50 ms has to be slack, not service time).
- Concurrency per replica should be tied to actual provisioned resources (for example, one in-flight request per vCPU core for a compute-bound serving path), not picked independently of the CPU budget, otherwise the CPU and throughput numbers won't reproduce each other.
- Headroom must cover both a safety margin against tail-latency blowup (target well under 100% utilization) and operational headroom (rolling deploys, node loss).
Worked example
Assumptions (illustrative, stated explicitly): average per-request service time W=20ms=0.02s, leaving roughly 30 ms of the 50 ms P95 budget as network and queueing slack; each replica is provisioned with C=10 concurrent in-flight requests, matched to 10 vCPU cores (one request per core).
Per-replica throughput via Little's Law: λreplica=C/W=10/0.02=500 QPS.
Replicas for raw throughput at the 10,000 QPS target: 10,000/500=20 replicas.
Add headroom for two separate reasons, both real costs, not one combined guess:
- Utilization headroom: target roughly 70% utilization to keep queueing delay small and protect the P95 tail: 20/0.7≈28.6→29 replicas.
- Rolling-update / failure headroom: reserve capacity equivalent to roughly 2 replicas being unavailable at any time: 29+2=31 replicas.
Resource budget per replica, tied directly to the concurrency assumption above rather than picked independently: 10 vCPU (matching C=10), plus memory for the served model (illustrative 4 GB) and runtime/framework overhead (illustrative 1 GB), rounded up with a small buffer to 6 GB RAM per replica.
Trade-offs & pitfalls
- The most common mistake in this kind of estimate is stating a per-replica throughput number that isn't derived from the stated service-time and resource assumptions, for example claiming 200 QPS per replica on 2 vCPU with 30 ms of CPU work per request implies a maximum of roughly 2/0.03≈67 QPS per replica, not 200; always check that the throughput, service time, and resource assumptions are mutually consistent before sizing the fleet on them.
- Sizing purely for average throughput and skipping the utilization headroom step will hit the QPS target on paper while missing the P95 latency target in practice, because queueing delay is non-linear near saturation.
- CPU and memory are independent constraints, a replica sized correctly for CPU-bound throughput can still be starved on memory if the model or working set doesn't fit, both must be checked, not just one.
- Validate every assumption with a real load test before committing capacity: ramp to the target QPS, hold it, and confirm the measured P95 matches the plan; if it doesn't, the service-time or concurrency assumption was wrong, not the arithmetic.
For a read-heavy workload with moderate writes, would you reach for a cache layer in front of the database or add read replicas? Walk through how you'd decide.
Sample Answer
Direct answer
For a read-heavy workload with moderate writes, for example a 90% read / 10% write split, the default lean should be read replicas, because they scale read capacity without adding a second consistency model to reason about. Add a cache on top only for a narrow, measured set of hot keys that replicas still cannot serve cheaply or quickly enough. The decision comes down to three questions: can the application tolerate replication lag or cache staleness, is the read traffic skewed enough that a small cache absorbs most of it, and is there enough engineering capacity to build correct cache invalidation.
Structured elaboration
| Dimension | Read replicas | Cache layer |
|---|---|---|
| Consistency | Eventual (replication lag); route read-after-write to the primary when needed | Explicit staleness via TTL (time-to-live) or invalidation logic |
| Operational complexity | Lower if using a managed database's built-in replicas (automated failover, monitoring included) | Higher, requires instrumenting invalidation, TTL tuning, and a new system to run |
| Cost model | Scales close to linearly with node count, often bundled into managed pricing tiers | Extra infrastructure, but can dramatically cut load on the underlying database for skewed traffic |
| Failure modes | Replication lag, split-brain on failover (two nodes each wrongly believe they are the primary, so both accept writes); mitigate with lag monitoring and routing critical reads to primary | Cache stampede on mass invalidation (many requests miss the cache at the same instant and all hit the database at once), stale reads if TTL is too generous; mitigate with request coalescing (merging those simultaneous identical requests into one database call instead of many) and short TTLs |
Decision rule: if managed read replicas are available with a replication lag the application tolerates, start there. Add a cache only where a specific, measured hot-key or hot-query pattern needs sub-database latency or needs to shed load the replicas can't absorb cheaply.
The absorbed framing of a 90% read / 10% write split is the same decision restated: the write share matters because every write still has to land on the primary and propagate down to every replica. At 10% writes this is a non-issue; if the write share climbed toward 40-50%, the datastore choice itself would need revisiting (see write-heavy architecture reasoning), not just the cache-versus-replica question.
Worked example
Assume a baseline read load of 10,000 requests per second (RPS) and, illustratively, that each database read replica sustainably serves 2,000 RPS at acceptable latency.
Without a cache: replicas needed =10,000/2,000=5 replica nodes (plus the primary handling writes).
With a cache in front, assume an 80/20 access skew (a common real-world pattern: 20% of keys account for 80% of reads) and a 90% cache hit rate on that hot 20%:
DB-bound reads=(0.8×10,000×(1−0.9))+(0.2×10,000)=(8,000×0.1)+2,000=800+2,000=2,800 RPS
Replicas needed with the cache in place: ⌈2,800/2,000⌉=2 replicas.
That's a drop from 5 replica nodes to 2 from caching just the hot 20% of keys, which is why a targeted cache is usually layered on top of replicas rather than chosen instead of them: it earns its operational cost only where the skew is large enough to matter.
Trade-offs & pitfalls
- Adding a cache first because it "feels faster," without first measuring read skew, risks solving an already-adequate problem while introducing invalidation bugs for no real gain.
- Replicas trade consistency for scale: if a user reads immediately after writing in the same session, that read must be routed to the primary or to a lag-aware router, or the user will see stale data from their own write.
- Cache stampede on a mass invalidation event can hit the primary at exactly the worst moment, right after the thing that made the cache go stale in the first place; request coalescing and staggered TTLs guard against this.
- A self-run cache cluster is a second system to operate, patch, and monitor; a managed database's built-in replicas usually cost less operational attention than they save, which is why replicas are the default and the cache is the exception.
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.