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.
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.
As the lead backend engineer, you must choose between three short-term options to cut P95 latency by 30% within a fixed budget: a vertical database upgrade, adding read replicas, or introducing a caching layer. What metrics and profiling steps would you use to evaluate each option? Describe your experimental rollout (A/B or canary), rollback plan, and the long-term maintainability implications of each choice.
Sample Answer
Direct answer
Under a fixed budget and a 30% 95th-percentile (P95, the latency value below which 95% of requests complete) target, I would profile first to find which resource is actually saturated, then pick the option whose lever matches that bottleneck rather than defaulting to the biggest hammer: a caching layer for read-heavy, repeatable queries; read replicas when the primary is saturated by read volume specifically; and a vertical upgrade only when the workload is genuinely resource-bound with no obvious inefficiency to fix first. Whichever option is chosen, I would roll it out behind a canary with an explicit rollback trigger, because a change aimed at cutting P95 by 30% is exactly the kind of change that can regress tail latency if the assumption behind it is wrong.
Structured elaboration
Metrics and profiling per option
- Vertical database upgrade: check CPU utilization, I/O wait, disk throughput, active connections, and slow-query logs (e.g.,
pg_stat_statements) to confirm the primary is resource-saturated rather than running inefficient queries that a bigger machine won't fix. - Read replicas: check the read-to-write ratio, replication lag under current load, and which endpoints are read-dominated. This option only helps if reads, not writes, are what's driving primary saturation and P95.
- Caching layer: check which endpoints repeat the same query for many requests (cacheability), current origin queries per second (QPS, queries per second) on those endpoints, and how tolerant the data is of a short time-to-live (TTL, the duration a cached value is considered valid before it must be refreshed).
Decision framework
- If profiling shows CPU or I/O saturation on otherwise well-optimized queries: vertical upgrade is the fastest lever, but it has a hard ceiling and doesn't reduce load, it just buys headroom.
- If reads dominate write volume and the primary's read load is the driver: read replicas are the natural fit, since they scale read capacity horizontally instead of scaling one machine up.
- If a meaningful share of requests are repeatable (same query, same or slowly-changing result): caching usually gives the largest P95 improvement per dollar, because it removes load from the database entirely rather than adding capacity to serve it.
Rollout: canary and rollback
Route a small percentage of traffic to the changed path first (a canary), and compare P95, error rate, and (for replicas) replication lag against the unmodified baseline before widening. The specific traffic-shifting mechanics, how you carve off exactly 5% of requests at a load balancer, are the same progressive-delivery machinery used for any staged rollout; the part specific to this decision is what you measure and what threshold triggers a rollback. Define the rollback trigger before starting: for example, P95 regressing beyond baseline, or error rate rising above an agreed service-level objective (SLO, an internal target for how the system should perform). Rollback should be a single reversible action: a feature flag to bypass the cache, a connection-pool switch back to the primary-only read path, or a DNS/config revert to the pre-upgrade instance, not a multi-step manual procedure under pressure.
Worked example
Assume, as a planning input rather than a measured fact, a current P95 of 900 ms. A 30% cut targets:
900×(1−0.30)=630 ms
Now evaluate the caching option concretely. Assume origin traffic on the targeted read endpoints peaks at 5,000 QPS, and the caching layer reaches an 80% hit ratio there (an illustrative target, not a guarantee: it must be validated against real access patterns before committing budget to it). Origin load after caching:
5000×(1−0.80)=1000 QPS reaching the database
That is an 80% reduction in database load on those endpoints, which is the mechanism that produces the P95 win: most of the remaining latency on a cache hit is the cache round-trip, not the database.
Now the read-replica option on the same traffic. Assume 75% of primary traffic is reads and 25% is writes. Moving all cacheable reads to replicas removes up to:
5000×0.75=3750 QPS off the primary
leaving the primary serving roughly 1,250 QPS of writes plus any reads that must stay on the primary for freshness reasons. Both options move a comparable share of load off the critical path here; the real decision hinges on which one the profiling data actually supports, and caching wins on cost when hit rate is achievable, while replicas win when the workload isn't cacheable but is still read-dominated.
Trade-offs & pitfalls
- A vertical upgrade is the easiest to roll back (revert to the old instance size) but has a hard scaling ceiling and doesn't reduce load on the system, it just delays the next capacity conversation.
- Read replicas introduce replication lag: a client that writes and immediately reads its own write from a replica can see stale data unless read-your-writes routing is handled explicitly. That routing logic is new code that has to be maintained.
- A caching layer is usually the biggest win per dollar when hit rate is high, but it adds an invalidation problem: incorrect TTLs or missed invalidation on writes create a second source of truth that can silently diverge from the database.
- Picking the option that matches the profiling data, not the one that is fastest to implement, is what separates a durable fix from a fix that gets undone in the next incident review. All three options are non-exclusive long-term: a mature system typically ends up using all three, sequenced by where the bottleneck actually was.
Why does connection pooling matter for a service running at scale? Describe best practices for managing both database and HTTP connection pools: pool size, max open connections, idle timeouts, connection lifetime, and behavior under a spike in load. How would you test and tune these settings before production?
Sample Answer
Direct answer
Connection pooling matters at scale because opening a new database or HTTP connection is expensive relative to a request (TCP handshake, and for a database, authentication and session setup), so reusing a small set of warm connections instead of creating one per request lowers latency and prevents the backend from being overwhelmed by connection churn. The core sizing problem is that a pool is a per-instance setting but the backend has a fleet-wide connection ceiling, so pool size has to be planned across the whole fleet, not tuned in isolation on one instance.
Structured elaboration
Why pooling matters at scale, mechanically
- Connection setup cost: a TCP handshake, TLS negotiation (for HTTP), and for a database, authentication plus session/state initialization, all add latency if paid on every request.
- Backend resource limits: every open connection holds memory and, for a database, often a whole backend process or thread; a backend with a hard maximum connection count can be pushed into refusing connections or degrading badly under connection churn even if query volume itself is modest.
- Reuse turns a per-request cost into a one-time cost amortized across many requests on the same warm connection.
Sizing pools: the fleet-wide constraint
The number one mistake is sizing a pool as if the instance owns the whole backend. It does not; every other instance is drawing from the same ceiling:
pool_size_per_instance≤⌊number_of_app_instancesdb_max_connections⌋If the database allows 500 total connections and the service runs behind 20 instances, each instance's pool must stay at or below ⌊500/20⌋=25 connections, or a fleet at full pool utilization exceeds the database's ceiling and starts getting connection refusals, exactly when load is highest and refusals hurt the most. This constraint must be revisited every time the fleet is resized by autoscaling, which is the part teams most often forget: a pool size tuned for 20 instances silently becomes unsafe the moment autoscaling adds a 21st.
Core pool parameters
| Parameter | What it controls | Tuning guidance |
|---|---|---|
| Pool size (min/max) | How many connections are kept open per instance | Bounded above by the fleet-wide formula above; bounded below by enough to avoid queuing under normal load |
| Max open/concurrent connections | Hard ceiling the pool will not exceed even under burst demand | Set to protect the backend, not just to satisfy the busiest moment; excess demand should queue or fail fast, not force more connections open |
| Idle timeout | How long an unused connection stays open before being closed | Long enough to avoid re-opening connections for normal traffic gaps; short enough to release resources during genuine lulls |
| Max connection lifetime | Forces a connection to be recycled after a set duration regardless of use | Keeps the pool from silently holding stale or half-broken connections open indefinitely; also spreads out reconnections instead of all connections expiring together |
| Acquisition timeout | How long a request will wait for a pooled connection before failing | Should fail fast rather than block indefinitely, so an overload turns into fast, visible errors instead of a pile of hung requests |
Behavior under a load spike: the connection-storm problem
The specific failure mode worth naming: a deploy, a failover, or a sudden traffic spike can cause many instances to simultaneously reconnect or spin up new pooled connections at once, a connection storm, which can itself exceed the database's connection ceiling even though steady-state pool sizing was correct. This has a process/thread-model dimension too: a backend that spawns one OS process or thread per connection (a common relational-database architecture) pays a much higher per-connection memory and context-switch cost under a storm than one built around lightweight connection handling, which changes how conservatively you should size db_max_connections in the first place. Mitigations: stagger reconnects with jitter (small random delays) instead of reconnecting all instances at once, keep pool warm-up gradual rather than instantaneous on instance startup, and prefer acquisition timeouts with backoff over unbounded retry storms.
Testing and tuning before production
- Load-test at realistic peak concurrency and burst shape, not just average throughput, since spikes and connection storms are what actually break pool sizing.
- Vary the number of app instances in the test to confirm the fleet-wide formula holds at the target autoscaling range, not just at today's instance count.
- Watch active/idle/wait-count and wait-time metrics from the pool itself, plus backend-side connection and CPU/IO metrics, and tune size, idle timeout, and lifetime to minimize wait time while keeping the backend under its ceiling.
- Explicitly test the failure path: kill connections mid-flight, simulate a slow backend, and confirm acquisition timeouts and backpressure behave as designed rather than hanging.
Worked example
A service runs 20 instances against a database capped at 500 total connections. Using the formula above, each instance is capped at 25 pooled connections. During a load test that simulates a rolling deploy (all 20 instances restarting within a short window), every instance attempts to rebuild its pool of 25 connections at once: 20×25=500 simultaneous reconnect attempts against a ceiling of exactly 500, with zero margin for any connection still draining from the old instances. Adding jittered reconnect delays and reducing per-instance pool size to 20 (giving 20×20=400, leaving 100 connections of headroom during a rollover) eliminates the connection-storm failures observed in the unthrottled test.
Trade-offs & pitfalls
- Sizing a pool against a single instance's peak load, without dividing by the fleet size, is the most common and most damaging mistake; it works until autoscaling adds instances, then fails exactly under peak traffic.
- A pool with no acquisition timeout turns backend overload into cascading request pile-ups instead of fast, visible failures; for services making many short-lived connections, pairing the pool with a circuit breaker (a resilience pattern that stops sending requests to a struggling dependency, covered under high-availability patterns rather than here) prevents that pile-up from spreading further upstream.
- Idle timeouts set too aggressively cause needless reconnection churn during normal traffic dips; set too loosely, they let leaked or stale connections accumulate unnoticed.
- Connection leaks (code paths that acquire a connection and never release it, often on an error path) are the quiet failure mode: the pool looks correctly sized until leaked connections slowly starve it, and only a saturation metric with alerting catches this before an outage.
An OLTP application is running at 20,000 requests per second and beginning to suffer from write contention. Propose a database-scaling roadmap: a short-term move (read replicas), a medium-term move (partitioning strategies), and a long-term move (sharding). Describe your rebalancing approach, migration steps, and how you'd maintain consistency and monitor progress along the way.
Sample Answer
Direct answer
Facing write contention at 20,000 requests per second, I would sequence the fix by how much architecture change it costs: read replicas first (days to weeks, no schema change, relieves the primary of read load so it can dedicate more capacity to writes), then partitioning of the hottest write tables (weeks to months, splits one table into many within the same database), then full sharding (months, splits the database itself across independent nodes) only once partitioning alone can no longer keep write throughput ahead of demand. Each step is validated against the projected load, not just the current one, since the roadmap has to hold as the write rate keeps climbing.
Structured elaboration
flowchart LR
A["Today: 20,000 req/s,\nwrite contention on one primary"] --> B["Short-term:\nread replicas"]
B --> C["Medium-term:\ntable partitioning"]
C --> D["Long-term:\nsharding"]
D --> E["Sustains projected\n100,000 writes/s,\n4,000,000 reads/s"]
Short-term: read replicas. Write contention is often partly caused by reads competing with writes for the same primary's CPU, memory, and I/O, even though replicas cannot scale write throughput directly. Deploying asynchronous read replicas and routing read-heavy, non-latency-critical queries away from the primary frees up headroom the primary can spend on writes. This is purely a read-scaling move (a read replica applies an asynchronously-copied stream of the primary's writes and serves reads from that copy, but never accepts writes itself, so adding more replicas relieves the primary of read load without ever adding write capacity); it buys time, it does not solve write contention at its root.
Medium-term: partitioning. Within the same database, split the hottest write tables so each write touches a smaller table with less lock contention and smaller indexes to maintain:
- Range partitioning (by time, for example) suits append-heavy, time-ordered writes and lets old partitions be archived or dropped cheaply.
- Hash or list partitioning suits writes that need to spread evenly rather than cluster by time.
- Implementation: analyze which tables and rows are actually hot, add partitioning without blocking writes (most modern databases support online partitioning operations), backfill in batches, and update queries to include the partition key so the database can prune irrelevant partitions instead of scanning all of them.
- This buys another step of headroom without yet paying the operational cost of a fully distributed database, but it does not remove the ceiling of a single physical database instance.
Long-term: sharding. Once partitioning within one database can no longer keep pace, split the data across independent database instances (shards), each an independent primary handling its own slice of writes in parallel. This is the only step of the three that removes the single-primary write ceiling entirely, and it is the most expensive: it needs a shard-key choice that spreads writes evenly (a hash of a stable, high-cardinality identifier, for example hash(user_id), is the usual default: hashing maps each write to a pseudo-random shard regardless of which user or time period is currently most active, so no single shard absorbs a disproportionate share of the write stream the way it would if writes were assigned by a correlated attribute like signup date), a routing layer, and a migration plan, plus a strategy for the rare cross-shard transaction. Shard-key selection at this stage is also where hotspot risk gets decided: a poorly chosen key concentrates writes on one shard no matter how many shards exist, so the same hotspot-rebalancing discipline used to fix a hot shard after the fact (detect via per-shard load metrics, then move the offending key ranges onto their own shard or split them further) should inform the key choice up front rather than being treated as a problem to solve only once it appears.
Rebalancing, migration, and consistency, applied to this roadmap
- Each transition (single database to partitioned, partitioned to sharded) uses the same safe pattern: backfill the new structure while the old one keeps serving, dual-write once backfill catches up, verify with checksums or row counts, then cut over in stages (a percentage of traffic or a cohort of users at a time) rather than all at once.
- Consistency: keep transactions single-partition or single-shard wherever possible; where a write must touch two shards, prefer an idempotent, retryable pattern (such as a saga: a sequence of local transactions with compensating steps) over a synchronous distributed transaction, since the latter adds coordination latency exactly where you are trying to remove contention.
- Monitoring during migration: replication/backfill lag, row-count or checksum parity between old and new structures, per-partition or per-shard write latency and error rate, and a rollback trigger if divergence exceeds a set tolerance.
Sizing the roadmap against projected load
Capacity planning for this service projects growth to 100,000 writes per second and 4,000,000 reads per second. For concreteness, assume (explicitly, as a planning input rather than a measured fact) that a single post-partitioning shard primary can sustain roughly 10,000 writes per second, and a single read replica can sustain roughly 40,000 reads per second. The number of write shards needed at the projected load:
shards=⌈10,000 writes/s per shard100,000 writes/s⌉=10The number of read replicas needed at the projected load:
replicas=⌈40,000 reads/s per replica4,000,000 reads/s⌉=100Ten shards is a normal, operable sharded topology. One hundred read replicas is not: at that replica count, the coordination and infrastructure cost of pure replication stops being the right lever, and the roadmap needs a caching layer in front of the database, for example an in-memory read cache keyed by the same primary lookup the application already uses, to absorb the bulk of that 4,000,000 reads/second before it ever reaches a replica, rather than trying to solve all of it with replica count alone. As an illustrative planning input, not a measured fact: if that cache achieves roughly a 90% hit ratio on read traffic, it absorbs 0.90×4,000,000=3,600,000 reads/second, leaving 4,000,000−3,600,000=400,000 reads/second that still need to reach a replica. At the same 40,000-reads-per-second-per-replica planning assumption used above, that works out to 400,000/40,000=10 replicas, an operable count in the same range as the 10 write shards, instead of the 100 replicas a cache-free design would require.
Worked example
A checkout-adjacent online transaction processing (OLTP) service sits at 20,000 requests/second today, with write contention visible as rising lock-wait time on the primary. Applying the roadmap: read replicas ship in week 2, immediately reducing primary CPU spent on reads and giving writes more headroom without touching the schema. Table partitioning on the hottest write table ships by month 2, cutting average write latency on that table because each write now touches a smaller partition and a smaller index. By month 8, write volume has grown enough that partitioning alone is no longer sufficient (write latency is climbing again despite partitioning), triggering the sharding project, sized using the ceiling formulas above against the 100,000 writes/second and 4,000,000 reads/second planning target, landing on 10 shards plus a caching layer to keep the read-replica count out of the triple digits.
Trade-offs & pitfalls
- Read replicas are frequently reached for first because they are the cheapest step, but shipping them and declaring the write-contention problem solved is a common wrong turn: they relieve pressure, they do not add write capacity.
- Partitioning without first identifying the actually-hot tables wastes the migration effort on tables that were never the bottleneck; profile before partitioning.
- Sizing a roadmap only against today's load (20,000 requests/second) instead of the projected target undercounts the shard and replica counts needed and forces a second migration soon after the first; size against the planning horizon, as shown above.
- Cross-shard transactions are the main correctness risk introduced by the long-term step; a design that assumes they will stay rare, without an explicit compensating-transaction pattern for when they are not, tends to produce silent inconsistency under load rather than a clean failure.
Describe the cache-aside, read-through, write-through, and write-behind cache topologies, and when you'd reach for each. For every topology, explain the read/write flow and the latency and consistency trade-offs, and give a concrete example use case such as session storage, a product catalog, or a leaderboard.
Sample Answer
Direct answer
Cache-aside, read-through, write-through, and write-behind differ in who is responsible for moving data between the cache and the database, and on which side of a write the durability guarantee sits. Cache-aside and read-through both leave reads on a lazy-load path (the application checks cache first, loading from the database only on a miss), while write-through and write-behind differ from each other in whether a write is confirmed synchronously to the database (write-through) or acknowledged immediately and flushed later (write-behind). The choice comes down to how much staleness and how much write latency the use case can tolerate.
Structured elaboration
Cache-aside (lazy-loading)
- Read flow: application checks the cache; on a miss it reads the database itself, populates the cache, and returns the result.
- Write flow: application writes the database, then explicitly invalidates or updates the affected cache entry.
- Latency and consistency: fast on hits, a first-read penalty on misses; a small race window exists between the database write and the cache invalidation, so a concurrent reader can briefly see stale data.
- Failure handling: if the cache crashes or is flushed, the application transparently falls back to reading the database on every subsequent miss; nothing is lost, throughput just degrades to database-only speed until the cache warms back up. This graceful-degradation property is cache-aside's main durability advantage over the write-behind pattern below.
- Example: a product catalog, where reads dominate heavily and a brief staleness window after a price or description edit is acceptable.
Read-through
- Read flow: the application asks the cache client for a key; the cache itself is responsible for loading from the database on a miss (rather than the application doing it), which centralizes the loading logic instead of duplicating it in every caller.
- Write flow: typically paired with write-through or write-behind, since read-through only defines the read side.
- Latency and consistency: same miss-cost profile as cache-aside; the benefit is code simplicity, not a different consistency guarantee.
- Failure handling: same graceful degradation as cache-aside on a cache failure, since the loader logic still falls through to the database.
- Example: session lookups, where centralizing the load-on-miss logic in the cache layer avoids repeating it across every service that needs a session.
Write-through
- Read flow: same as cache-aside or read-through.
- Write flow: the application writes to the cache, and the cache synchronously writes through to the database before acknowledging the write as successful.
- Latency and consistency: higher write latency than write-behind, because every write pays the database round trip, but the cache and database are never out of sync from the client's perspective.
- Failure handling: if the origin (database) write fails, the cache write must be rolled back or the entry invalidated, or the cache would hold a value the database never actually has; a cache-side crash after a successful database write just loses the cached copy, not any data, since the database is always the durable copy of record.
- Example: a shopping cart, or the broader e-commerce cart use case generally, where the user needs to see their own change reflected immediately and correctly, and the write rate is low enough that the added latency is acceptable.
Write-behind (write-back, asynchronous)
- Read flow: reads are served from cache, populated by writes or by a read-miss loader.
- Write flow: the application writes to the cache, which acknowledges immediately and batches or asynchronously flushes accumulated writes to the database.
- Latency and consistency: the lowest write latency of the four, since the client never waits on the database, but the database is only eventually consistent with the cache, and durability now depends entirely on the cache surviving until its next flush.
- Failure handling, in depth: this is the pattern with a real durability analysis to do, not just a caveat. If the cache crashes before a batch of writes has been flushed, every write in that unflushed batch is lost, a crash-recovery data-loss window whose size is exactly the flush interval (a 5-second flush interval means up to 5 seconds of writes are at risk on any crash). Recovery after a crash means restoring from the last durable flush and accepting that anything after it is gone, unless the cache itself is backed by a write-ahead log or replicated before acknowledging, which narrows the window but adds back some of the latency write-behind was chosen to avoid. Batching writes together to flush also introduces write amplification when a key is updated many times before its flush: instead of every individual write reaching the database, only the final value per flush cycle is written, which reduces database load but means intermediate values are never durably recorded at all, not just delayed.
- Example: high-throughput counters, such as a leaderboard, or an analytics-ingestion pipeline, where the flush interval is analogous to a batching window and losing a small, bounded amount of the most recent data on a rare crash is an acceptable trade for sustaining a write rate the database could not absorb directly. The durability trade-off is the same shape in both cases: low per-write latency and reduced database load, in exchange for a bounded, quantifiable data-loss window plus write amplification on frequently-updated keys.
Choosing between them
| Priority | Best fit |
|---|---|
| Simple, read-heavy, tolerant of brief staleness | Cache-aside |
| Same as above, but want loading logic centralized | Read-through |
| Need the cache and database to never disagree, can accept write latency | Write-through |
| Need the lowest possible write latency, can accept a bounded, understood data-loss window | Write-behind |
Worked example
An e-commerce cart service needs the cart total to be correct the instant a user adds an item (write-through: synchronous database write, cache mirrors it, no staleness). The same platform's "trending products" counter, incremented on every product view across millions of views per day, uses write-behind: each view increments the cached counter and returns instantly, and the counter is flushed to the database on a fixed interval. If the cache process crashes, the trending counter loses at most that interval's worth of increments, which is invisible on a counter fed by millions of events, whereas losing even one cart write would be a customer-facing correctness bug. The same reasoning applies directly to an analytics-ingestion pipeline: individual event counts can tolerate the same bounded loss window, but a financial ledger update could not, which is why write-behind is scoped to the specific fields where that trade-off is safe rather than applied to a whole write path uniformly.
Trade-offs & pitfalls
- Treating write-behind's flush interval as a tuning knob without stating the resulting data-loss window explicitly is a common gap in design reviews; the window should be a named, deliberate number, not an implicit side effect of the batch size chosen for throughput.
- Cache-aside's invalidate-after-write race window is small but real under high concurrency; if a use case cannot tolerate any stale read, write-through is the safer default despite its latency cost.
- Read-through's benefit is purely architectural (centralized loading logic); choosing it under the assumption that it also improves consistency versus cache-aside is a misunderstanding of what the pattern actually changes.
- For any pattern, a cache failure that is treated as a full outage rather than a graceful fallback to the database is a design gap, not an inherent property of caching: cache-aside, read-through, and write-through can all fall back cleanly; only write-behind has real, quantifiable data at risk on a cache crash.
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.