Scalability Patterns and Techniques Questions
Scaling a system to handle growth in traffic and data: horizontal versus vertical scaling, statelessness, sharding and partitioning strategies, read replicas, and connection pooling. Covers capacity estimation, identifying bottlenecks, and the tradeoffs each scaling axis introduces. The general toolkit for taking a design from thousands to millions of users.
During a large traffic spike, your cloud autoscaler hit a quota limit and the service breached its SLOs. As the incident commander, what immediate mitigations would you take (manual scaling, throttling), how would you communicate with your cloud provider and stakeholders, and what medium-term fixes (quota monitoring, predictive scaling) would you put in place? How would you update runbooks and alerts to prevent a repeat?
Sample Answer
Direct answer
As incident commander, the first move is to stabilize without waiting on the cloud provider: manually add capacity wherever quota headroom still exists, and shed or queue the traffic you can't serve so the SLO (service level objective, the target you've committed to for availability or latency) breach doesn't get worse. In parallel, escalate to the provider and give stakeholders a clear, honest status. The incident is not resolved by root-causing; it's resolved by capping the damage now and then building the quota monitoring, predictive scaling, and updated runbooks that make sure a known quota ceiling never again gets discovered mid-spike.
Structured elaboration
Immediate mitigations (first 0-30 minutes)
- Manually provision capacity in a region, account, or instance family that still has quota headroom, preferring larger instance types over more instances if the limiting quota is instance count rather than vCPU.
- Throttle non-critical traffic at the edge: return 429 (rate limited) or 503 (unavailable) with a
Retry-Afterheader for low-priority requests (bulk exports, background jobs), preserving capacity for the traffic that actually matters. - Apply admission control instead of dropping traffic uniformly: if only a fraction of demand can be served, let requests in up to capacity (first-come or a fair queue with an ETA) rather than randomly failing a percentage of everyone's requests. This is what separates "everyone gets a slow, unfair experience" from "most users are unaffected and the rest see a clear queue."
- If a queue sits in front of a write-heavy path, let it absorb the burst rather than pushing writes straight through to the database at spike rate.
Communication
- Cloud provider: open a severity-1 support case immediately with the quota metric, current usage, and requested limit; this is not a channel to wait on for immediate relief, so treat it as a parallel track, not the mitigation itself.
- Internal: post a structured status (what's affected, current mitigation, next update time) to a single incident channel on a fixed cadence, not ad hoc.
- Stakeholders and, if customer-facing, a status page: state the actual impact honestly and give a concrete next-update time rather than a resolution promise you can't back.
Medium-term fixes
The underlying gap this incident exposes is that scaling was purely reactive with no advance knowledge of demand. The fix is to design ahead for the traffic patterns you can actually anticipate, which fall into a few recognizable shapes:
- A short, extreme, scheduled burst (for example, a flash sale expected to run around 100,000 requests/second for roughly 10 minutes): because the timing is known in advance, pre-scale capacity ahead of the event rather than relying on the autoscaler to react to it live, put a queue in front of the write path to smooth the burst instead of hitting the database at peak rate, and stage the pre-scale as a monitored, reversible step (canary the added capacity, keep an automatic rollback if error rates rise) rather than a one-way commit.
- A sharp, less-precisely-timed spike (for example, a marketing email that drives a 10x jump within about 30 minutes): the exact start time is fuzzier than a scheduled sale, so the plan needs a short-term component (aggressive reactive scaling plus throttling as a backstop for the first few minutes) and a separate medium-term component (pre-notifying the team before large sends go out, and pre-warming capacity ahead of known send windows so the reactive layer isn't starting from a cold baseline).
- An extreme, unscheduled-feeling spike far above baseline (for example, roughly 100x normal traffic during a flash sale): at that magnitude, no autoscaler reacts fast enough on its own, so the design has to include pre-warmed capacity sized to the expected peak and admission control that treats all waiting users fairly (a first-come or randomized queue with a visible position or ETA) instead of an uncontrolled scramble where whoever's request happens to land first wins and everyone else gets errors.
Across all three shapes, the common fix is the same: stop treating "quota is sufficient" as an assumption and start treating it as a monitored, tested constraint, with permanent quota increases requested ahead of realistic peak-plus-buffer, not discovered during an incident.
Runbook and alert updates
- Add a dedicated quota-exhaustion playbook: exact commands for manual provisioning in each region/account, the throttling levers available, the provider escalation contact path, and a decision matrix for when to degrade vs. scale vs. fail over.
- Add quota-utilization alerts at conservative thresholds (for example, flagged at 60/75/90% of the current limit as an illustrative staging, not a universal standard) so the team requests an increase before hitting the wall, not after.
- Add an alert specifically for "autoscaler issued a scale-out that did not result in additional serving capacity within an expected window," which is a different failure than "no scale-out was attempted" and needs its own signal.
- Schedule runbook drills (tabletop or live) that specifically simulate a quota ceiling being hit, since a runbook that has never been rehearsed against this exact failure mode is unlikely to be followed correctly under real pressure.
Worked example
Assume, as stated planning inputs rather than measured facts: steady-state traffic of 5,000 requests/second (RPS) served by 50 instances, giving a baseline capacity ratio of
50 instances5,000 RPS=100 instanceRPSA flash-sale spike hits 100,000 RPS for about 10 minutes, which is a total request volume of
100,000 sreq×600s=60,000,000 requestsAt 100 RPS/instance, serving the full spike needs 100,000 / 100 = 1,000 instances. If the account's instance quota is capped at 200, the achievable capacity at that ceiling is
200 instances×100 instanceRPS=20,000 RPSwhich is only a fifth of demand, so the fraction of traffic that has to be shed or queued once the quota ceiling is hit is
100,000100,000−20,000=0.80an 80% shortfall. That number is the case for admission control over random shedding: dropping 80% of requests indiscriminately produces a bad experience for everyone, while admitting exactly the 20,000 RPS the fleet can serve and fairly queuing the rest (with a visible wait, not a silent failure) turns the same shortfall into a bounded, predictable degradation instead of a chaotic one. It's also the case for the medium-term fix: a permanent quota request sized to at least the 1,000-instance peak, not the 200-instance historical average, is what prevents this specific ceiling from being hit again.
Trade-offs & pitfalls
- An emergency quota increase request is not instant relief; the mitigation plan cannot depend on the provider responding within the incident window, only the medium-term fix (a pre-approved higher baseline quota) removes that dependency.
- Manually bringing up capacity in a different region can violate data-residency or added-latency assumptions the service normally relies on; that trade-off needs to be made consciously during the incident, not discovered afterward.
- Throttling everyone equally punishes both low- and high-value traffic the same way; admission control that's blind to request importance is only marginally better than dropping randomly.
- Alerting only on "the autoscaler failed to scale" misses the earlier, more useful signal: quota utilization climbing toward its ceiling before a scale-out attempt ever fails.
Why does connection pooling matter for a service running at scale? Describe best practices for managing both database and HTTP connection pools: pool size, max open connections, idle timeouts, connection lifetime, and behavior under a spike in load. How would you test and tune these settings before production?
Sample Answer
Direct answer
Connection pooling matters at scale because opening a new database or HTTP connection is expensive relative to a request (TCP handshake, and for a database, authentication and session setup), so reusing a small set of warm connections instead of creating one per request lowers latency and prevents the backend from being overwhelmed by connection churn. The core sizing problem is that a pool is a per-instance setting but the backend has a fleet-wide connection ceiling, so pool size has to be planned across the whole fleet, not tuned in isolation on one instance.
Structured elaboration
Why pooling matters at scale, mechanically
- Connection setup cost: a TCP handshake, TLS negotiation (for HTTP), and for a database, authentication plus session/state initialization, all add latency if paid on every request.
- Backend resource limits: every open connection holds memory and, for a database, often a whole backend process or thread; a backend with a hard maximum connection count can be pushed into refusing connections or degrading badly under connection churn even if query volume itself is modest.
- Reuse turns a per-request cost into a one-time cost amortized across many requests on the same warm connection.
Sizing pools: the fleet-wide constraint
The number one mistake is sizing a pool as if the instance owns the whole backend. It does not; every other instance is drawing from the same ceiling:
pool_size_per_instance≤⌊number_of_app_instancesdb_max_connections⌋If the database allows 500 total connections and the service runs behind 20 instances, each instance's pool must stay at or below ⌊500/20⌋=25 connections, or a fleet at full pool utilization exceeds the database's ceiling and starts getting connection refusals, exactly when load is highest and refusals hurt the most. This constraint must be revisited every time the fleet is resized by autoscaling, which is the part teams most often forget: a pool size tuned for 20 instances silently becomes unsafe the moment autoscaling adds a 21st.
Core pool parameters
| Parameter | What it controls | Tuning guidance |
|---|---|---|
| Pool size (min/max) | How many connections are kept open per instance | Bounded above by the fleet-wide formula above; bounded below by enough to avoid queuing under normal load |
| Max open/concurrent connections | Hard ceiling the pool will not exceed even under burst demand | Set to protect the backend, not just to satisfy the busiest moment; excess demand should queue or fail fast, not force more connections open |
| Idle timeout | How long an unused connection stays open before being closed | Long enough to avoid re-opening connections for normal traffic gaps; short enough to release resources during genuine lulls |
| Max connection lifetime | Forces a connection to be recycled after a set duration regardless of use | Keeps the pool from silently holding stale or half-broken connections open indefinitely; also spreads out reconnections instead of all connections expiring together |
| Acquisition timeout | How long a request will wait for a pooled connection before failing | Should fail fast rather than block indefinitely, so an overload turns into fast, visible errors instead of a pile of hung requests |
Behavior under a load spike: the connection-storm problem
The specific failure mode worth naming: a deploy, a failover, or a sudden traffic spike can cause many instances to simultaneously reconnect or spin up new pooled connections at once, a connection storm, which can itself exceed the database's connection ceiling even though steady-state pool sizing was correct. This has a process/thread-model dimension too: a backend that spawns one OS process or thread per connection (a common relational-database architecture) pays a much higher per-connection memory and context-switch cost under a storm than one built around lightweight connection handling, which changes how conservatively you should size db_max_connections in the first place. Mitigations: stagger reconnects with jitter (small random delays) instead of reconnecting all instances at once, keep pool warm-up gradual rather than instantaneous on instance startup, and prefer acquisition timeouts with backoff over unbounded retry storms.
Testing and tuning before production
- Load-test at realistic peak concurrency and burst shape, not just average throughput, since spikes and connection storms are what actually break pool sizing.
- Vary the number of app instances in the test to confirm the fleet-wide formula holds at the target autoscaling range, not just at today's instance count.
- Watch active/idle/wait-count and wait-time metrics from the pool itself, plus backend-side connection and CPU/IO metrics, and tune size, idle timeout, and lifetime to minimize wait time while keeping the backend under its ceiling.
- Explicitly test the failure path: kill connections mid-flight, simulate a slow backend, and confirm acquisition timeouts and backpressure behave as designed rather than hanging.
Worked example
A service runs 20 instances against a database capped at 500 total connections. Using the formula above, each instance is capped at 25 pooled connections. During a load test that simulates a rolling deploy (all 20 instances restarting within a short window), every instance attempts to rebuild its pool of 25 connections at once: 20×25=500 simultaneous reconnect attempts against a ceiling of exactly 500, with zero margin for any connection still draining from the old instances. Adding jittered reconnect delays and reducing per-instance pool size to 20 (giving 20×20=400, leaving 100 connections of headroom during a rollover) eliminates the connection-storm failures observed in the unthrottled test.
Trade-offs & pitfalls
- Sizing a pool against a single instance's peak load, without dividing by the fleet size, is the most common and most damaging mistake; it works until autoscaling adds instances, then fails exactly under peak traffic.
- A pool with no acquisition timeout turns backend overload into cascading request pile-ups instead of fast, visible failures; for services making many short-lived connections, pairing the pool with a circuit breaker (a resilience pattern that stops sending requests to a struggling dependency, covered under high-availability patterns rather than here) prevents that pile-up from spreading further upstream.
- Idle timeouts set too aggressively cause needless reconnection churn during normal traffic dips; set too loosely, they let leaked or stale connections accumulate unnoticed.
- Connection leaks (code paths that acquire a connection and never release it, often on an error path) are the quiet failure mode: the pool looks correctly sized until leaked connections slowly starve it, and only a saturation metric with alerting catches this before an outage.
Explain the difference between horizontal partitioning (sharding) and vertical partitioning for scaling a dataset. What criteria would you use to choose a shard key, how would you plan and execute a resharding or rehashing operation in a cloud environment, and what operational challenges (rebalancing, hotspots, migration windows) should you anticipate?
Sample Answer
Direct answer
Horizontal partitioning (sharding) splits rows across nodes to scale storage and throughput; vertical partitioning splits columns or tables by function to isolate workloads and shrink row size. They solve different problems and are often used together. Choosing a shard key well up front matters more than almost anything else about a sharded system, because a bad key is expensive to discover late and disruptive to fix.
Structured elaboration
| Horizontal partitioning (sharding) | Vertical partitioning | |
|---|---|---|
| What's split | Rows, across nodes | Columns or tables, by function |
| Solves | Storage and throughput ceiling of one node | Row size, workload isolation (e.g., separating a hot auth table from a rarely-touched analytics table) |
| Typical key | A row-level shard key (user ID, order ID) | A functional boundary (which columns or tables belong together) |
| Does it reduce rows per node? | Yes, directly | No; a node with a subset of columns can still have every row |
Shard-key selection criteria
- Evenness: the key should distribute load uniformly across shards; a skewed key concentrates traffic on a subset of shards regardless of how many shards exist.
- Query-pattern alignment: colocate data that's typically accessed or joined together on the same shard, since cross-shard joins are expensive or impossible.
- High cardinality: a key with few distinct values (like a boolean flag) can't spread load across many shards no matter how it's hashed.
- Stability: a key that rarely or never changes for a given row avoids the operational cost of moving a row between shards when its key value changes.
Planning and executing a reshard
Prefer consistent hashing (with virtual nodes, so each physical shard owns many small hash ranges rather than one large contiguous range) over naive modulo-based hashing, because of how differently the two schemes behave when the shard count changes. Under modulo hashing (hash(key) mod N), changing N from 4 to 5 shards invalidates the mapping for nearly every key, since very few keys land on the same shard under both divisors. Under consistent hashing, adding the 5th shard to a 4-shard ring moves only the keys that fall between the new shard's position and its nearest predecessor:
versus, for the naive scheme:
N+1N=54=80%That gap (roughly a fifth of keys moving instead of roughly four-fifths) is the practical reason resharding plans lean on consistent hashing: the migration is proportional to the capacity added, not to the total dataset size.
The migration itself, in a cloud environment, typically runs online: stream ongoing changes from old to new shard layout (change data capture, CDC, capturing row-level changes as an ordered feed) while backfilling historical data in the background, validate the new layout against the old with checksums or row counts, then cut reads and writes over, ideally behind a feature flag or routing layer that can revert quickly if validation fails.
Operational challenges
- Rebalancing: moving data between shards is I/O- and network-intensive; throttle the migration and run it in the background rather than as a blocking step, and schedule any unavoidable cutover for low-traffic periods.
- Hotspots: even a well-chosen key can develop a hot shard as traffic patterns shift; mitigate with key-splitting (subdividing an overloaded key's range) or a cache layer absorbing read pressure in front of the hot shard.
- Migration windows: aim for a fully online migration so there's no hard window at all; if some downtime is unavoidable, keep it as short and as clearly scoped as possible, and communicate it like any other planned maintenance.
Worked example
A user-activity table sharded by user_id using consistent hashing across 4 shards needs to add a 5th shard as write volume grows. Using the calculation above, roughly 20% of keys need to move to the new shard (versus roughly 80% under naive modulo hashing), which is the concrete reason the migration is scoped as "move a fifth of the data," not "re-shuffle nearly everything." The team streams changes to the affected key ranges via CDC while backfilling their historical data, validates row counts on both old and new shard for the migrated ranges, then flips routing for those keys to the new shard, all without a hard maintenance window.
Trade-offs & pitfalls
- A monotonically increasing key (an auto-incrementing ID, or a timestamp) is a classic hotspot trap: every new row lands on whichever shard currently owns the "latest" range, so all write traffic concentrates on one shard no matter how many shards exist.
- Horizontal partitioning enables cross-node scale but makes cross-shard joins and multi-row transactions expensive or unavailable; that cost is inherent to sharding, not a bug to be engineered away.
- Vertical partitioning helps with row size and workload isolation but does nothing for a table whose row count is the actual bottleneck; picking the wrong partitioning axis for the actual problem wastes the migration effort.
- Consistent hashing reduces data movement on resize but doesn't guarantee even load by itself; a skewed key distribution can still produce uneven shard load even with a well-behaved hashing scheme.
When should you reach for asynchronous processing in a cloud architecture? Give examples of tasks that belong on a queue, the benefits and trade-offs (latency, throughput, complexity), and how you'd choose between a simple message queue, a pub/sub system, and a streaming platform for different workloads.
Sample Answer
Direct answer
Reach for asynchronous processing when a task doesn't need to block the user's response, has variable or long duration, benefits from being retried independently of the request that triggered it, or needs to scale at a different rate than the request path that produces the work. The trade-off is always the same shape: lower perceived user latency and higher burst-absorption capacity, in exchange for longer end-to-end completion time and real added complexity (idempotency, ordering, observability).
Structured elaboration
Tasks that belong on a queue: long-running jobs (video transcoding, PDF generation, model training), background work (email/SMS delivery, webhook retries, batch extract-transform-load jobs), rate-limited external calls (payment-gateway retries respecting a provider's own limits), and fan-out work (thumbnail generation, notification distribution to many recipients).
Benefits versus trade-offs:
| Dimension | Sync | Async |
|---|---|---|
| Perceived latency | Higher (user waits for the full task) | Lower (user gets an immediate acknowledgment) |
| End-to-end completion | Bounded by the request | Can be longer; work finishes on its own timeline |
| Throughput under burst | Limited by request-handling capacity | Higher; a queue absorbs bursts the request path can't |
| Complexity | Lower | Higher (idempotency, ordering, retry/dead-letter queue (DLQ), status tracking, monitoring) |
Choosing the messaging technology:
- Simple message queue (point-to-point task queue with retry semantics, one consumer per message): use for straightforward background jobs where ordering isn't required.
- Pub/sub (one message, many independent subscribers): use for event-driven fan-out across decoupled services, where each subscriber reacts to the same event differently.
- Streaming platform (ordered, replayable log, supports stateful consumers): use when ordering, replay, or exactly-once-style processing semantics matter, such as event sourcing or analytics pipelines.
Decision checklist: need ordering or replay, use streaming; need to fan out to many independent consumers, use pub/sub; need a simple worker queue with retries and nothing fancier, use a message queue.
Worked example
Migrating an existing synchronous pipeline to async: a concrete pilot. Say order-confirmation emails are currently sent synchronously inside the checkout request handler, adding real latency to checkout and coupling checkout availability to the email provider's uptime. The migration:
- Introduce a queue behind the checkout handler; the handler publishes an "order confirmed" message and returns to the user immediately instead of waiting on the email call.
- Define an idempotency key (the order ID) so that if the message is redelivered (a realistic possibility under at-least-once delivery), the consumer can detect and skip a duplicate send rather than emailing the customer twice.
- Decide on ordering: for this task, per-order ordering doesn't matter across different orders, so no partition-by-key requirement beyond what's needed for the idempotency check itself.
- Add a bounded retry (say, 3 attempts with backoff) and a DLQ for emails that still fail, so a transient provider outage doesn't silently drop confirmations.
- Pilot on a low-traffic checkout path first, and communicate the user-facing change explicitly: the email is no longer guaranteed to have been sent by the time the checkout page renders, which may require a small copy change ("confirmation email on its way" rather than implying it already arrived).
Precompute layer for a legacy synchronous aggregation service. A service that computes an expensive aggregate (say, a rolling usage summary) synchronously on every request can be converted to a precompute pattern: an async job recomputes the aggregate on a schedule or on a triggering event and writes it to a fast-read store; the read path serves from that store. On a cache miss (a query for a window that hasn't been precomputed yet, common right after rollout), the service falls back to computing it synchronously that one time and warms the cache with the result, so misses shrink over time instead of recurring. A one-time backfill job for historical windows at rollout avoids a long tail of slow cache-miss requests immediately after launch.
Batch versus online inference for ML serving. The same sync-versus-async trade-off shows up in model serving: batch (async) inference runs predictions over a group of inputs on a schedule or triggered job, giving higher throughput and lower cost per prediction at the expense of staleness equal to the batch interval; online (sync) inference serves a live request through the model on the request path, giving the freshest possible result but bounded by a per-request latency budget and a harder cost curve to scale at high request volume. The choice depends on whether the feature consuming the prediction can tolerate staleness measured in minutes or hours, or genuinely needs a fresh result per request.
Trade-offs & pitfalls
- Skipping idempotency is the single most common mistake in this pattern. At-least-once delivery is the realistic default for most queueing systems; a consumer that isn't safe to run twice on the same message will eventually double-charge, double-send, or double-write something.
- Ordering requirements are easy to miss until they bite. If two messages about the same entity can be processed out of order and produce a different, wrong final state, that entity's messages need to be routed to preserve order (typically via a partition key), which the simplest queue setups don't guarantee by default.
- The user-facing latency change needs to be communicated, not just implemented. Moving a task from synchronous to asynchronous changes what "done" means from the user's perspective; shipping that silently creates confusing, inconsistent-feeling product behavior even when the backend change is technically correct.
- A precompute layer with an unbounded cache-miss fallback can silently degrade into "synchronous, but slower and more complex." If misses stay common because a backfill never ran, or the precompute schedule can't keep up with query patterns, you've added the complexity of async without capturing its latency benefit.
Explain the basic queueing-theory concepts behind capacity planning: arrival rate, service rate, utilization, and the M/M/1 queue. Walk through a simple numeric example showing how a small increase in utilization can produce a disproportionate increase in average latency.
Sample Answer
Direct answer
Queueing theory formalizes something every engineer has felt intuitively: as a system gets busier, wait time does not grow in proportion to how busy it is, it grows much faster the closer utilization gets to 100%. The core quantities are the arrival rate (how fast work shows up), the service rate (how fast a server can finish work), and utilization (the ratio of the two). The M/M/1 queue is the simplest model that makes this precise, and its formulas show exactly why "we're at 80% capacity, that's fine" and "we're at 95% capacity, that's fine" are very different claims.
Core concepts
- Arrival rate (λ): the average number of requests arriving per second.
- Service rate (μ): the average number of requests a single server can complete per second, equivalently 1 divided by the mean time to handle one request.
- Utilization (ρ): the fraction of capacity in use, ρ=λ/μ for a single server. The system is only stable (queue length stays bounded over time) if ρ<1; at or above ρ=1, work arrives faster than it can be finished and the queue grows without bound.
- M/M/1 queue: a single server where arrivals follow a Poisson process (arrivals are independent and memoryless) and service times are exponentially distributed (also memoryless). "M" stands for "Markovian" (memoryless) on both the arrival and service side; "1" means a single server. This memorylessness is what makes the formulas below solvable in closed form, which is why M/M/1 is the standard first model even though real service-time distributions are rarely exactly exponential.
M/M/1 formulas (steady state)
ρ=μλAverage number of requests in the system (waiting plus being served):
L=1−ρρAverage time a request spends in the system (waiting plus service):
W=μ−λ1The denominator μ−λ is the key structural fact: as λ approaches μ, that denominator approaches zero and W diverges, which is the mathematical reason latency blows up near saturation rather than growing smoothly.
Worked numeric example
Take a single server with μ=100 requests/second.
λAλBλC=50=80=90ρAρBρC=50/100=0.5=80/100=0.8=90/100=0.9WAWBWC=1/(100−50)=0.02s=20ms=1/(100−80)=0.05s=50ms=1/(100−90)=0.10s=100msGoing from 50% to 80% utilization (a 1.6x increase in arrival rate) increases average latency from 20ms to 50ms, a 2.5x increase. Going from 80% to 90% utilization (only a 1.125x increase in arrival rate) still nearly doubles latency again, from 50ms to 100ms, a 2x increase on its own, and 5x relative to the 50% case. A modest-looking rise in utilization near the high end produces a disproportionate jump in latency; the same modest rise near the low end (say, 20% to 30%) would barely move the needle.
How this informs capacity planning
- Do not operate near ρ≈1. Pick a target maximum utilization based on the latency the service level objective (SLO) actually tolerates, commonly somewhere in the 50-70% range for workloads with real variability, leaving room for the nonlinear region above that.
- To hit a target average latency W, the required spare capacity is μ−λ≥1/W; this converts a latency target directly into a minimum service-rate requirement.
- Real service-time distributions are usually not exactly exponential, and real arrivals are usually not exactly Poisson (traffic bursts, batches, and diurnal patterns all violate the memoryless assumption), so treat the M/M/1 formulas as a directionally correct first estimate, then validate with a load test against the real distribution rather than trusting the closed-form number exactly.
- Multiple servers (an M/M/c queue) change the curve's shape but not the underlying lesson: latency still grows nonlinearly as aggregate utilization approaches saturation, it just takes more concurrent servers to keep the whole system away from that region.
Trade-offs and pitfalls
The most common mistake is reading "80% utilization" as "20% headroom, therefore safe," without accounting for how close that already is to the steep part of the latency curve; on this M/M/1 model, the jump from 80% to 90% utilization alone nearly doubles latency. A second pitfall is capacity planning entirely off average utilization while ignoring burstiness: even if average ρ looks comfortable, a service whose real arrivals come in bursts (not the smooth Poisson process the model assumes) can spend meaningful time at effective utilization far above the average, and M/M/1's steady-state numbers say nothing about that transient behavior.
Unlock Full Question Bank
Get access to all Scalability Patterns and Techniques interview questions and detailed answers.
Sign in to ContinueJoin thousands of developers preparing for their dream job.