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.
As an engineering manager, how would you coach a team through deciding where to draw microservice boundaries? Discuss bounded contexts, ownership, data ownership, coupling, how many teams are involved, deployment cadence, and operational cost. Walk through a concrete example of splitting a monolith into two services and your reasoning.
Sample Answer
Direct answer
As an engineering manager, I would not hand the team a boundary; I would coach them through a structured discussion of who owns what data and decisions today, where the coupling actually lives, and how many teams would be involved in a change, and let the boundary fall out of those answers rather than be drawn from a diagram first. The goal is a boundary the team can defend six months later, not just one that looks clean on a whiteboard.
Structured elaboration
1. Start from bounded contexts, not from code structure.
A bounded context is a part of the business domain with its own consistent vocabulary and rules (in an e-commerce system, "order" means something different to fulfillment than it does to billing, even though both teams say the word "order"). I'd have the team map where the domain's vocabulary actually changes meaning or where business rules diverge, since that is usually a much better boundary signal than where the code happens to have a file split today.
2. Anchor the boundary in data ownership, not just function.
Ask: who is the single source of truth for this data, and who else merely reads it? A service boundary that does not match a clear data-ownership boundary tends to produce two services that must coordinate on every write, which defeats the purpose of splitting them. If two candidate services both need to be the authority on the same data, that is a signal the boundary is drawn in the wrong place.
3. Weigh coupling honestly, in both directions.
Coupling shows up as call frequency (how often does one service need to call the other synchronously to do its job) and as change frequency (do these two areas of functionality tend to change together in the same pull requests, historically). High synchronous call frequency across a candidate boundary is a sign the split will trade in-process function calls for slower, less reliable network calls without a real decoupling benefit.
4. Size the boundary to the number of teams involved, not the other way around.
A service boundary that spans two teams' ownership creates a shared release calendar and a shared on-call surface, which erodes the main benefit of splitting in the first place: independent deployment cadence. I'd ask the team to name, concretely, who would own each side of the split going forward, because "we'll figure out ownership later" is the most common way a clean-looking split turns into a coordination tax within a quarter.
5. Weigh deployment cadence and operational cost against the coupling savings.
Every new service is a new thing to deploy, monitor, and be paged for. A boundary is only worth drawing if the independent-deployment benefit (teams shipping on their own schedule, without blocking each other) outweighs the added operational surface (a new service to run, new failure modes at the network boundary between them, new latency in the request path).
Worked example
A concrete coaching example I'd walk a team through: a monolithic order-management service currently handles both order placement (validating cart contents, applying pricing, writing the order record) and order fulfillment (routing to a warehouse, tracking shipment status, handling returns).
graph LR
subgraph Before[Monolith]
M[Order Management Service]
end
subgraph After[Split by bounded context]
O[Order Placement Service]
F[Order Fulfillment Service]
end
M -->|split| O
M -->|split| F
O -->|order created event| F
Walking the framework: order placement and fulfillment use the word "order" with different meaning (a placement-time order is a pricing and payment record; a fulfillment-time order is a shipment and inventory record), a bounded-context signal. Ownership splits cleanly: the placement team owns pricing and payment correctness, the fulfillment team owns warehouse and carrier integration correctness, and neither needs to be the source of truth for the other's data. Coupling across the boundary is low-frequency and asynchronous (fulfillment reacts to an "order created" event rather than calling placement synchronously mid-request), so splitting the network call out does not introduce a latency-sensitive dependency. Two teams already exist along roughly this line, so the split matches real organizational ownership rather than creating a new coordination requirement. That combination, not any single factor alone, is what makes this a good place to cut.
Trade-offs & pitfalls
The biggest trap I coach teams away from is drawing the boundary around a technical concern (like "the reporting queries are slow, let's split those out") instead of a bounded context, which tends to produce a service that still needs to synchronously coordinate with its sibling on every write and adds a network hop for no real autonomy gain. A close second is splitting before ownership is settled: a service with no single owning team accumulates a shared release calendar and shared on-call, which is a worse outcome than the monolith it replaced. My job as the manager is to slow the team down enough to check data ownership and coupling honestly before the boundary is drawn, since a boundary is far more expensive to undo after two teams have built around it than it is to get right up front.
Define edge caching and origin caching in plain terms for a cross-functional audience. For a photo-sharing app with 50 million daily active users and highly bursty traffic, which caching layers would you prioritize: CDN, regional caches, or application cache, and why? Briefly describe your invalidation strategy, how much staleness you'd accept, and the performance metrics you'd monitor.
Sample Answer
Direct answer
Think of caching as putting copies of data closer to the people who want it, at progressively larger "stores": a content delivery network (CDN) keeps copies at edge locations near users worldwide, a regional cache keeps copies in a handful of data-center regions, and an application cache keeps hot data in memory right next to the servers that handle requests. For a photo app with 50 million daily active users and bursty traffic, the priority order is CDN first, regional cache second, application cache third, because most of the cost and latency risk comes from serving the same popular photos over and over to a global audience, and the CDN is what absorbs that at the lowest cost per request.
Why this order, in plain terms
- CDN (edge caching), highest priority. Photos are, once uploaded, mostly unchanging. Storing copies at CDN points of presence around the world means a user in another country gets the photo from a nearby server instead of round-tripping to wherever the app's servers actually run. This is what makes bursty traffic (a post suddenly going viral) survivable: the CDN absorbs the spike instead of it hitting the origin servers directly.
- Regional caches, second priority. These sit closer to the origin than the CDN, in each region where the app runs servers. They catch requests the CDN missed (a photo nobody has viewed recently in that area) and reduce how often a request has to cross regions to reach wherever the primary data lives, which matters for both speed and cost.
- Application cache, third priority. This is in-memory data held directly by the app servers, mainly for things that change more often or need to be assembled per-request, like a user's session or a feed's metadata (like counts, captions), rather than the photo bytes themselves.
Invalidation strategy and how much staleness to accept
- The photos themselves: treat as immutable. When a user replaces a photo, give the new version a new URL rather than overwriting the old one in place; that sidesteps invalidation entirely, since the old URL simply stops being referenced. Cache these for a long time.
- Thumbnails and resized versions: shorter cache lifetimes, or serve a slightly stale version while a fresh one is generated in the background, since users rarely notice a resize regenerating a few seconds late.
- Metadata (likes, comment counts): short cache lifetimes and update the cache on write, since this changes constantly and users do notice when a like count looks frozen.
- Deletions and privacy actions: the one case that should never be "eventually consistent." If a user deletes a photo, that removal should propagate immediately, not wait out a cache expiration, because a deleted-but-still-cached photo is a privacy and trust problem, not just a UX nitpick.
As a rule of thumb: the more a piece of data resembles "a fact that was published once," the longer a cache can hold it; the more it resembles "a live counter or a permission decision," the shorter that window needs to be.
A concrete walkthrough: one photo, start to finish. Say a user uploads photo.jpg at 2:00:00 PM; it's cached at the CDN edge with a long time-to-live (TTL, how long a cached value stays valid before it's considered expired) since photo bytes are treated as immutable. At 2:05:00 PM that same user deletes the photo. The delete triggers an immediate invalidation call that purges the object from the CDN, the regional cache, and any application-cache entry referencing it, rather than waiting for the normal TTL to lapse. By roughly 2:05:02 PM, all three layers have confirmed the purge; a friend who opens that user's profile at 2:05:03 PM sees no photo at all, instead of the deleted image loading one more time from a stale edge copy. Contrast that with the like-count metadata next to the same photo: if it uses a 5-second cache lifetime instead of immediate invalidation, a viewer might briefly see a like count that is a few seconds behind reality, an acceptable trade-off for that specific piece of data, unlike the deleted-photo case above.
Metrics to monitor
- Cache hit ratio at each layer (CDN, regional, application): the single best signal that the tiering is doing its job.
- Origin request rate: should stay low and flat even during traffic spikes if the CDN and regional caches are absorbing load correctly.
- Latency at the 95th and 99th percentile (the response time that 95% and 99% of requests beat), since averages hide the slow outliers users actually complain about.
- How quickly a deletion or update propagates through the cache layers, since that is the metric that catches a privacy-invalidation bug before a user does.
Trade-offs and pitfalls
The main trade-off is staleness versus cost and speed: caching for longer serves more traffic cheaply but risks showing outdated content, while caching for a shorter time keeps things fresher but pushes more load back to origin servers, which is expensive at 50 million daily users. The most common mistake in this design is treating every kind of data the same way, for example, applying one blanket cache duration to both the photo bytes (safe to cache for a long time) and the like counter next to it (which looks broken if it is stale for more than a few seconds). The second most common mistake is not treating deletions as a special, urgent case: a cache design that is otherwise well-tuned for performance can still create a real privacy incident if a deleted photo keeps serving from cache for its normal time-to-live.
As an engineering manager, describe a simple capacity-planning approach for a service expected to grow 3x in traffic over the next 12 months. What inputs would you gather, such as current QPS and P95 CPU/memory per instance? Walk through the key calculations for forecasting instance or shard counts, and how you'd turn that forecast into hiring, infrastructure, or autoscaling decisions.
Sample Answer
Direct answer
Anchor the plan on a per-instance capacity number you can actually benchmark, not a guess: measure current queries per second (QPS, queries per second) and the 95th-percentile (P95, the value below which 95% of observations fall) CPU and memory per instance, project the 3x traffic target onto that per-instance capacity to get a target instance count, and only then work out what that delta costs in infrastructure spend versus what it costs in engineering time and headcount. Those are two different questions: "how many more instances" is usually a budget and autoscaling-configuration decision, while "does the architecture even support that many instances cleanly" is the one that turns into a hiring conversation.
Structured elaboration
Inputs to gather
- Current peak QPS and its trend over recent months, not just a single snapshot.
- P95 CPU and memory utilization per instance at current peak load; P95 rather than average, because average hides the moments the system is actually under stress.
- A benchmarked (not assumed) maximum sustainable QPS per instance, measured under realistic load, not theoretical hardware limits.
- Current autoscaling configuration: minimum and maximum instance counts, and how long a new instance takes to become ready (cold-start time), since that affects how much buffer you need above the bare-minimum forecast.
- Recruiting lead time for the team, if the forecast implies new engineering work rather than just more of the same infrastructure.
Key calculation
required instances=⌈QPS per instancepeak QPS×growth factor×(1+safety buffer)⌉
Assume, as a planning input rather than a measured fact, a current peak QPS of 3,000, a benchmarked capacity of 150 QPS per instance, a 3x growth target, and a 20% safety buffer for headroom above the raw forecast:
⌈1503,000×3×1.20⌉=⌈15010,800⌉=⌈72⌉=72 instances
For comparison, today's instance count under the same 20% buffer:
⌈1503,000×1.20⌉=⌈24⌉=24 instances
Instance count scales linearly with traffic here (from 24 to 72, a 3x increase matching the 3x traffic target), because per-instance capacity was held constant. That linearity check is itself useful: if the projected instance count did not scale roughly with the traffic multiplier, it would signal that something other than raw compute, a shared dependency like a database connection ceiling, is the real constraint, not instance count.
Turning the forecast into decisions
| Lever | What it addresses | When it's the right call |
|---|---|---|
| Autoscaling configuration | Routine, gradual demand within the existing architecture | The projected instance count fits comfortably within what the current design already tolerates; mostly a cost and configuration conversation |
| Infrastructure spend | Buying more of what you already run | The 72-instance target is a straightforward extension of the current stateless, horizontally-scaled design |
| New engineering work (headcount) | A structural limit the current design won't clear, for example a shared database that can't take 3x the connections, or a single component that isn't horizontally scalable | Profiling shows the bottleneck isn't instance count but a shared dependency; this needs a project (sharding, a caching layer, async processing) and a timeline, not just more servers |
If the forecast requires new engineering work, translate the estimated effort into a hiring ask against your team's actual recruiting lead time (commonly a few months for a senior engineer, a planning assumption you should validate against your own team's recent hiring, not a fixed constant) rather than assuming headcount can be added instantly once budget is approved.
Trade-offs & pitfalls
- The formula assumes per-instance capacity stays constant as load grows; if the bottleneck is actually a shared resource (a database, a single-instance cache, a rate-limited third-party API), adding instances past that point doesn't help and the linear projection will be wrong in a way the math alone won't reveal.
- Skipping the safety buffer and rounding down "to save cost" removes exactly the headroom meant to absorb the difference between a forecast and reality; a moderate buffer is worth its cost until you have data suggesting otherwise.
- Treating this as a one-time calculation rather than a recurring check misses the point: re-run it with fresh telemetry each quarter, because both the QPS-per-instance benchmark and the growth trend can shift as the product and traffic mix change.
- Converting a capacity gap directly into a headcount number without first checking whether it's actually an autoscaling or budget problem leads to over-hiring for what could have been solved by turning a dial.
You are a Technical Product Manager for a cloud developer platform. Define horizontal scaling versus vertical scaling in concrete terms, then give two product scenarios (one favoring horizontal, one favoring vertical) and explain the trade-offs in cost, downtime risk, operational complexity, observability, and developer experience. How would you influence engineering's choice, and what metrics would you monitor to validate it?
Sample Answer
Direct answer
Horizontal scaling means running more copies of a service side by side (more instances behind a load balancer) so the same work is split across a wider set of machines. Vertical scaling means making one existing machine bigger (more CPU, memory, or disk on the same box). As a technical product manager, the question I'd push engineering on isn't "which is better" in the abstract; it's "which one fits this specific service's constraints right now," because the two options carry very different cost, risk, and speed-to-ship trade-offs.
Structured elaboration
Definitions, concretely:
- Horizontal scaling: going from 1 app instance to 5 instances, each handling a fifth of the traffic, coordinated by a load balancer.
- Vertical scaling: taking that same single instance and moving it to a larger machine, for example doubling its CPU and memory.
Scenario favoring horizontal: a multi-tenant API serving many short, independent requests.
- Cost: higher baseline (more machines running), but better cost efficiency per request once traffic is high and steady.
- Downtime risk: lower; instances can be replaced one at a time without taking the whole service down.
- Operational complexity: higher upfront; needs load balancing and service discovery in place, which is infrastructure work, not a business decision by itself.
- Observability: needs request-level and fleet-level visibility (how is load distributed across instances), not just one machine's health.
- Developer experience: scaling is "add another instance," which is fast to execute once the infrastructure exists, but the service has to be stateless first (see below).
Scenario favoring vertical: a legacy or stateful component that can't easily be split, such as a single-process cache or an analytics-ingest service holding state in memory that isn't designed to be split across machines.
- Cost: a bigger machine has worse cost-per-unit-capacity at the high end, but it's often the fastest way to buy headroom without an engineering rewrite.
- Downtime risk: higher; resizing frequently requires a restart, and there's a single point of failure the whole time.
- Operational complexity: lower day-to-day (one machine to watch), but scaling further is capped by the largest machine available and harder to automate safely.
- Observability: narrower, focused on that one machine's CPU, memory, and health.
- Developer experience: no code changes required, which is attractive under deadline pressure, but it's a deferral of the real fix, not a substitute for it.
Signals that should trigger a move from vertical to horizontal, even if the team's instinct is to keep resizing:
- You're already near the largest machine size available, or the next size up costs disproportionately more for a shrinking capacity gain, so vertical simply runs out of room as a lever.
- A resize requires downtime, and that downtime window is now colliding with real user traffic instead of fitting inside a quiet maintenance period, meaning the "safe" vertical option has stopped being safe.
- Growth has become spiky rather than steady. A single bigger machine can absorb a slow, predictable climb, but it can't add capacity fast enough for short-lived spikes the way a set of instances that scale out and back in can.
User-visible impact during the transition. Moving from one big machine to several smaller ones is not free for users if the service was holding state in memory (a logged-in session, an in-progress upload). Unless that state is externalized to a shared store first, users can be logged out or lose in-progress work mid-cutover. There is also typically a short window of uneven response times while new instances warm up behind the load balancer, before their health checks stabilize.
How I'd influence engineering's choice:
- Translate the business need into concrete decision criteria: expected request volume, response-time targets, cost ceiling, and how soon this needs to ship.
- Ask directly whether the service is stateless (safe to run many identical copies) or stateful (holds data on one machine that would need to move first); this single question usually decides which path is realistic, more than a general cost debate does.
- Propose starting with the cheapest safe option (often vertical, if there's headroom left) while scoping the refactor that horizontal scaling requires, rather than treating it as an all-or-nothing choice.
- Get explicit agreement on a timeline: at what point does the team commit to the horizontal path even if vertical is still technically an option, so the decision doesn't get re-litigated every time a resize buys another few months.
Worked example
A concrete story: a developer-platform API starts on a single, reasonably large instance. Over several months, traffic grows steadily and the team resizes the instance twice, each time buying a few months of headroom with a short maintenance-window restart. On the third approaching resize, the team discovers they're already near the largest instance size the cloud provider offers for that machine family, and the next tier up costs far more for a proportionally smaller capacity increase. That's the vertical-headroom-exhausted signal firing. At the same time, product has just launched a feature that drives short, unpredictable traffic spikes around specific events rather than steady growth, which is the spiky-growth signal. Together, these push the team to invest in making the service stateless (moving session data out of the process and into a shared store) so it can run behind a load balancer as multiple instances, even though that refactor takes real engineering time the earlier vertical resizes didn't.
Trade-offs & pitfalls
- Horizontal scaling isn't free just because it's more "modern." It requires the service to be stateless first; skipping that step and scaling horizontally anyway produces inconsistent behavior (a user's session data only living on one of several instances) that is worse than staying vertical until the refactor is actually done.
- A string of "just one more vertical resize" decisions can quietly become the expensive path, if nobody is tracking how close the team is to the largest available machine size.
- The transition itself has a user-visible cost that's easy to leave out of the plan. Budgeting the refactor without budgeting for the state-externalization work, or without warning users about a rockier-than-usual cutover window, turns a well-reasoned architecture decision into a rough surprise for customers.
- The right metrics for validating this decision are the same ones that should have driven it: response-time percentiles (the response time under which a given percentage of requests complete), utilization per instance, cost per request, and how often scaling events happen; a decision that isn't being watched with these after the fact is a guess, not a validated choice.
A service uses in-memory caches to reduce database load, but it's experiencing cache thrashing during traffic spikes. As the engineering manager reviewing this with your team, what are the likely causes, and what fixes would you expect the team to propose? Consider eviction-policy tuning, tiered caching (edge, regional, origin), request coalescing, warm caches, and circuit breakers. How would you quantify the expected reduction in database load for a given improvement in hit rate?
Sample Answer
Direct answer
As the engineering manager reviewing this with the team, cache thrashing during a traffic spike almost always traces back to one of a few root causes: the cache is too small or the time-to-live (TTL) too short for the actual working set, the eviction policy is discarding items that are still hot, the cache starts cold right when the spike hits, or a thundering herd of requests floods the origin the moment a popular key expires. The fixes the team should propose map one-to-one onto those causes, and the value of any fix should be provable in terms of database load reduction before it's called done, not just "latency felt better."
Likely causes and the fixes to expect
- Cache too small for the working set, or eviction discarding hot-but-sparse keys: least-recently-used (LRU) eviction can wrongly evict an item that is popular but was not accessed in the last few seconds, in favor of something accessed once and never again. Tuning toward least-frequently-used (LFU) or a hybrid, or simply increasing cache size to match the real working set, addresses this directly. Adding admission filtering (only caching a key after it's been requested more than once) avoids wasting cache space on one-off keys in the first place.
- Tiered caching gaps: a single cache tier means every miss goes straight to the database. Edge, regional, and origin-level tiers each reduce origin hits multiplicatively; a gap in the tiering (for example, no regional cache) means the origin absorbs load that a middle tier should have caught.
- Request coalescing missing: without it, a popular key expiring under load produces one database query per concurrent reader instead of one query total. This is the specific mechanism behind a thundering herd, and it is usually the single highest-leverage fix when thrashing coincides with traffic spikes rather than steady-state load.
- Cold caches at spike time: if the cache was never warmed (after a deploy, a restart, or an autoscaling event added new instances with empty local caches), the first wave of spike traffic pays full origin cost. Predictive or scheduled pre-warming for known hot keys addresses this.
- No backpressure on the database itself: backpressure means signaling an overloaded downstream component, here the database, to slow down or reject new work instead of letting requests queue up indefinitely; without a circuit breaker (a guard that tracks recent failures and temporarily stops sending calls to a struggling dependency once a threshold is crossed, so it fails fast instead of piling on) or a hard cap on concurrent database calls, a wave of cache misses can take the database down entirely rather than just running slow, turning a performance problem into an outage.
Quantifying the expected database load reduction
Database load is, to a first approximation, proportional to the cache miss rate, so it scales with (1−hit rate). If a fix raises the hit rate from h1 to h2, the resulting database load relative to before is:
DB loadoldDB loadnew=1−h11−h2As an illustrative example (not a claim about this specific service, these are example inputs to show the method): if the team's current hit rate is h1=70% and the proposed fixes are projected to raise it to h2=90%:
1−0.701−0.90=0.300.10=0.333That means database load after the fix is expected to be about 33.3% of what it was before, a roughly 66.7% reduction. The team should measure the actual before/after hit rate in production (or in a load test with realistic traffic replay) and plug the real numbers into this same ratio, rather than trusting a projected hit-rate improvement without confirming it landed.
As the reviewing manager, what to ask for
- Instrumentation before the fix ships: per-key access distribution, miss-storm rate (misses per second, not just overall hit rate), and cache occupancy, so the team has a baseline to compare against.
- A pilot on one service or one traffic shard rather than a global rollout, so the hit-rate improvement (and any regression) shows up in a comparison the team can actually read.
- A specific number to report back: measured database queries per second before and after, converted into the same hit-rate ratio shown above, not just "it feels faster."
Trade-offs and pitfalls
The trap in this kind of review is accepting "we improved the eviction policy" as sufficient without a before/after measurement, since eviction tuning, request coalescing, and cache warming can each look individually reasonable while still leaving the actual bottleneck (often the missing coalescing) unaddressed. A second trap is prioritizing the fix that's easiest to implement (often eviction-policy tuning) over the fix with the largest leverage during a spike specifically (usually request coalescing, since it caps the herd effect directly rather than just improving steady-state hit rate).
Unlock Full Question Bank
Get access to all 9 Scalability Patterns and Techniques interview questions and detailed answers.
Sign in to ContinueJoin thousands of developers preparing for their dream job.