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.
For a read-heavy product catalog service, weigh the trade-offs between replicating a full cache to every region versus partitioning (sharding) cache entries by product or region. Consider read latency, cache-miss patterns, memory and network cost, consistency, and rebalancing complexity, then recommend an approach for a global retailer that sees traffic bursts from multiple regions.
Sample Answer
Direct answer
For a global retailer with bursty, multi-region traffic, neither pure full replication nor pure partitioning wins outright: full replication gives the best latency and simplest rebalancing but pays for it in memory and cross-region sync cost, while partitioning is cheaper but concentrates risk into hotspots when demand shifts. The right default is a hybrid: keep a small, region-local cache of the hottest slice of the catalog fully replicated in every region for latency, and back it with a sharded cache for the long tail, promoting items into the local cache when a region's traffic to them justifies it.
Structured elaboration
| Dimension | Full replication (every region holds the whole cache) | Partitioned (sharded by product or region) |
|---|---|---|
| Read latency | Best: any product is a local hit | Good only when the request lands on a local shard; a remote shard adds a cross-region hop |
| Cache-miss pattern | Only on first global write or expiry; predictable | Lower miss rate per shard for that shard's hot items, but a burst on one product can overload the single shard that owns it |
| Memory & network cost | High: full catalog held N times, one per region, plus cross-region invalidation traffic | Lower: no duplication of the catalog, and update broadcasts are smaller |
| Consistency | Async replication is simplest and typical; synchronous replication for strong consistency adds real latency | Simpler for the shard that owns a given item, since there's one writer path, but reads from other regions still need a remote call or a replication mechanism |
| Rebalancing complexity | Low: adding a region just means standing up another full copy | Higher: partition migrations and consistent-hashing-style reassignment are needed; hotspots require live re-sharding or targeted replication |
Why a global retailer with bursty traffic needs the hybrid, not either extreme
Bursty, multi-region traffic on a retail catalog is rarely uniform: a small set of products (a flash sale, a viral item) drive a disproportionate share of reads at any given time, and which products are hot can shift quickly. Pure partitioning puts that risk on a single shard, since consistent-hashing-style assignment (products and cache shards are placed as points on a circular hash space, so only nearby points move when shards are added; the mechanics of the hash ring itself are covered in more depth under load balancing's consistent-hashing pattern, and what matters here is the caching consequence) doesn't know a key is about to become hot until it already is. Pure full replication avoids that risk entirely but pays a flat memory and cross-region sync tax for the entire long tail of the catalog, most of which is rarely read in any given region.
The hybrid keeps region-local, fully replicated caches sized to each region's actual working set (the products that region's users actually read), backed by a sharded cache holding the full catalog. A traffic-based promotion rule (an item crossing a per-region hit-rate threshold gets pushed into that region's local cache) handles the shifting-hotspot case without requiring the whole catalog to be replicated everywhere.
Worked example
Assume, as a planning input rather than a measured fact, a product catalog sized at 50 GB, served across 6 regions.
Full replication cost:
50GB×6regions=300GB total cache memory
Partitioned cost (no duplication, split evenly across 6 shards, plus a replication factor of 2 within each shard for availability rather than for cross-region latency):
6shards50GB≈8.3GB per shard,50GB×2=100GB total with the availability replica
The partitioned approach uses roughly a third of the memory of full replication (100 GB versus 300 GB) at this illustrative catalog size. The hybrid sits between the two: if each region's working set is, say, 10% of the catalog (5 GB), replicating just that slice to all 6 regions costs:
5GB×6=30GB
on top of the 100 GB sharded backing store, for roughly 130 GB total, a fraction of full replication's 300 GB while still giving most reads (the ones hitting each region's working set) a local hit.
Trade-offs & pitfalls
- The hybrid's promotion rule needs a threshold and a demotion path; without demotion, the region-local cache grows unbounded as items get promoted but never removed, eventually approaching full replication's cost anyway.
- Cross-region invalidation is still required for the sharded backing store even in the hybrid; underestimating that traffic (versioned, pub/sub-style invalidation messages rather than synchronous broadcasts) is a common way the "cheaper" option ends up not being cheaper.
- A single globally hot product (a flash sale item) can still overload the shard that owns it even with promotion in place, if promotion reacts slower than the traffic spike; this is the scenario that specifically motivates proactive cache warming ahead of known events rather than purely reactive promotion.
- A content delivery network (CDN, a network of edge servers that cache content close to users) is a natural complement for static product assets (images, descriptions) but doesn't solve the dynamic pricing/inventory caching problem this comparison is about; don't conflate the two layers.
- Getting the region-local cache's time-to-live (TTL, how long a cached value is considered valid before refresh) too long trades staleness (wrong price or stock shown) for the latency win; too short and the hybrid starts behaving like the sharded-only design under load.
Describe a practical approach to capacity planning for a brand-new cloud service that has no historical traffic data. How would you make an initial workload estimate, decide on safety margins and headroom, plan for elastic capacity, and define the metrics and experiments you'd run to validate your assumptions after launch?
Sample Answer
Direct answer
With no historical traffic, you do not guess a single number: you build a workload estimate from comparable analogs and top-down business inputs, wrap it in an explicit safety margin, put it behind elastic capacity so the estimate does not have to be exact, and then replace the estimate with real data as fast as possible after launch through staged rollout and monitored experiments.
Structured elaboration
1. Build an initial estimate from two independent angles and reconcile them.
- Top-down: start from a business number you do have (invited users, marketing reach, sales pipeline) and multiply down to requests. This is the only lever available with zero history.
- Analog: find the closest comparable system you or the industry already operates (a similar feature, a similar-sized customer base, a similar product category) and scale its known request-per-user rate to your expected user count.
- Reconcile the two. If they disagree by more than roughly 2-3x, that gap itself is useful information: it tells you where your uncertainty is concentrated and what to instrument first.
2. Convert the estimate into a load shape, not just a total.
A daily total hides the number that actually threatens the system: peak requests per second (RPS, requests per second). Apply a peak-to-average ratio to account for daily cycles and, for a launch specifically, a possible synchronized spike (a launch email, a push notification, a press mention) that behaves nothing like organic steady traffic.
3. Set headroom deliberately, and say why.
Headroom on a zero-history estimate covers two different kinds of error: normal variance (traffic is noisier than a smooth average implies) and estimate error (the whole model could be wrong). Treat these as multiplicative: a peak-shape multiplier for the first, then a separate safety-margin multiplier for the second. Document both numbers as assumptions, not facts, so whoever revisits capacity later knows which parts were guessed.
4. Plan for elastic capacity so the estimate does not have to be right.
Because pre-launch numbers are inherently soft, favor a design where compute scales out automatically (for example an Auto Scaling group, ASG, sized with a low minimum and a generous maximum) over one where you provision a fixed fleet sized to the estimate. Stateless request handlers are what make this possible: any instance can pick up any request, so the ASG can add or remove capacity without session-affinity constraints. Identify the one component that will NOT scale elastically as fast as the rest (usually the database or a rate-limited third-party dependency) and size or protect that one deliberately, since it becomes the real ceiling regardless of how large the compute fleet grows.
5. Define what you will measure and how you will validate the assumption after launch.
Before launch, decide: the metrics that reveal reality (RPS, P95/P99 latency [95th-percentile/99th-percentile], error rate, queue depth, database connection saturation), the rollout mechanism that limits blast radius while those metrics come in (percentage-based ramp or canary release to a small traffic slice first), and the trigger for pausing the ramp (an explicit threshold on any of the above, decided in advance rather than improvised under pressure).
Worked example
Assume, as planning inputs rather than measured facts:
- 10,000 users are active on day one (from a marketing pre-registration count, discounted for expected activation rate).
- Each active user generates 15 requests over the day (from an analog product's per-user request rate).
- A peak-to-average ratio of 4x, reflecting a synchronized launch announcement rather than smooth organic arrival.
- A safety margin of 2x on top of the peak, to absorb estimate error since there is no history to validate the inputs against.
That "14 RPS" is not a forecast you defend, it is a starting point for the ASG's scaling policy and a number you replace with observed data within the first days of traffic.
Trade-offs & pitfalls
Over-provisioning a fixed fleet to the safety-margin number wastes money for a launch that may undershoot; under-provisioning without elastic headroom risks a visible outage on the day traffic is most scrutinized. The middle path (a small guaranteed baseline plus autoscaling) is usually right, but it only works if the service is stateless and the true bottleneck (often the database, not the request tier) is identified and protected separately, since databases scale far less elastically than compute. The most common senior-vs-junior tell is whether the candidate treats the initial number as a fact to defend or as an assumption to instrument and correct quickly after launch.
A single incoming request fans out to 50 parallel downstream calls. Each downstream call has a P95 latency of about 100ms, and the downstream system caps out at 1,000 RPS. If your service needs to handle 200 incoming RPS, is the downstream a bottleneck? Show your calculations, then propose architectural changes such as batching, caching, or queueing to reduce the downstream load, and explain the trade-offs.
Sample Answer
Direct answer
Yes, the downstream system is a clear bottleneck, by roughly 10x. At 200 incoming requests per second (RPS) with a fan-out of 50 calls each, the service generates 10,000 downstream calls per second, but the downstream system only accepts 1,000 RPS. Fixing this requires either reducing the number of downstream calls per incoming request (caching, batching, coalescing) or decoupling the response from the downstream work (async processing), not simply adding more capacity on your own side.
The math
Throughput check. Fan-out multiplies the incoming rate directly:
200 RPS×50 calls/request=10,000 downstream calls/s required
required:capacity=10,000:1,000=10:1
Required load is 10 times the downstream cap. That alone confirms a bottleneck: no amount of retrying or connection pooling on your side changes a system that is already saturated at its own ceiling.
Concurrency cross-check (Little's Law). It helps to sanity check the same conclusion from a different angle: how many downstream calls must be in flight simultaneously, not just per second. Little's Law relates throughput X and average time-in-system R to the average number of concurrent items N:
N=X×R
At the P95 latency (the response time that 95% of requests come in faster than) of about 100 ms (0.1 s), the concurrency the fan-out actually demands is:
Nrequired=10,000 RPS×0.1s=1,000 concurrent downstream calls
Whereas the downstream system, operating at its own stated cap with that same latency, is only structured to sustain:
Ncapacity=1,000 RPS×0.1s=100 concurrent downstream calls
Both views agree: you need about 10x the concurrency the downstream system is built to hold. This cross-check matters in an interview because it shows the bottleneck isn't just a rate-limit number on a dashboard, it is a real resource constraint (connections, threads, or queue slots) that a naive retry loop would make worse, not better, by piling on more concurrent attempts against an already-saturated system.
Mitigation options and what each one requires
Different mitigations close the 10x gap in different ways. It is worth deriving the minimum each one needs before choosing, rather than picking whichever sounds most familiar:
| Technique | How it reduces load | Minimum needed to close the gap | Key trade-off |
|---|---|---|---|
| Caching | Cache hits never reach downstream | Hit rate h such that (1−h)×10,000≤1,000⇒h≥0.9 | Staleness, invalidation complexity, only works for cacheable/idempotent reads |
| Batching or aggregation | Combines many logical calls into one downstream request | Batch factor b such that b10,000≤1,000⇒b≥10 | Adds wait-to-accumulate latency; downstream must expose a batch API (application programming interface, an endpoint accepting many items in one call) |
| Request coalescing (in-flight dedup) | Collapses concurrent identical requests into one call | No guaranteed factor; only helps if requests genuinely repeat the same key in a short window | Needs a singleflight-style layer (lets only the first caller for a key actually fetch it, while others waiting on that key reuse its result); zero benefit if requests are for distinct keys |
| Async queue with a worker pool | Decouples the caller's response from when downstream work completes | Does not reduce total required calls; still needs a sustained drain rate at or below 1,000 RPS or the backlog grows without bound over time | Higher end-to-end latency, needs a durable queue, changes the service-level agreement (SLA) from synchronous to eventual |
| Admission control / graceful degradation | Sheds or simplifies requests before they generate 50 downstream calls each | Reduces load by exactly whatever fraction is shed or simplified | Visible feature loss to some fraction of users |
The queueing row is the one candidates most often get wrong: a queue is a shock absorber for bursts, not a source of extra downstream capacity. If the arrival rate into the queue is sustained above the rate downstream can drain (1,000 RPS here), the backlog and its latency grow without bound over time; it only helps if the 10,000 RPS demand is a transient spike layered on top of a steady-state average that downstream can actually absorb.
Worked example: combining two mitigations
A single technique often has to hit an aggressive threshold alone (90% cache hit rate, or a batch factor of 10). Combining two moderate mitigations is usually more realistic. Assume, as an illustrative starting point (not a measured figure), a cache hit rate of 80% and a batch factor of 3 for the remaining traffic:
uncached calls=(1−0.8)×10,000=2,000 RPS
after batching by 3=32,000≈666.7 RPS
666.7 RPS is below the 1,000 RPS cap, with about 33% headroom. This is a useful pattern to point out explicitly: two moderate, individually achievable improvements (an 80% hit rate is realistic for many read-heavy access patterns; batching 3 calls together is a small API change) can beat needing one extreme, harder-to-sustain number from a single technique.
Trade-offs and pitfalls
- Treating "add a queue" as the fix without checking the sustained drain rate. A queue converts an overload into a growing backlog; it does not remove the overload.
- Choosing a hit rate or batch factor that meets the cap with zero margin. Production traffic is bursty and cache hit rates drift, so design to a threshold with headroom, not the exact breakeven point.
- Applying caching or batching uniformly across all 50 downstream calls when only some of them are actually cacheable or batchable in practice; the real achievable reduction is bounded by whichever calls are eligible.
- Retrying failed downstream calls without first fixing the 10x overload. Retries against an already-saturated system amplify load and can turn a slow degradation into a full outage.
- Skipping the concurrency cross-check. Throughput alone can hide a resource-exhaustion story (thread pools, connection limits) that shows up as timeouts before the RPS counter ever looks alarming.
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.
Design connection-pool management for a microservice that serves 500 concurrent request handlers and makes both database and upstream HTTP calls. How would you size the DB connection pool versus the HTTP client pool, choose timeouts and max lifetimes, integrate with circuit breakers, and test pool behavior under a simulated production spike?
Sample Answer
Direct answer
Size each pool from the request rate implied by the concurrency, not from the concurrency number itself: work out how many requests per second actually need a database connection or an outbound HTTP connection at steady state, multiply by how long each connection is held (Little's Law), then add an explicit safety margin. Timeouts must be shorter than the end-to-end service level objective (SLO), and connections need a maximum lifetime shorter than any upstream infrastructure's idle-connection timeout so sockets get rotated before they are silently dropped. A circuit breaker sits in front of each pool so a failing dependency fails fast instead of exhausting the pool while everyone queues for a resource that will not respond.
Sizing the pools
Step 1: convert concurrency to throughput. 500 concurrent handlers is a snapshot of how many requests are in flight, not a rate. Little's Law relates them: concurrency equals arrival rate times average time in system (L=λW). If the average end-to-end request latency is known (or targeted), the implied arrival rate falls out directly.
Step 2: split by dependency. Not every request touches the database, and not every request calls the same upstream. Measure (from application performance monitoring, not guesswork) what fraction of requests hit each dependency and how long that call takes, then apply Little's Law again, this time scoped to that dependency, to get the concurrent connections it needs.
Step 3: timeouts and lifetimes.
- Connect timeout: short, on the order of a few hundred milliseconds for calls inside the same network, since a connection that cannot even establish is not going to succeed on retry within budget.
- Request/read deadline: a slice of the overall SLO, sized so that even if this call takes its full budget, the caller still has time left to respond.
- Idle timeout: closes unused connections after they've sat idle for a while, freeing capacity for other requests.
- Max connection lifetime: rotates connections periodically so no single socket outlives an upstream load balancer's or NAT device's own idle-reap window, which would otherwise cause requests to fail against a connection that looks alive locally but was already dropped upstream.
- Pool checkout timeout: kept short. If a caller cannot get a connection quickly, failing fast and shedding load beats piling requests into an ever-growing queue.
Step 4: circuit breakers and bulkheads. Wrap each pool (or each upstream host) with its own breaker: once the failure rate or consecutive-failure count crosses a threshold, the breaker opens and callers get an immediate fallback or error instead of waiting out a timeout against a resource that is already failing. A bulkhead means giving a critical dependency its own pool, separate from a non-critical one, so a slow or saturated non-critical dependency cannot starve connections a critical path needs. The state-machine mechanics of breakers (closed / open / half-open, backoff between probes, how they interact with retries and cascading failure) are a deeper topic in their own right; what matters here is that the pool and the breaker are wired to the same dependency boundary, so a breaker trip and a pool-exhaustion alert are talking about the same failure.
Step 5: test it. Load test by ramping to the target concurrency and then spiking beyond it, watching pool wait time and saturation rather than just end-to-end latency. Separately, inject faults (added latency, forced errors) on the database and the upstream to confirm the breaker actually trips and the fallback path is exercised, not just that it exists in code. Run a soak test long enough to exercise the max-lifetime rotation at least once, to catch connection leaks that a short test would miss.
Worked example
Assumptions used below (explicitly stated, not measured facts):
- Concurrency L=500 handlers (given).
- Target average end-to-end latency W=200ms (an illustrative SLO target, not a measured value).
- 30% of requests make a database call, averaging 20 ms per call (illustrative).
- 50% of requests call an upstream HTTP service across 5 upstream hosts, averaging 50 ms per call (illustrative).
- Safety factor of 2x on the computed steady-state connection counts (illustrative, tune from observed variance).
Database pool:
DB call rate=0.3×2500=750 calls/s DB connections needed=750×0.02s=15With the 2x safety factor, target pool size = 30, capped at whatever the database server's own connection limit allows (if the DB caps total connections lower than 30 times the number of service instances, add a proxy/pooler in front of the database rather than growing the per-instance pool further).
HTTP client pool:
HTTP call rate=0.5×2500=1250 calls/s HTTP connections needed=1250×0.05s=62.5→63Spread across 5 upstream hosts: 63/5≈12.6→13 per host minimum, times the 2x safety factor = 26 per host target. (For HTTP/2 upstreams that multiplex many logical requests over one connection, this per-connection count would be far smaller; the formula is the same, only the effective "one connection handles many concurrent calls" changes the divisor.)
flowchart LR
C[Request handlers x500] --> CB1{Circuit breaker: DB}
C --> CB2{Circuit breaker: upstream HTTP}
CB1 -->|closed| DBPool[DB connection pool target 30]
CB1 -->|open| FailFast1[Fail fast / fallback]
DBPool --> DB[(Database)]
CB2 -->|closed| HTTPPool[HTTP client pool per host 26]
CB2 -->|open| FailFast2[Fail fast / fallback]
HTTPPool --> Upstream[Upstream service]
Trade-offs and pitfalls
Over-provisioning a database pool is not free: every open connection costs the database server memory and, for some engines, a background process, so "just make the pool bigger" is not a safe default. Under-provisioning causes callers to queue for a checkout instead of failing fast, which turns a capacity problem into a latency-cascade problem for every caller of the service, not just the ones touching the saturated dependency. A common ordering mistake is setting an inner call's timeout longer than the budget the outer caller has left, which wastes the outer request's remaining time waiting on a call that can no longer help it finish in budget; timeouts should shrink, not grow, as you move deeper into a call chain. Finally, a load test that only ramps to the steady-state target and never spikes past it will not catch pool exhaustion, since the whole point of a pool is that it behaves fine until it doesn't.
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.