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.
Your relational database allows a maximum of 500 concurrent connections. Your application runs on 25 identical JVM instances plus 3 background job workers. Walk through how you'd compute a safe default connection-pool size per instance, the factors you'd weigh, and your recommended pool size. What would you monitor, and what would you do if you saw connection saturation in production?
Sample Answer
Direct answer
Work backward from the database's hard ceiling, not forward from a per-instance number that feels generous. Reserve headroom for admin connections, failover, and bursts, divide what remains across every process that opens a connection (application instances and background workers alike), and floor the result. For 25 application instances (each running as its own process, commonly a JVM, a Java Virtual Machine, instance) plus 3 background workers against a 500-connection ceiling, that works out to a safe default of 14 connections per process, with room to tune down from there once you have real utilization data.
Structured elaboration
The formula
- Total connection-opening processes: 25 application instances + 3 background workers = 28.
- Reserve headroom for database maintenance tooling, admin connections, and burst capacity: a common starting assumption, stated here as a planning input rather than a measured fact, is 20%, leaving 80% of the ceiling as usable pool capacity.
- Safe per-process pool size: divide usable capacity across all processes and round down, so no single process's pool can push the fleet over the ceiling if every process is fully utilized simultaneously.
usable=500×0.80=400
per-process pool=⌊28400⌋=⌊14.29⌋=14
Factors that push the number up or down
- Transaction duration: long-running transactions hold a connection for longer, so a workload with slow queries needs a smaller pool per process (or faster queries) to avoid connections queuing up behind a few slow holders.
- Background workers are often burstier than request-serving instances; giving them the same per-process budget as application instances can be wasteful most of the time and insufficient during a burst. A separate, smaller steady-state budget with the ability to borrow from an external pooler (below) handles this better than a single uniform number.
- Read-heavy workloads can shift some of this load to read replicas, reducing the connection pressure on the primary specifically.
- Connection acquisition overhead and idle timeout settings affect how much of the pool is actually "available" at any moment versus tied up in idle-but-not-yet-reaped connections.
Monitoring and response to saturation
- Track: active connections per process, pool wait count and wait time, database-side connection count versus the 500 ceiling, and application-level timeout errors on connection acquisition.
- Alert when pool utilization sustains above roughly 75%, or when any request has to wait for a connection for more than a short, explicitly agreed threshold.
- On saturation: first apply backpressure (queue or reject new work rather than let requests pile up waiting on connections), then check whether the pool size itself is oversized relative to actual concurrent demand (overcommitted pools waste database-side capacity even when the app-side pool looks "full" only occasionally), then consider moving read traffic to a replica, and only then treat raising the database's connection ceiling or adding an external pooler as the longer-term fix.
Worked example
The formula above holds at moderate fleet size, but it breaks down as the fleet grows faster than the database's connection ceiling. Assume, again as a planning input, a fleet that has grown to 200 identical instances against a database with a 1,000-connection ceiling. Applying the same approach:
usable=1000×0.80=800
per-instance pool=⌊200800⌋=4
Four connections per instance is thin: if a single instance briefly needs more than four concurrent database calls, for example during a request burst, it will queue locally even though the database as a whole has spare capacity sitting idle in other instances' unused pool slots. This is also exactly the scenario where a rolling deploy becomes dangerous: if all 200 instances restart in a short window and each immediately tries to open its full pool of 4 connections, that is a connection storm, a spike toward the 1,000-connection ceiling driven by simultaneous reconnects rather than steady-state load, and it can exhaust the ceiling before the deploy even finishes.
The fix at this scale is usually an external pooler such as PgBouncer, placed between the 200 application instances and the database, multiplexing many app-side connections onto a much smaller database-side pool. PgBouncer's two relevant modes make different trade-offs:
- Session pooling mode: a database connection is assigned to one client connection for its entire session (until it disconnects). This preserves session-scoped features like prepared statement caching and
SETcommands, but it does not reduce database-side connections below the number of concurrently active client sessions, so it does not solve the 200-instances-against-1,000-connections problem by itself. - Transaction pooling mode: a database connection is held only for the duration of a single transaction, then released back to the pool immediately for another client to use. This lets a small database-side pool (for example, 50 to 100 connections) serve far more app-side connections than session mode can, at the cost of breaking anything that relies on session state persisting across transactions, such as advisory locks held between transactions or
SET-scoped session variables.
Given the transaction-pooling trade-off, a deploy in this configuration should also stagger restarts (for example, capping simultaneous restarts to a small percentage of the fleet with a ramp-up delay between batches) rather than restart all 200 instances at once, regardless of whether an external pooler is in place.
Trade-offs & pitfalls
- Sizing the pool from "what feels safe" rather than the database's actual ceiling is the most common mistake: it works until a deploy, a traffic spike, or a new service instance pushes the fleet over the limit all at once.
- A pool that is too large per instance doesn't help throughput once the database itself is saturated; it just moves the queue from the application to the database, where diagnosing it is harder.
- An external pooler adds an operational dependency and a new failure mode (the pooler itself can become a bottleneck or single point of failure) in exchange for solving a problem the simple per-instance formula cannot solve at high instance counts.
- Any pool-sizing decision should be re-validated with load testing and real monitoring data, not treated as a one-time calculation; traffic mix and query patterns change the safe number over time.
Design sharding for a real-time pub/sub service that has a very large number of channels. Compare client-side sharding, where clients pick their partition, against server-side partitioning, where the broker assigns it. Discuss rebalancing cost, fairness, latency, and client churn from mobile users with intermittent connectivity, and recommend an approach for a client base that is mostly mobile and frequently offline.
Sample Answer
Direct answer
For a pub/sub service with a huge number of channels and a client base that is mostly mobile and frequently offline, prefer server-side partitioning over pure client-side sharding. Client-side sharding (clients hash the channel to pick their partition themselves) is operationally cheap but turns every flaky-network reconnect into a fresh routing decision made independently by thousands of clients. Server-side partitioning (the broker assigns and owns the mapping) keeps that complexity centralized, where it can buffer subscription state through a drop and rebalance in a controlled, batched way instead of leaving it to uncoordinated client behavior.
Structured elaboration
The two placement strategies
- Client-side sharding: the client computes
partition = hash(channel_id) % n(or a similar scheme) and connects directly to the broker owning that partition. No broker-side coordination service is needed. - Server-side partitioning: clients connect to a stateless gateway or routing tier; a partition-metadata service tells the gateway which broker shard owns the channel, and the broker fleet can move ownership without the client ever recomputing anything.
Compare on the dimensions that matter here
- Rebalancing cost: under client-side sharding, changing shard count means every client must recompute its mapping, and if the mapping uses naive modulo hashing, nearly all clients reconnect to a different shard at once. Server-side partitioning can migrate ownership of a subset of channels while the gateway keeps existing client connections open, so the client never sees the rebalance directly.
- Fairness: client-side hashing has no visibility into actual load, so it can leave one shard hot while another idles. A server-side metadata service can rebalance based on measured load per shard, not just key distribution.
- Latency: client-side sharding has one fewer hop if the client connects straight to the right broker. Server-side partitioning adds a routing lookup, but that lookup can be cached at the gateway and is usually a small fraction of a mobile network's own round-trip time.
- Client churn from intermittent connectivity: this is the deciding factor for a mostly-mobile, frequently-offline client base. Client-side sharding treats every reconnect as a brand-new decision, so a client that drops and reconnects seconds later re-derives its partition and may miss messages published in between. Server-side partitioning lets the broker hold the subscription state (with a bounded grace period) so a quick reconnect resumes the same logical subscription instead of starting over.
A brief note on mechanics: whichever side owns the mapping, production systems generally implement the underlying hash-to-shard assignment with a consistent-hashing-style ring rather than plain modulo, specifically to bound how many keys move when shard count changes: nodes and channels are placed as points on a circular hash space (the ring), and virtual nodes give each physical shard many small points on it instead of one, so only the moved points' channels need to migrate. The ring construction and virtual-node tuning that make that work: giving each physical shard many small points scattered around the ring instead of one large contiguous range means more virtual points per physical shard smooths load distribution further but adds more routing metadata to maintain, while fewer points is cheaper to track but risks uneven load; the point relevant here is just that the choice of hashing scheme, not just who runs it, determines how disruptive a rebalance is.
Worked example
Assume, as a planning input rather than a measured fact, a fleet of 50 broker shards holding about 2 million channels between them, and that the service needs to grow to 60 shards to handle load. Under a consistent-hashing-style ring, the expected fraction of channels that move when going from n to n + m shards is approximately:
n+mm=6010≈16.7%
Under naive modulo hashing (partition = hash(channel_id) mod n), changing the divisor n reassigns nearly every key, so close to 100% of channels remap even though only 10 shards were added. That gap, roughly 17% of clients reconnecting elsewhere versus nearly all of them, is why a server-side partitioning layer that controls the hashing scheme matters more than which side technically "does" the sharding: a mobile client base that reconnects constantly on its own cannot absorb a near-100% remap event without a visible spike in dropped or duplicated messages.
Trade-offs & pitfalls
- A server-side metadata/routing tier becomes a new critical-path dependency. If it is a single instance rather than a redundant, health-checked fleet, it becomes the system's weakest point exactly when you added it to improve reliability.
- A mass reconnect event (an app release, a regional network outage recovering) can still overwhelm a well-designed server-side system if all clients retry at once; client SDKs need jittered backoff regardless of which sharding approach is chosen.
- Client-side sharding is not always wrong: for a small number of large, stable shards and a client base with reliable connectivity (internal services, desktop clients), its simplicity and lack of a coordination service can be the better trade.
- Buffering subscription state server-side to smooth over churn has a cost: unbounded buffering for offline clients grows broker memory. A grace period with a hard cap forces an explicit decision about how "offline" a client can be before it must fully resubscribe.
flowchart LR
subgraph Server-side partitioning
C1["Mobile client"] --> GW["Stateless gateway"]
GW --> MD["Partition metadata service"]
MD --> GW
GW --> B1["Broker shard 1"]
GW --> B2["Broker shard 2"]
end
Design an operational plan and technical implementation to reshard a live sharded database with minimal downtime. Cover choosing the new shard key or shard count, the data-migration strategy (online migration, dual writes, change-data-capture), routing updates, throttling the migration, validation steps, and your rollback procedure.
Sample Answer
Direct answer
Resharding a live database with minimal downtime is an online migration: keep the old topology serving traffic while you copy data to the new one under a throttle, use dual writes or change-data-capture (CDC) to keep the copy current, validate it against the source before trusting it, and cut routing over in a small, reversible step rather than a single big-bang switch. The hard parts are rarely the copy itself; they are choosing a shard key/count that won't need redoing soon, keeping writes ordered correctly during the overlap window, and having a rollback that is actually exercised, not just documented.
Structured elaboration
Choosing the new shard key or shard count
- Analyze real access patterns (write hotspots, cardinality, range-vs-hash trade-offs) from production telemetry rather than guessing; target per-shard headroom (commonly keeping any shard well under its saturation point) so the new topology doesn't need revisiting immediately.
- For a hot single shard suffering write contention, the two concrete techniques are rekeying (choosing a shard key that spreads the contended writes, if the current key is the root cause) and range-splitting (dividing the shard's existing key range into two or more ranges, each becoming its own shard, when the key itself is fine but the range has outgrown one node).
- For a lookup-based versus algorithmic shard-mapping layer: an algorithmic mapping (consistent hashing, which places nodes and keys as points on a circular hash space called a ring, so a topology change only reassigns the small slice of keys near it, or the simpler hash-mod-N) needs no metadata lookup and scales trivially, but it distributes tenants roughly evenly by key, which fails badly when tenant sizes vary by orders of magnitude (a large enterprise tenant next to hundreds of tiny ones). A lookup-based mapping (an explicit tenant-to-shard table) costs a metadata hop per request but lets you place large tenants on dedicated shards and pack small tenants together deliberately. For a multi-tenant system with that kind of size skew, prefer lookup-based mapping, or a hybrid: algorithmic mapping for the default pool of ordinary tenants, with a lookup override for the small number of outsized ones.
Data-migration strategy
- Online copy under CDC: stream changes from the source shard(s) via a CDC pipeline (e.g., a log-based change stream) into the target topology while a background job bulk-copies existing data in ranges. The CDC stream and the bulk copy have to be reconciled (apply CDC events that land inside an already-copied range, buffer or replay ones that land ahead of the copy cursor).
- Dual writes with a quiesce window: to guarantee transaction ordering is preserved for records that are actively being migrated, briefly pause (quiesce) writes to the specific key range being cut over, drain the in-flight CDC backlog for that range so source and target agree, flip routing for just that range, then resume writes. This bounds the ordering risk to a small, deliberately-chosen window instead of trying to reason about ordering across an open-ended dual-write period for the whole dataset.
- Split-and-merge at scale: for a large, contended shard, resharding is a split (one shard's range divides into two, each getting its own routing-table entry) rather than a full-dataset migration; the reverse operation, merge, combines two under-loaded shards back into one when growth was over-provisioned. At terabyte scale, both operations are driven by range boundaries and routing-table updates rather than moving the whole dataset, which keeps a single split or merge bounded in time and blast radius.
- A cross-topology special case: migrating a stateful store (for example, a session store) from a single-region deployment to a multi-region sharded cluster follows the same dual-write skeleton, but adds a proxy layer in front of both topologies so the application never talks to the store directly during the migration; the proxy performs parity verification (comparing responses from the old and new paths for the same key) before fully cutting traffic over, which catches migration bugs before they're user-visible.
flowchart TD
A[Old shard topology] -->|CDC stream| B[Change pipeline]
A -->|throttled bulk copy| C[New shard topology]
B --> C
C --> D[Shadow reads: compare old vs new]
D -->|parity confirmed| E[Routing cutover, range by range]
D -->|mismatch| F[Pause cutover, reconcile]
E --> G[Quiesce window per range: drain, flip, resume]
G --> H[Old topology retired]
Routing updates
Introduce a migration-aware routing layer (gateway or client-side routing table) so application code never encodes shard topology directly. Roll the new routing out in stages: internal/test traffic, a small canary of production traffic, then full cutover, with each stage able to fall back to the previous routing table without a code deploy.
Throttling
Throttle the copy and CDC-apply workers adaptively against replication lag, tail latency, and source-shard CPU/I/O, not a fixed rate, so migration traffic backs off automatically when it's competing with production load. As a worked illustration of sizing a per-shard step against an operational target: suppose the team's service-level agreement (SLA) is for each per-shard rebalance step to complete in well under two minutes even across a 200+ node cluster, and a given shard holds roughly 50 GB of data with the copier capped at 500 MB/s per shard to protect production I/O. The data-transfer portion alone takes
which leaves about 20 seconds of headroom under the 120-second target for updating ring/routing metadata and completing leader election (the process by which the new shard's replica-set nodes agree on which one of them becomes the active primary/writer) for the new shard's replica set, before the step is considered done. That headroom is the number to watch: if metadata propagation or leader election routinely eats into it, the fix is a faster metadata path, not a bigger copy-bandwidth cap.
Automating rebalancing across many shards
When rebalancing hundreds of shards simultaneously rather than one at a time, the automation has to cap total concurrent migrations (global throttle, not just per-shard), and it must keep secondary-index consistency in scope explicitly: an index that's rebuilt or re-pointed on a different schedule than its base table will silently return wrong results for the window in between, so index cutover needs to be part of the same per-range quiesce step as the base data, not a follow-up job.
Validation
- Row counts and per-range checksums between source and target.
- Shadow reads: serve a fraction of real production reads from the new topology in parallel with the old one (without returning the shadow result to the client) and compare, which surfaces correctness bugs under real traffic patterns before any client depends on the new path.
- Client-side cache invalidation during cutover: if clients or edge caches hold entries keyed by the old routing, cutting a range over without invalidating those cached entries serves stale-routed responses; the cutover step has to include a cache-invalidation signal (a version bump on the routing key, or explicit invalidation) for anything caching downstream of the router.
- End-to-end application-level smoke tests against the new topology before removing fallback to the old one.
Rollback
Keep dual writes (or the CDC stream) active until validation passes; if a problem surfaces, stop applying new CDC events, revert routing to the old topology for the affected range, and resume from the last confirmed-consistent checkpoint rather than restarting the whole migration. Because cutover happens range-by-range behind a quiesce window, rollback scope is bounded to the ranges already flipped, not the whole dataset.
Trade-offs & pitfalls
- A migration that skips the quiesce window to avoid any write pause will have an ordering gap for records touched during the exact moment of cutover; a short, well-instrumented pause on a narrow key range is a better trade than an unbounded consistency risk.
- Algorithmic shard mapping is simpler to build and reason about, but for multi-tenant systems with orders-of-magnitude tenant-size variance, defaulting to it produces shards that are balanced by tenant count and wildly unbalanced by load; verify the mapping choice against actual per-tenant load, not tenant count.
- Skipping shadow reads to save migration time is the single most common way correctness bugs escape into production, because checksums catch data-copy errors but not query-behavior differences (an index that returns results in a different order, a range boundary edge case).
- Treating secondary-index rebuild as a background job that can lag the base-table cutover is a frequent source of silent incorrect reads immediately after a migration; keep index and base-data cutover atomic per range.
Explain how read replicas for relational databases improve read throughput. Describe the common replication modes (asynchronous versus semi-synchronous) and the operational pitfall of replication lag. What monitoring and safeguards would you put in place to detect and handle a lagging replica?
Sample Answer
Direct answer
Read replicas are read-only copies of a primary relational database that let you route read-heavy traffic away from the primary, so read throughput scales roughly with the number of replicas instead of being capped by one machine's capacity. The two common replication modes trade off write latency against durability: asynchronous replication is fast but can lag, semi-synchronous replication waits for at least one replica to acknowledge before confirming a write, trading some write latency for a stronger durability guarantee. Replication lag, the gap between a write landing on the primary and appearing on a replica, is the operational pitfall that follows directly from choosing asynchronous replication for speed.
Structured elaboration
Why read replicas scale reads
A single primary database has a ceiling on how many queries per second (QPS, the standard measure of database or API load) it can serve before CPU, memory, or I/O saturates. Since most application workloads are read-heavy relative to writes, adding replicas that each hold a full copy of the data lets read queries fan out across many machines while writes still funnel through the one primary that owns correctness. This is a read-scaling pattern specifically: it does nothing for write throughput, which is bounded by the primary alone (write scaling is a separate problem, addressed by partitioning or sharding rather than replicas).
Replication modes
- Asynchronous: the primary commits and returns success to the client without waiting for any replica to apply the change. Write latency stays low and unaffected by replica health, but a replica can fall arbitrarily behind under load, and if the primary fails before a replica caught up, those last writes are lost from that replica's perspective.
- Semi-synchronous: the primary waits for acknowledgment from at least one replica (that the write was received, not necessarily fully applied) before confirming the commit to the client. This bounds the worst-case data loss to writes that hadn't yet reached any replica, at the cost of added write latency and a risk that a slow replica introduces a stall on every write.
This is standard terminology in an online transaction processing (OLTP) context, meaning a workload of many small, individual reads and writes (as opposed to large analytical scans); read replicas are one of the first tools reached for once a single OLTP primary starts to strain under read load.
Replication lag as the operational pitfall
Lag arises from network delay, I/O contention on the replica, or the replica processing a backlog of changes slower than the primary produces them. Its consequence is stale reads: a client that just wrote data may query a replica and not see its own write, or two clients may observe the data in different states depending on which replica they hit.
Read-routing design to minimize stale reads while maximizing throughput
The application layer, not just the database, needs a policy for which reads are allowed to be stale:
- Reads that must reflect the client's own very recent write (a user viewing the profile they just edited) should go to the primary, or to a replica only after confirming its lag has caught past that write's position.
- Reads that tolerate a small staleness window (a public dashboard, a search index, an analytics report) should go to replicas by default, since that is where the throughput gain comes from.
- A hybrid policy, sometimes called read-your-writes routing, pins an individual client to the primary (or to a replica known to be caught up) for a short window right after that client writes, then lets subsequent reads fall back to any replica.
Monitoring and safeguards
| What to watch | Why |
|---|---|
| Replication lag (seconds and/or log position gap) | Direct measure of staleness risk; the number a routing or alerting decision should key off |
| Replica apply rate versus primary write rate | Rising divergence predicts lag will keep growing rather than catch up |
| Replica CPU/IOPS (input/output operations per second)/network | Identifies whether the replica itself is the bottleneck causing lag |
| Query load on replicas (especially long-running analytical queries) | A single expensive query can starve the replication-apply thread and cause a lag spike |
Safeguards built on that monitoring: alert when lag crosses a threshold tied to the application's staleness tolerance; throttle or move expensive ad hoc/analytical queries off replicas that also serve latency-sensitive reads; and, for any workflow that promotes a replica (to primary, during a failure), require lag to be at or near zero before promotion, since promoting a lagging replica means accepting the unreplicated writes as lost. That promotion and failover mechanics belong to the high-availability side of the system, not to the read-scaling pattern itself, but the monitoring described here is exactly what feeds that decision when it happens.
Worked example
A social-media-style application serves 9,000 reads per second and 1,000 writes per second against a single primary that is now CPU-saturated on reads. Adding 3 asynchronous read replicas and routing all reads except "read-your-own-write" cases to a round-robin pool across them reduces the read load on the primary from 9,000 QPS to roughly 0 (reads move off entirely), leaving the primary handling only the 1,000 writes/second plus the small share of reads that require read-your-writes freshness. Each replica now carries roughly 9,000 / 3 = 3,000 reads/second on average, well within a single replica's typical headroom, illustrating the linear-ish scaling read replicas provide as long as write volume itself stays within what one primary can sustain.
Trade-offs & pitfalls
- Read replicas scale reads only; teams sometimes reach for them to fix a write-contention problem, which they cannot, because writes still funnel through one primary.
- Asynchronous replication's low write latency is attractive, but skipping the read-routing design above (treating every replica as equally fresh) is the most common way stale reads leak into user-facing behavior.
- Semi-synchronous replication reduces data-loss risk but can introduce write stalls if the acknowledging replica itself becomes slow; it shifts risk from data loss to latency, it does not eliminate risk.
- Promoting a lagging replica during an incident, without checking lag first, can silently drop the most recent committed writes; this is a data-loss event dressed up as a recovery action.
Explain how a CDN works and when you'd reach for one in a global application. Cover edge caching, cache-control headers, TTL strategy, surrogate keys, origin failover, cache invalidation, and how you'd handle dynamic versus static content (signed URLs, edge logic). What are the cost and operational trade-offs?
Sample Answer
Direct answer
A content delivery network (CDN) is a fleet of geographically distributed edge servers that cache and serve content close to the requesting user, cutting latency and offloading the origin server. You reach for one in a global application whenever a meaningful share of traffic is cacheable and users are spread across regions, since the CDN removes both the network-distance cost and the repeated origin load for the same content served over and over.
Structured elaboration
Edge caching and cache-control
Edge nodes store a response and decide how long to keep it based on HTTP caching headers the origin sets, most commonly Cache-Control (for example public, max-age=3600 to allow shared caching for one hour) and Vary (to tell the CDN that responses differ by a request header, such as Vary: Accept-Encoding, so it must cache separate copies per variant). ETag lets the edge or the client do a conditional request and get a cheap "not modified" response instead of re-downloading unchanged content.
TTL strategy
Static, versioned assets (images, bundled JS/CSS with a hash in the filename) can take long time-to-live (TTL) values, hours to days, since a content change simply ships under a new URL. Semi-dynamic content (an API response that changes occasionally, a personalized-but-cacheable fragment) needs a short TTL, seconds to minutes, and benefits from stale-while-revalidate (serve the slightly-stale cached copy immediately while refreshing it in the background) to keep latency low without serving badly outdated data. As an illustrative split, not a universal rule: a content-heavy site might find roughly 70% of its requests are static and 30% are dynamic or personalized, which is a useful mental model for deciding where TTL and invalidation effort should concentrate, since the static 70% is where aggressive caching pays off with the least risk.
Surrogate keys and invalidation
A Surrogate-Key (or Surrogate-Control) header lets the origin tag a response with one or more logical keys (a product ID, a category, a build version) so that a single purge call can invalidate every cached object sharing that tag, without the origin needing to know every individual URL that resulted from it. This is the mechanism that makes targeted, safe invalidation of a large object set practical: purging by surrogate key when a product's price changes clears exactly the cached responses for that product, not the whole cache. Purge and invalidate are distinct operations: a purge removes the object outright (next request is a full miss), while marking an object stale lets the edge serve it once more while fetching a fresh copy in the background, which is gentler on the origin during a large invalidation event.
Origin failover
Configuring a primary and secondary origin pool, with edge-side health checks, lets the CDN route around a failing origin automatically. Origin shielding, routing all edge cache misses through one designated shield location before they reach the origin, protects the origin from a "thundering herd" of simultaneous misses across many edge nodes after a mass invalidation or cold start.
Dynamic versus static content
Static assets are served straightforwardly from edge cache with long TTLs. Dynamic or personalized content needs a different approach:
- User-uploaded static assets (profile photos, attachments) are cacheable like any static asset once uploaded, but the origin (typically an object store, not the application server) needs its access secured, commonly signed URLs or a signed cookie, so the CDN can serve the object publicly at the edge while the origin itself stays access-controlled rather than open to the world.
- Signed URLs and edge logic: for content that must be authorized per-user but is still worth caching, a signed URL with an expiry lets the edge validate the request without a round trip to the application, and edge compute (code that runs at the edge rather than only at the origin) can handle that validation, run A/B test bucketing, or assemble a personalized response from cached fragments, all close to the user.
Effectiveness metrics
The metrics that tell you whether a CDN deployment is actually working: cache hit ratio (the share of requests served from edge without reaching origin), the TTL distribution actually observed across your cached object types (confirming static assets are getting long TTLs and dynamic ones short, as designed), and origin request rate (the traffic the origin actually sees, which is what a CDN exists to reduce). These are the metrics to instrument and watch, not numbers to assume; a CDN misconfigured with no-store on cacheable responses can show a healthy-looking deployment with a near-zero hit ratio.
Cost and operational trade-offs
- Caching reduces origin compute and egress cost but adds CDN service fees and the operational surface of managing cache rules, purges, and edge logic.
- Edge compute adds real capability (auth, personalization, A/B logic close to the user) at the cost of a new runtime to test, deploy, and debug, plus more vendor-specific surface.
- Aggressive TTLs cut origin load and latency the most but raise the risk of serving stale content, especially for personalized or frequently-changing data; conservative TTLs are safer but leave more traffic hitting the origin. Tuning this trade-off is what the metrics above are for.
Illustrative cache-hit and origin-load estimate
To reason about the trade-off concretely rather than by assumption, here is a worked, fully-pinned example (all inputs stated, not measured or claimed as a real published result):
Assume a video-thumbnail service receives 1,000,000 thumbnail requests per day, and, before any CDN, all of them hit the origin. Suppose an aggressive edge-TTL policy is expected to reach a 90% cache hit ratio, while a more conservative TTL policy (chosen for a use case sensitive to staleness) is expected to reach only 60%. The origin request rate under each policy:
origin requests=total requests×(1−hit ratio) aggressive: 1,000,000×(1−0.90)=100,000 origin requests/day conservative: 1,000,000×(1−0.60)=400,000 origin requests/dayThe aggressive policy cuts origin load 4x further than the conservative one (400,000 / 100,000), which is the shape of the trade-off: every point of hit-ratio improvement is a direct, proportional reduction in origin traffic and cost, and the right policy for a given asset type depends on how much staleness that asset can tolerate, weighed against that origin-load reduction.
Worked example
A global content platform serves images through a CDN with Cache-Control: public, max-age=604800 (roughly 7 days) and a Surrogate-Key per content ID. When an image is replaced, the origin issues a single purge by surrogate key rather than needing to know every resized/format variant URL the CDN generated from it, all of which share that key. Meanwhile the platform's semi-dynamic JSON API for "current viewer count" uses Cache-Control: public, max-age=5, stale-while-revalidate=30 so the edge serves a slightly-stale count instantly while refreshing in the background, keeping origin load low without users perceiving staleness beyond a few seconds.
Trade-offs & pitfalls
- Caching personalized or dynamic content without a correctness plan (the wrong
Varyheader, or none at all) is the most common CDN mistake: it silently serves one user's personalized response to another. - A CDN can worsen correctness for content that changes faster than its TTL if invalidation is not wired up; caching is not a substitute for proper invalidation, it is a mechanism that makes fast, correct invalidation more necessary, not less.
- Origin shielding helps against mass cache-miss storms but adds a hop and a new potential bottleneck of its own if the shield location is undersized relative to peak miss traffic.
- Treating hit ratio as the only success metric misses half the picture: a high hit ratio on the wrong content (over-caching data that needed to be fresh) is a correctness bug wearing a good-looking dashboard.
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.