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.
Describe stateless versus stateful service designs, and explain why statelessness enables easier horizontal scaling. Include strategies to externalize state (databases, caches, session stores), and discuss scenarios where a stateful service is genuinely necessary, for example leader election or long-lived sticky connections.
Sample Answer
Direct answer
A stateless service keeps no client- or session-specific data in process memory between requests: every request carries (or looks up) everything needed to handle it. A stateful service holds that context in memory across requests, such as an open socket or an in-memory session. Statelessness enables horizontal scaling because any instance can serve any request, so a load balancer can distribute traffic with no per-node reconciliation and instances become disposable: add, remove, or replace them freely.
Structured elaboration
Why statelessness is the enabler, mechanically
- Interchangeability: since no instance holds unique data, a request routed to any instance gets the same result, which is what makes simple, even load distribution possible.
- Painless scaling events: adding capacity means starting new instances with no data migration; removing capacity means stopping instances with no data loss, because nothing lived there that mattered.
- Simpler recovery: a crashed stateless instance is replaced, not repaired; there is no in-memory state to reconstruct.
- Simpler rolling deploys: instances can be cycled one at a time without draining session state first.
Strategies to externalize state
- Durable application state moves to a database (relational or NoSQL) instead of process memory.
- Fast shared access to hot data moves to a distributed cache (an in-memory key-value store shared across instances) rather than a per-instance local cache.
- Session data specifically has two common paths: a centralized session store, or client-held tokens (such as a signed JSON Web Token, JWT) that let the server stay stateless entirely because the client presents its own state on every request.
- Large binary objects move to an object store rather than local disk.
- Asynchronous or queued state (work in flight, not yet durable) moves to a message queue.
- Small amounts of shared coordination state (configuration, leader pointers, locks) move to a dedicated coordination service built for that job.
Where a stateful service is genuinely necessary
| Scenario | Why statelessness does not fit | Typical mitigation |
|---|---|---|
| Leader election | A cluster needs exactly one active decision-maker at a time; that role is inherently shared, coordinated state, not a per-request fact | Use a purpose-built coordination service to hold and arbitrate the leader pointer, and design fast, automated re-election on failure |
| Long-lived connections (WebSocket, gRPC streams) | The connection itself is state; a mid-stream request cannot be freely handed to a different node without breaking the stream | Route stream traffic with connection affinity (a load-balancing mechanic, out of scope here) and design the client to reconnect and resume cleanly on node loss |
| Low-latency in-memory computation | Externalizing every read adds network latency that some workloads cannot absorb | Keep a local cache as a performance optimization layered on top of an externalized source of truth, not as the only copy |
| Stateful stream processing (windowed aggregation) | The computation's correctness depends on an accumulating window that must survive across events | Checkpoint local state to durable storage continuously so a replacement node can resume from the last checkpoint instead of from scratch |
The SRE angle: operating a stateful service you cannot avoid. When a stateful service is genuinely required, the operational burden shifts from "how do I scale it" to "how do I keep it reliable despite being pinned." That means treating the stateful nodes as a smaller, more carefully managed subset of the fleet: automated health checks tied to a fast failover or re-election path, regular drills that exercise that failover before it is needed under real incident pressure, and monitoring that specifically watches for the failure modes unique to statefulness (a stuck leader, a connection that never drains, a checkpoint that falls behind). The mitigation is never "make it stateless anyway"; it is "shrink the stateful surface to the minimum and instrument that minimum heavily."
Worked example
A web application currently stores logged-in session data in each server's memory, so a user's second request must land on the same server (sticky routing) or their session appears to vanish. To make the tier horizontally scalable: move session data out of process memory into a shared session store, or better, switch to a signed token that the client holds and presents on each request so the server does not need to look anything up at all. Either change means any instance can now serve any request, sticky routing is no longer required, and the fleet can be scaled up or down purely on load, with new instances immediately able to serve full traffic with zero data migration.
Trade-offs & pitfalls
- Do not confuse "stateless service" with "no state exists": the state still exists, it has just moved to a system designed to hold it reliably. That external system becomes a new dependency and a new potential bottleneck, so it needs its own scaling plan.
- Client-side tokens remove server lookups but push size and revocation concerns onto the client and the token design; server-side session stores keep revocation simple but add a network hop and a shared-store scaling problem.
- Treating "sticky sessions" as a free fix for statefulness is a common wrong turn: it works, but it silently reintroduces a form of per-node state ownership and undermines the failure-isolation benefit that horizontal scaling was meant to provide.
- For the genuinely stateful cases, skipping the failover-drill discipline is the most common production failure: the coordination logic works in testing and then fails silently the first time it is needed in an incident, because it was never exercised under realistic conditions.
Your service is stateful and stores per-user ephemeral session data in memory. Walk through three practical approaches to scale it horizontally without breaking correctness: sticky sessions, active-active replication, and externalizing the state. For each, describe the operational complexity, failure modes, and performance characteristics.
Sample Answer
Direct answer
There are three practical ways to scale a stateful service that keeps per-user ephemeral data in memory: keep the state local and pin each user to one instance (sticky sessions), replicate the state across instances so any of them can serve a user (active-active replication), or move the state out of the instances entirely into a shared store (externalization). Sticky sessions are the cheapest to build but the most fragile; externalization is the most operationally mature and is the default recommendation for anything beyond a small fleet; active-active replication is the highest-complexity option, justified mainly when you need both local-memory latency and multi-node resilience together.
Structured elaboration
1. Sticky sessions (session affinity at the load balancer)
- How it works: the load balancer hashes a stable client identifier and routes all of that client's requests to the same instance, which keeps the session in local memory. The routing mechanism itself (which hashing scheme, which health checks) is load-balancer territory; what matters for scaling the service is the consequence: correctness now depends on that one instance staying up.
- Operational complexity: low to moderate. Configure affinity on the load balancer, size instance capacity for peak per-node session count, add health checks and connection draining so deploys don't silently drop sessions.
- Failure modes: an instance failing loses every session pinned to it unless you add a fallback (short-time-to-live (TTL) backup write, or accept the loss and force re-authentication). Autoscaling and rolling deploys cause rebalancing, which cold-starts sessions on whichever instance a client lands on next. Hot users concentrated on one instance cause uneven load that autoscaling on aggregate metrics won't catch.
- Performance: the best case for latency, since reads and writes never leave process memory. Effective capacity is limited by the least-loaded-instance problem: total fleet capacity is high, but any single hot instance can be saturated while others are idle.
2. Active-active replication
- How it works: session state is synchronously or asynchronously replicated across a set of instances (via a gossip protocol, where nodes periodically exchange state with random peers; a consensus protocol, where nodes vote to agree on a value before it counts as committed; or conflict-free replicated data types, data structures built to merge concurrent updates without conflicts) so more than one node holds a current copy and can serve a read or accept a write.
- Operational complexity: high. You are now running a small distributed system inside your application tier: membership management, conflict resolution for concurrent writes to the same session, and monitoring for replication health.
- Failure modes: network partitions are the central risk. Without a consensus mechanism enforcing a single writer or a quorum, a partition produces split-brain, where two nodes both believe they hold the authoritative copy of a session and diverge. Even without a partition, replication lag means a read on one node can be stale relative to a very recent write on another.
- Performance: write latency goes up (synchronous replication waits on other nodes; asynchronous replication is fast but reintroduces the staleness problem it was meant to avoid). Read latency can stay low if reads are served from local replicas. Write throughput is bounded by replication and, if used, consensus overhead.
3. Externalizing state (shared cache or durable store)
- How it works: instances stop holding session state at all. Every request reads from and writes to a shared store (Redis, Memcached, or a durable database), which makes every instance interchangeable.
- Operational complexity: moderate. You now operate one more piece of infrastructure (a cache or database cluster: sizing, eviction policy if it's a cache, persistence if it needs to be durable) and the application changes to fetch and save state instead of holding it.
- Failure modes: the store itself becomes the thing that must not go down, so it needs its own high availability, replicas, and a defined fallback if it's briefly unreachable. Cache-invalidation bugs (serving stale state after a write, or evicting state that should have survived) are the main correctness risk. Network dependency between every app request and the store adds a new failure surface that didn't exist with local memory.
- Performance: a network hop per read/write adds latency relative to local memory, but it is predictable and independent of which instance serves the request, so it scales cleanly with the store's own capacity. A local read-through cache in front of the store recovers most of the latency for repeat reads.
Worked example
Consider a live-chat product where each open chat session tracks typing indicators, unread counts, and the last 20 messages in a rolling buffer, all ephemeral and rebuilt from durable message history if lost. The team starts with sticky sessions during the MVP: single-digit instance count, acceptable risk of a lost typing-indicator state on deploy. At 30 instances and daily autoscaling swings, rebalancing during scale-events starts dropping a visible fraction of active typing indicators on every scale-up, and hot users (large group chats) overload single instances that aggregate metrics don't flag as unhealthy. The team considers active-active replication to keep the local-memory latency but rejects it: the correctness bar for typing indicators is low (occasional staleness is invisible to users), so the complexity of consensus or CRDTs isn't earning its keep. They externalize instead: session state moves to a Redis cluster, keyed by chat id, with a short TTL so abandoned sessions clean up automatically. Every instance becomes interchangeable, rebalancing no longer drops state, and the added latency (one Redis round trip per message event) is well under the product's tolerance for a chat UI.
Trade-offs & pitfalls
- Choosing active-active replication for its latency profile without a real need for both local-memory speed and multi-node write availability at once is the most common overreach; most services that reach for it would be better served by externalization plus a local read-through cache.
- Assuming sticky sessions "just work" at scale and only discovering the rebalancing-churn problem after autoscaling is already live is a frequent operational surprise; load-test the rebalance path, not just steady-state.
- With externalized state, forgetting that the store is now a dependency of every request (not an optional side-system) leads to under-provisioning its availability relative to the application tier it's supposed to be simplifying.
- With active-active replication, treating eventual consistency as "good enough" without checking whether the specific session data has a correctness requirement (e.g., a shopping cart total versus a typing indicator) can silently introduce bugs that only show up under network partition, which is exactly when they're hardest to debug.
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.
You need to vertically scale a production stateful database (increase CPU and memory on the primary instance) while minimizing downtime and preserving data consistency. Walk through the runbook you would execute: pre-checks, rolling steps, fallback options, and monitoring to verify success. Assume cloud-managed instances and the ability to create a temporary read replica to help with the cutover.
Sample Answer
Direct answer
Use the temporary read replica as the mechanism that turns an in-place resize (which can mean real downtime) into a controlled cutover: provision the replica already at the larger CPU/memory spec, let it fully catch up to the primary, briefly pause writes, promote the replica to primary, and repoint the application. The actual write-unavailability window is bounded by how long it takes to drain in-flight writes and flip the connection target, not by the resize operation itself, which is why this pattern minimizes downtime even though it isn't strictly zero-downtime.
Structured elaboration
Pre-checks
- Confirm the exact target CPU/memory spec against measured load, not a guess, and confirm a maintenance window and a communicated service level objective (SLO, the measurable target for allowed downtime or latency impact) for the operation.
- Take a fresh on-demand snapshot immediately before starting, independent of the replica strategy, as a last-resort fallback.
- Verify current replication lag baseline and that the environment supports creating a same-region replica sized larger than the current primary.
- Confirm the automation (scripts, IaC) for promotion and connection-string cutover has been tested outside of this incident, not written live.
Rolling steps
- Create a read replica provisioned at the new, larger instance spec. Let it catch up and monitor replication lag until it's negligible and stays that way for a sustained period, not just a single low reading.
- Briefly quiesce writes on the current primary (put the application into a short read-only or write-paused mode).
- Promote the replica to primary. Because it was fully caught up at the moment of promotion, this preserves the data that existed at quiesce time.
- Repoint the application's write target to the newly-promoted primary (a connection-string or routing change, ideally something that doesn't require an application redeploy).
- Resume writes and run smoke tests against critical read and write paths.
- Rebuild redundancy: the original (smaller) instance can be resized and re-added as a replica, or replaced, restoring the topology's normal read-replica count.
Fallback options
- If the replica fails to catch up before the maintenance window closes, abort the promotion; investigate whether replication is network- or I/O-bound, and either wait for a longer window or address the bottleneck before retrying.
- If a data mismatch or unexpected inconsistency is detected after promotion, the fallback is the pre-operation snapshot, not the old primary (which may now be behind); restore from snapshot to a fresh instance if this happens.
Monitoring
- Before: replication lag trend over a real observation window, not a single point-in-time check, plus baseline CPU/memory/connection counts to compare against post-cutover.
- During: replication lag right up to the promotion moment, since promoting a replica that's meaningfully behind means losing whatever writes happened after its last applied transaction.
- After: error rates, write and read latency, and a targeted data check (row counts or checksums on a few critical tables, or confirming the most recent known transactions are present) rather than assuming success from the absence of alarms.
Worked example
A safe promotion policy might require replication lag to stay under 1 second for a sustained 5-minute window before promotion is allowed (a stated operational threshold for this runbook, not a universal rule that applies to every workload). The actual write-pause duration during cutover isn't a fixed number worth quoting as a general fact, since it depends on the application's connection pool behavior, not on the database resize itself: it's bounded by how long the app takes to drain in-flight writes and how quickly its clients reconnect to the new endpoint after the connection target flips, which is exactly why this pattern is described as minimizing downtime, not eliminating it. If the application's reconnect and retry logic is slow or missing, the same database-side runbook produces a much longer perceived outage even though the database steps themselves didn't change.
Trade-offs & pitfalls
- This is not truly zero-downtime: promoting a replica that isn't fully caught up loses whatever writes landed after its last applied transaction, so the safety of the whole procedure hinges on verifying lag is genuinely near zero at the moment of promotion, not assuming it.
- If the application cannot tolerate even a brief write-pause, this pattern isn't sufficient on its own; a true zero-downtime requirement needs a different approach, such as a proxy layer that queues writes during cutover.
- Before building a custom replica-promotion runbook, check whether the cloud provider's native "modify instance class" operation already performs an equivalent internal promote-and-swap; if it does, it may be simpler and better-tested than a hand-rolled version of the same idea.
- The replica-promotion mechanics here (catching up, promoting, cutting over) are a practical means to a scaling end; the deeper mechanics of replication modes and failover consensus are a related but distinct topic from the scaling procedure itself.
Compare range-based, hash-based, and directory-based sharding strategies for partitioning write-heavy user data. Discuss the pros and cons of each in terms of hotspot formation, rebalancing complexity when adding nodes, support for range queries, and impact on secondary indexes and transactions.
Sample Answer
Direct answer
Range-based sharding groups keys into contiguous ranges (good for range queries, prone to hotspots on the newest range); hash-based sharding spreads keys evenly by hashing them (avoids that hotspot but destroys ordering, so range queries scatter across every shard); directory-based sharding uses a lookup service to map each key to a shard explicitly (most flexible for steering around hotspots, at the cost of operating and scaling that lookup service). For write-heavy user data, the right starting point is almost always hash-based, layered with directory-style overrides for the specific keys that turn out to be hot in practice.
Structured elaboration
At the simplest level, the choice comes down to one question: does your access pattern need range queries, or does it need even load distribution? Those two goals are usually in tension, which is why the three strategies exist.
| Dimension | Range-based | Hash-based | Directory-based |
|---|---|---|---|
| Hotspot formation | High risk: a shard holding the newest or most-active range absorbs disproportionate write traffic | Low risk for aggregate load, since keys spread evenly; a single very-hot key still hits one shard regardless of hashing | Lowest risk in practice: a hot key can be explicitly moved or replicated to relieve its shard |
| Rebalancing when adding nodes | Moderate to costly: splitting a range means physically moving a large contiguous block of data | With plain modulo hashing, adding a node reshuffles most keys; with consistent hashing, only keys near the new node's position move | Cheapest to reason about: update the directory entry for the affected keys; the underlying data still has to move, but only for the keys chosen, not a bulk range or hash-bucket shift |
| Range queries | Efficient: a scan over an ordered key range touches only the shards holding that range | Poor: hashing destroys order, so a range scan has to fan out to every shard and merge results | Depends on the directory's mapping scheme: efficient if it preserves contiguous mappings, expensive if mappings are arbitrary |
| Secondary indexes | Simple if the index key correlates with the shard range; otherwise needs scatter-gather (fan the query out to every shard, then merge the partial results into one answer) | Harder: a global secondary index generally needs either per-shard indexes aggregated at query time, or a separate index service | Can centralize index ownership through the same directory, but that adds another responsibility to the lookup service |
| Transactions | Good for range-local transactions; anything spanning ranges needs a distributed transaction | Single-key transactions stay on one shard; multi-key transactions usually span shards unless keys are deliberately co-located | Best control: related keys can be deliberately mapped to the same shard to keep transactions local, but only as reliably as the directory is kept consistent |
Alternate shard-key dimensions. User ID is not the only reasonable key for user data. Two common alternatives, each changing the hotspot and locality picture:
- Partitioning by tenant (for a multi-tenant product): every row for one customer organization lives together, which keeps that customer's queries and transactions shard-local, but a single very large tenant can become a hot shard on its own regardless of hashing, and tenant sizes are rarely uniform.
- Partitioning by geographic region: keeps a region's data physically close to the users generating it and can satisfy data-residency requirements, but regions do not generate even load (one region's business hours are another's overnight lull), so this dimension needs the same hotspot vigilance as range-based sharding by time.
Either alternative needs a request-routing design that resolves a given request to the correct partition before any query executes: for tenant or region keys this is usually a straightforward lookup (tenant ID and region are typically known before the query is issued), but it is still an explicit routing decision, not something a plain hash of an opaque ID gives you for free.
A simpler way to hold the whole comparison in your head: shard-key choice is really a load-distribution and data-locality decision. If you need data that is accessed together to live together (range scans, tenant-local transactions, region-local latency), you are choosing locality over even distribution, and range- or directory-based sharding fits. If you need writes spread as evenly as possible and do not need that locality, hash-based sharding fits.
Worked example
A write-heavy user-activity table (one row per user action, high insert rate) is first sharded by signup-date range, since that was the easiest key to reach for. Under load, the shard holding "this week's" range becomes visibly hotter than the others, since all new activity lands there, exactly the hotspot risk in the table above. Switching the shard key to hash(user_id) spreads new writes evenly across all shards, because a hash has no notion of "this week." The cost: a query like "all activity in the last 7 days across all users" now has to fan out to every shard and merge results, whereas it used to touch one or two. If that query is rare and the write hotspot was the actual production problem, the hash-based trade is the right one; if that range query is core to the product, a directory-based approach (route the most recent, hottest range's writes across several shards explicitly, while keeping most historical data range-organized) is worth the added operational complexity instead.
Trade-offs & pitfalls
- Consistent hashing (not plain modulo hashing) is what makes hash-based sharding practical to rebalance; naive
hash(key) % Nreshuffles nearly everything whenNchanges, which erases the rebalancing advantage hashing is chosen for. - Directory-based sharding trades data-movement cost for a new single point of scaling and consistency risk: the directory service itself must be highly available and kept correct, or every lookup through it is wrong.
- Secondary indexes are the detail teams most often discover too late: a design that looks fine for primary-key access can require an expensive scatter-gather or a whole extra index service the moment a "find by email" or "find by status" query shows up.
- Choosing tenant or region as the shard key without checking size distribution first is a common wrong turn: one outsized tenant or one high-traffic region can silently become the new single hot shard, defeating the purpose of sharding at all.
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.