Database Performance Tuning and Scaling Questions
System-level performance work beyond a single query: configuration and resource tuning, capacity planning, handling large data volumes, and scaling read and write throughput. Covers identifying bottlenecks, growth management, the vertical-versus-horizontal scaling decision, materialized views for expensive queries, and maintenance work such as vacuuming, index rebuilds, bulk loads, and safe schema changes on large tables. Tests whether a candidate can keep a database healthy as load grows.
You are running PostgreSQL with a 500GB dataset on a host with 200GB of total memory. How would you determine reasonable values for shared_buffers and work_mem? What would you monitor to know if you got it right (cache hit ratio, page read times, OS page cache usage), how would you avoid swapping, and how would you roll the change out safely in production?
Sample Answer
Direct answer
Since 500 GB doesn't fit in 200 GB of RAM, size shared_buffers to hold the actively hot
slice of the data, not the whole dataset, a widely used starting point is about 25% of
RAM (roughly 50 GiB here), leaving the rest to the operating system's own page cache,
which Postgres relies on as a second layer. Size work_mem conservatively, because it's
granted per sort or hash operation, not as one shared pool, and a single complex query
can request several of these at once, multiplied again by parallel workers. This matters most under
a typical OLTP workload (many short, concurrent transactions serving live application traffic, as
opposed to a few large analytical queries), since a high global default gets multiplied across every
concurrently active connection at once.
Reasoning through the two settings
shared_buffers: the 25% guideline is a documented starting point, not a law. Past
that point, growing it further competes with the OS page cache for the same physical
RAM (a page cached in both layers is partially wasted memory), while increasing
background-writer and checkpoint (the periodic point where Postgres flushes recently modified
in-memory pages to disk) overhead proportional to how much dirty data can
accumulate. It requires a full server restart to change, since it's allocated as shared
memory at startup, which shapes the rollout plan below.
work_mem: budgeted per sort, hash join, or merge operation. A single query with two
sorts and a hash join can request three allotments simultaneously; parallel workers each
get their own. The real ceiling isn't total_RAM / connection_count, it's headroom
divided by realistic peak concurrent operations. Setting it too high is a classic
out-of-memory or swap risk under concurrent load, not a free win for one query's sort.
Unlike shared_buffers, it can be changed per session or per role without a restart.
Rollout plan: because work_mem is a session/role-level setting, trial a higher value
on just the reporting role first, watch memory behavior under real concurrency, then
promote it to the default once validated. Because shared_buffers needs a restart, plan
a maintenance window, keep the previous value on hand, and know reverting also needs a
restart, it isn't a hot config reload either direction.
What to monitor: the buffer cache hit ratio (below), whether sorts are spilling to
disk (the query plan will say so explicitly, see below), OS-level free memory and swap
activity (a host that starts swapping is making everything slower, not just one
under-provisioned query), and page read/write timing if track_io_timing is enabled.
Avoiding swapping: never let shared_buffers + (work_mem × realistic peak concurrent operations) + other server overhead approach physical RAM. A work_mem value that's
safe at 5 concurrent users can be the thing that swaps the box at 50; set a deliberate
ceiling per role rather than raising it until one report "feels fast."
Worked example (executed on a real PostgreSQL 16 instance)
Defaults on a fresh container (small demo values, shown only to illustrate the commands,
not as a sizing recommendation for a 200 GB host):
SHOW shared_buffers; -- 128MB
SHOW work_mem; -- 4MB
SHOW effective_cache_size; -- 4GB
The concrete signal a query plan gives you when work_mem is too low, run against a
500,000-row, 87 MB table, ordering by a text column:
SET work_mem = '64kB';
EXPLAIN (ANALYZE, BUFFERS) SELECT * FROM bench ORDER BY payload;
Sort Method: external merge Disk: 26008kB
Worker 0: Sort Method: external merge Disk: 26448kB
Worker 1: Sort Method: external merge Disk: 27024kB
SET work_mem = '64MB';
EXPLAIN (ANALYZE, BUFFERS) SELECT * FROM bench ORDER BY payload;
Sort Method: quicksort Memory: 36981kB
Worker 0: Sort Method: quicksort Memory: 32689kB
Worker 1: Sort Method: quicksort Memory: 34701kB
"External merge" with a Disk: size means the sort didn't fit in its work_mem
allotment and spilled to disk, extra I/O the query wouldn't otherwise pay. "Quicksort"
with a Memory: size means it stayed fully in memory. That line in a plan is a direct,
observable signal that a hot query's work_mem is undersized; it's what you'd grep for
in a slow-query log, not just a number you infer from theory.
Cache hit ratio, executed against the same instance:
SELECT sum(heap_blks_hit) AS hit, sum(heap_blks_read) AS miss,
round(100.0*sum(heap_blks_hit)/nullif(sum(heap_blks_hit)+sum(heap_blks_read),0),2) AS pct
FROM pg_statio_user_tables;
-- hit=555554 miss=0 pct=100.00
The demo database fits entirely in this small container's memory, so 100% is expected
here and isn't the target for the 500 GB scenario; on the actual host you'd track this
ratio against a target you've validated for your workload, since 500 GB genuinely
cannot be fully cached in 200 GB of RAM, that gap is exactly why this sizing decision is a
real trade-off rather than a solved problem.
Trade-offs and pitfalls
- Treating "25% of RAM" as a formula rather than a starting point: a mostly-sequential-
scan analytical workload wants a smallershared_buffersand a larger
effective_cache_size(a planner hint about assumed total cache, it allocates no memory
itself) instead. - Setting
work_memglobally high enough to satisfy the worst report risks an OOM under
real peak concurrency; a lower global default plus a per-role override for the specific
workload that needs bigger sorts is the safer shape. effective_cache_sizebeing misconfigured (too low or too high) silently misleads the
query planner's cost estimates even though it allocates no memory itself, worth being
explicit about since it's a common point of confusion withshared_buffers.shared_bufferschanges need a restart to apply and to revert; roll it out in a
maintenance window with the previous value ready, not as a live config reload.
You need a materialized view or aggregated table that stays fresh within 10 minutes over a dataset ingesting about 1TB per day. How would you implement incremental refresh so you are not recomputing the whole view every time, what would partition-level refresh buy you, and what do you do when the underlying data gets corrected after the fact through a backfill?
Sample Answer
Approach
Native Postgres materialized views do not support incremental refresh, REFRESH MATERIALIZED VIEW, with or without CONCURRENTLY, always recomputes the entire result. To hit a 10-minute freshness target over a dataset ingesting about 1 TB per day, build incremental refresh directly: a trigger-maintained summary table that updates only the rows affected by each write, so refresh cost scales with the volume of change since the last update, not with total dataset size.
The source table should itself be partitioned by time at this ingest rate (the same pattern used for retention). Combined with the trigger, a single write never touches a whole day's partition to update the summary, only the specific rows its own key maps to. If a full rebuild is ever needed, recovering from a bug, or onboarding this pattern onto existing history, doing it one partition at a time bounds the blast radius and lets the rebuild be resumed or parallelized instead of running one multi-terabyte query. Partition-level refresh is what makes that occasional full-rebuild cost tractable; it is not what keeps the steady-state refresh fast, the trigger is.
When a late-arriving correction changes a row that was already aggregated, the trigger's update branch must remove the old value's contribution before adding the new one. A naive trigger that only handles inserts will silently under- or over-count once corrections start arriving, a very easy defect to ship and not notice until the numbers are visibly wrong.
Code
CREATE TABLE daily_customer_totals (
customer_id int,
day date,
total_amount numeric(12,2) NOT NULL DEFAULT 0,
n_orders int NOT NULL DEFAULT 0,
PRIMARY KEY (customer_id, day)
);
CREATE OR REPLACE FUNCTION orders_incremental_refresh() RETURNS trigger AS $$
BEGIN
IF TG_OP = 'INSERT' THEN
INSERT INTO daily_customer_totals (customer_id, day, total_amount, n_orders)
VALUES (NEW.customer_id, NEW.created_at::date, NEW.amount, 1)
ON CONFLICT (customer_id, day) DO UPDATE
SET total_amount = daily_customer_totals.total_amount + EXCLUDED.total_amount,
n_orders = daily_customer_totals.n_orders + 1;
RETURN NEW;
ELSIF TG_OP = 'UPDATE' THEN
UPDATE daily_customer_totals
SET total_amount = total_amount - OLD.amount, n_orders = n_orders - 1
WHERE customer_id = OLD.customer_id AND day = OLD.created_at::date;
INSERT INTO daily_customer_totals (customer_id, day, total_amount, n_orders)
VALUES (NEW.customer_id, NEW.created_at::date, NEW.amount, 1)
ON CONFLICT (customer_id, day) DO UPDATE
SET total_amount = daily_customer_totals.total_amount + EXCLUDED.total_amount,
n_orders = daily_customer_totals.n_orders + 1;
RETURN NEW;
ELSIF TG_OP = 'DELETE' THEN
UPDATE daily_customer_totals
SET total_amount = total_amount - OLD.amount, n_orders = n_orders - 1
WHERE customer_id = OLD.customer_id AND day = OLD.created_at::date;
RETURN OLD;
END IF;
END;
$$ LANGUAGE plpgsql;
CREATE TRIGGER trg_orders_incremental
AFTER INSERT OR UPDATE OR DELETE ON orders
FOR EACH ROW EXECUTE FUNCTION orders_incremental_refresh();
Key points
The ON CONFLICT ... DO UPDATE upsert makes the insert branch safe under concurrent writers touching the same key. The update branch's subtract-then-add is what makes backfill corrections correct, the detail a naive implementation misses. An initial backfill over existing historical data must overwrite (DO UPDATE SET ... = EXCLUDED...), not skip on conflict, or any row already touched by the trigger before the backfill ran is left with an incorrect partial value forever, this exact bug was caught and fixed while building this demo: an initial ON CONFLICT DO NOTHING backfill left an undercount for every key the trigger had already partially updated.
Worked example
Real, executed output. A new $100.00 order inserted for an existing customer was picked up immediately with no full recompute, total_amount moved from 15,911.04 to 16,011.04 and n_orders from 55 to 56. Then a backfill correction, adding $50.00 to an existing order's stored amount:
UPDATE orders SET amount = amount + 50 WHERE id = 765;
The summary correctly reflected it, total_amount became 16,061.04 and n_orders stayed at 56, a correction to an existing order, not a new one. Verified directly against a fresh SUM/COUNT over the source rows for that same customer and day, which returned the identical 16061.04 / 56, confirming the incremental summary and a full recompute agree exactly.
Complexity
Each write costs O(1) additional work, one indexed upsert on the summary table's primary key, independent of total dataset size. A full rebuild, if ever needed, costs O(n) over whatever slice is rebuilt, which is why partition-level rebuilds bound that cost instead of paying O(n) over the entire multi-terabyte history at once.
Edge cases
Concurrent writers updating the same customer-and-day key simultaneously: ON CONFLICT DO UPDATE makes this safe at the row level since Postgres serializes conflicting upserts on the same key, but two concurrent update-branch subtract-then-add sequences on the same key are two separate statements each, not one atomic unit, so under heavy concurrent correction traffic on the same key this needs review for whether it should run inside a stricter transaction boundary. A delete of a row inserted before the trigger existed, with no corresponding insert having ever incremented the summary, would drive the summary negative, worth a floor check or a periodic reconciliation job. A correction that moves a row from one day bucket to another (changing created_at itself) is already handled correctly, since the trigger looks up the old day for the subtraction and the new day for the addition separately, but it is easy to get wrong if the day were stored as its own independent column instead of derived from created_at on each side of the trigger.
Trade-offs and pitfalls
This pattern is custom code per aggregation shape, unlike a native materialized view, every new summary needs its own trigger logic, that is the real cost of O(1) incremental refresh. Triggers add write-path latency and lock contention on the source table, every insert now also performs a summary-table upsert inside the same transaction, profile this under real write concurrency before assuming it is free. At very high write volume the trigger itself can become the bottleneck it was meant to avoid, since every source write now also pays for an extra upsert inside the same transaction; the usual next step up is a change-data-capture (CDC) pipeline, which reads the database's replication log to build the summary from a stream of row changes outside the source transaction entirely, trading the trigger's simplicity for a separate system to run and monitor. A bug in the trigger logic, as demonstrated above with the ON CONFLICT DO NOTHING backfill mistake, corrupts the summary silently while the source data stays correct, so this pattern needs its own periodic reconciliation check comparing the summary against a fresh aggregate on a sample, precisely because there is no database-native correctness guarantee the way a REFRESH-based materialized view has: however stale, a materialized view is always exactly correct as of some point in time, while a hand-rolled incremental one can be wrong at every point in time if the trigger logic has a bug.
Tell me about a time you optimized a database's performance in production. What was slow, what metrics did you start with (latency, throughput), what did you actually do (indexing, partitioning, denormalization, query rewrites, or something else), how did you validate the improvement, and what trade-offs did you accept?
Sample Answer
Direct answer
Pick one real incident with one measurable starting symptom, latency or throughput are
the usual candidates, name specifically what you tried and why each step was chosen over
the alternatives, and close with how you validated the improvement and what it cost you
elsewhere. A grab-bag of "things I've optimized over the years" reads as unfocused; one
incident, narrated end to end, reads as judgment.
A scaffold for the story
Situation. State the symptom in a measurable way: "the checkout endpoint's p95
(95th-percentile) database time was regularly exceeding 400 ms during traffic peaks" beats "the database
was slow." Latency and throughput are the two most common starting metrics, latency
because users feel it directly, throughput because it caps how much load the system can
absorb before latency degrades for everyone.
Task. What target made the incident worth fixing: an SLA being missed, a
capacity ceiling being approached, a specific complaint.
Action. Narrate three concrete techniques in the order you tried them, not a list.
The available categories a real answer usually draws from are indexing, query rewrites,
partitioning, denormalization (storing redundant, precomputed copies of data to avoid recomputing it
on every read), or connection/configuration tuning; pick whichever three
actually applied to your incident. A representative shape:
- Confirmed the query's plan with
EXPLAIN ANALYZE(the command that runs the query and
reports the real plan and timing, not an estimate) and added an index matching the
filter and sort the query actually used, verifying the plan flipped from a sequential
scan to an index scan. - Found the query was also pulling unused joined columns (
SELECT *where three columns
were used); rewrote it to select only what was needed and pushed a filter earlier in
the plan. - Noticed via the database's activity view that many requests were queued waiting for a
pooled connection, not for query execution; resized the pool and added a
statement_timeoutso one runaway query couldn't starve it.
Result. Validate against the same metric you started with, measured under
comparable load, not a different, easier one. Watch it hold up over days in production,
not just at deploy time, since delayed effects (maintenance overhead, data growth) can
erode a fix that looked clean on day one.
Trade-off accepted. A senior answer names a real cost. For example: the new index
sped up reads but added overhead to a high-insert table's write path; that was monitored
afterward and accepted as a small, budgeted latency increase on writes rather than
blocking the read fix on eliminating it entirely.
Trade-offs and pitfalls
- Listing several unrelated wins instead of one incident reads as a highlight reel, not
evidence of how you reason through a problem under pressure. - Skipping the "why this technique, in this order" reasoning makes it sound like guessing
that happened to work, which is the opposite of what the question is testing for. - Claiming an improvement without describing how it was validated, and against what
baseline, is the single most common gap; interviewers will ask. - Not naming a trade-off accepted reads as either inexperience or a curated story; real
fixes cost something somewhere.
A dashboard's average load time is 8 seconds, and the business needs it under 2. Walk through your remediation plan: how you would measure where the time is going, what quick wins you would try first, what medium-term and long-term architecture changes you would consider, and how you would know if a change made things worse.
Sample Answer
Direct answer
Don't start by changing anything. Measure where the 8 seconds actually goes, apply the
cheap low-risk fixes first, then the structural ones, and re-measure against the same
baseline metric before calling it done, with every step independently revertible. A
concrete example: on a comparably-shaped dashboard I traced a slow aggregate to two
co-occurring symptoms, table bloat and a skewed partition, that each looked fine in
isolation but compounded together.
The remediation plan
1. Measure first. Run EXPLAIN (ANALYZE, BUFFERS) (the command that actually
executes the query and reports the real plan, row counts, and buffer usage, as opposed to
plain EXPLAIN's estimate) on the dashboard's queries. Check the database's query-stats
extension for which statement accounts for the most total time. Check the dead-tuple
ratio (rows a DELETE or UPDATE left behind that haven't been reclaimed yet) in
pg_stat_user_tables. If the table is partitioned, confirm partition pruning is actually
limiting the scan to the relevant partition, and don't assume that alone means the query
is healthy.
2. Quick wins (hours, low risk). Add a missing index the plan reveals is needed.
Rewrite an obviously wasteful query, SELECT * pulling unused columns, an N+1 loop
issuing one query per row instead of one query for all of them. Add a statement_timeout
so one slow request can't starve the rest. If slight staleness is acceptable, add a short
server-side cache in front of the aggregate.
3. Medium-term (days to weeks). A materialized view (a query's result saved as its own table and refreshed on a schedule, instead of recomputed on every read) or precomputed rollup refreshed
on a schedule. Connection-pool tuning. VACUUM (reclaims space left behind by deleted or updated rows so it can be reused, though it does not shrink the file on disk) or reindex maintenance if bloat turns out to
be structural rather than a one-off. Moving the report off the primary database to a read
replica.
4. Long-term / architecture (weeks+). A real incremental pre-aggregation pipeline
that updates rollups at write time instead of computing them at read time. A partitioning
redesign if one partition has become the effective bottleneck. Moving the analytical
workload off the OLTP (online transaction processing, the primary system handling live
reads and writes) primary entirely, onto an OLAP (online analytical processing, a system
built for scanning and aggregating large amounts of historical data) engine designed for
this shape of query.
5. Knowing if a change made things worse. Track the same metric before and after,
ideally the 99th-percentile dashboard load time from real traffic, not a single manual
test run. Watch the query's estimated row count versus actual in the plan (a growing gap
signals stale statistics or a regression). Watch for new lock contention or replication
lag the fix itself introduced. Ship each change behind something independently revertible
(a config value, a feature flag) rather than as one combined deploy.
Worked example: index bloat and partition skew as co-occurring symptoms
I reproduced the bloat side of this on a real PostgreSQL 16 instance to ground the
numbers rather than assert them. A 500,000-row, 78 MB table had half of its rows deleted, executed against a real
PostgreSQL 16 instance:
DELETE FROM bench WHERE id % 2 = 0; -- 250,000 rows removed, half the table
-- after the stats collector flushes:
SELECT n_live_tup, n_dead_tup FROM pg_stat_user_tables WHERE relname='bench';
-- n_live_tup | n_dead_tup
-- ------------+------------
-- 250000 | 250000
A second delete removed half of the remaining rows, checked after the stats collector
caught up, and shows the cumulative dead-tuple count, not just this delete's share:
relname | n_live_tup | n_dead_tup | n_tup_del
---------+------------+------------+-----------
bench | 125000 | 375000 | 375000
n_dead_tup (375,000) counts every row version deleted since the last vacuum or
analyze, from both deletes combined, so it is not the same number as "how many dead
row versions are still physically sitting in the table right now": PostgreSQL prunes
some dead line pointers from a page opportunistically whenever it revisits that page
for any reason, including the second DELETE's own scan, not only during an explicit
VACUUM. A direct on-disk check with the pgstattuple extension right before running
VACUUM confirms only about 125,000 tuples were still physically dead at that point,
matching what VACUUM actually finds:
VACUUM (VERBOSE) bench;
tuples: 125020 removed, 125000 remain, 0 are dead but not yet removable
VACUUM removed 125,020 tuples, not 375,000, because roughly 250,000 of the
cumulative dead-tuple count had already been reclaimed by that opportunistic pruning
before VACUUM ever ran; a second VACUUM immediately after confirmed there was
nothing left to reclaim (0 removed, 125000 remain, where the "remain" figure is
VACUUM's report of the table's current live row count, not a dead-tuple count, so it
reads 125,000, matching the live rows still in the table, rather than 0; 0 remain
would only be printed if the table were actually empty). Despite reclaiming that space, the table's
physical size barely moved: 78 MB before and 78 MB after (81,928,192 bytes before,
81,936,384 bytes after, a difference of exactly one 8 KB page, an unmeasurable change
at the table's scale, not literally the same size the rounded MB figures suggest).
VACUUM without FULL
marks space reusable for future inserts; it does not return it to the operating
system. The table now holds 125,000 live rows, one quarter of the 500,000 rows it was
originally sized for, inside that same 78 MB file, so every scan of that table,
including the dashboard's aggregate, was reading roughly four times the pages the live
data needed, not twice.
Now add the partition angle: if this table is partitioned by month and the query filters
by date, pruning correctly narrows the scan to one partition, say the current month's,
which is also the one absorbing most of the day's writes and deletes (a "hot" partition).
Pruning working correctly (touching only 1 of 12 partitions) doesn't mean that partition
is healthy; if it's carrying a disproportionate share of the table's churn, it can be
bloated exactly like the example above while its 11 siblings are fine. That's why "the
plan shows pruning is working" is not the same conclusion as "the query is fast", you
still have to check the one partition doing all the work.
Trade-offs and pitfalls
- Reaching for an index before checking
EXPLAIN ANALYZErisks fixing the wrong problem;
bloat and stale statistics both produce slow plans that look, at a glance, like a
missing-index problem. - A cache "fixes" the symptom without addressing what staleness stakeholders can actually
accept; a silently stale number on a business dashboard can be worse than a slow but
correct one, so that's a product conversation, not just a technical one. - Precomputing the aggregate shifts cost to the write path and introduces a staleness
window; decide the acceptable window with the business rather than assuming. - Validating against a single idle-system test run and calling it done misses whether the
fix holds under real concurrent load and ongoing bloat regrowth; a fix that looks great
in staging can still fail the p99 (99th-percentile) target in production a week later.
You need to cut p95 read latency for a global user base from 120 milliseconds to under 40 milliseconds at 50,000 reads per second. Design a read-scaling architecture layering edge or CDN caching, in-memory caching, and read replicas. How would you decide what to cache and for how long, using a pattern like stale-while-revalidate, and what consistency guarantees would you be willing to give up to hit that target?
Sample Answer
Direct answer
Layer three caches in front of the primary: an edge / content delivery network (CDN) tier for data that is the same for everyone or close to it, a regional in-memory cache (Redis or similar) for personalized or fast-changing data, and regional read replicas as the fallback source once both caches miss. Whether this works comes down to one number: because a request from a distant user has to cross real physical distance to reach a single origin region, a cache miss alone can already exceed a 40 millisecond budget, so hitting p95 (95th-percentile) latency under 40 milliseconds at 50,000 reads per second requires the combined edge-plus-in-memory cache hit ratio to sit at roughly 95% or higher. The guarantee you give up to get there is strong, always-fresh reads for everyone except the user who just wrote the data: cached and replica-served reads become eventually consistent, bounded by a stated staleness window per data type, while the acting user's own write is routed around the cache so they never see their own action appear to have failed.
Structured elaboration
Why distance alone rules out an origin-only design
Signal propagation through fiber-optic cable travels at roughly two-thirds the speed of light in vacuum. That is a physical limit, not a database tuning problem: a round trip between two points thousands of kilometers apart costs real milliseconds no query optimization can remove. A single-origin-region architecture, however well-tuned the database is, cannot put a p95 under 40 milliseconds for users on the other side of the planet. That is the architectural reason edge and regional caching are not optional here, they are the only way to physically shorten the path for most requests.
Classifying data to decide what goes in which tier
| Data shape | Cache tier | Reasoning |
|---|---|---|
| Same for (almost) every user, changes rarely (catalog, public profile pages, static configuration) | Edge / CDN, long time-to-live (TTL) | Few distinct cache keys shared by many users means very high hit ratios are achievable at the point closest to the user |
| Personalized or computed per-user (a feed, an aggregate, a leaderboard) | Regional in-memory cache, shorter TTL | Too many distinct keys for edge caching to pay off; centralizing in one regional cache still avoids hitting the database on every read and is easier to invalidate correctly |
| The user's own just-written data | No cache; routed to the primary or a replica confirmed to have absorbed the write | This is the one case where staleness is not acceptable: a user should never read back something that looks like their own write silently failed |
| Anything that must be correct at the instant of read (for example, a balance immediately before authorizing a debit) | No cache, always the primary | An explicit, named exception, not something quietly caught by a general TTL |
Stale-while-revalidate (SWR) and why it, not active invalidation, should carry most of the load
Under stale-while-revalidate, a cache serves the value it already has immediately, even if its soft freshness window has passed, while kicking off a background refresh from the layer below (an in-memory cache refreshing from a replica, or an edge node refreshing from origin). The next request after the refresh completes gets the new value. The user never waits on the refresh. Compare that to actively pushing invalidation messages to every edge node whenever underlying data changes: that turns cache freshness into its own distributed consistency problem across every point of presence (PoP) in the network, and for most read-mostly data it is slower and more fragile than simply letting a short TTL expire. Reserve active invalidation for the small set of high-value keys where even a short TTL window of staleness is unacceptable (for example, a price that just changed).
Read replica routing and the freshness floor they set
Give each major region its own read replica, route reads to the nearest one, and continuously track each replica's replication lag (how far behind the primary it is). If a replica's lag exceeds the staleness budget promised by the cache tiers above it, fail toward the primary for that traffic rather than silently serving data staler than what was promised. Replication lag is also the floor beneath every promise this design makes: no cache TTL should be tighter than what the replica feeding it can actually guarantee.
Worked example
Two numbers actually gate this design: how close a hit has to be to beat 40 milliseconds, and how far a miss is from a distant origin region. The second one is physics, not opinion, so it is worth computing rather than assuming.
C_VACUUM = 299_792.458 # km/s, speed of light in vacuum
FIBER_INDEX = 1.47 # typical single-mode fiber refractive index
v_fiber = C_VACUUM / FIBER_INDEX # km/s, signal speed in fiber
routes_km = {
"cross-country (US east-west)": 4200,
"cross-continent (US-Europe)": 6000,
"antipodal (US-Singapore)": 15300,
}
for name, km in routes_km.items():
rtt_ms = 2 * km / v_fiber * 1000 # round trip: there and back
print(f"{name}: {rtt_ms:.1f} ms of propagation delay alone")
# p95 sizing: if a hit is comfortably under 40ms and a miss stays near today's
# 120ms baseline (the same WAN + database round trip that produced it), then
# for the 95th-percentile request to land under 40ms, at most 5% of requests
# may miss.
miss_budget_pct = 5.0
print(f"required combined edge+in-memory hit ratio >= {100 - miss_budget_pct:.0f}%")
read_qps = 50_000
target_hit_ratio = 0.97 # 2 points of margin above the 95% floor
miss_qps = read_qps * (1 - target_hit_ratio)
regions = 5
print(f"miss traffic reaching replicas: {miss_qps:,.0f} reads/sec, "
f"~{miss_qps/regions:,.0f} reads/sec per region across {regions} regions")
write_tps = 5_000
print(f"write target: {write_tps:,} TPS (transactions per second), unaffected by "
f"any of the caching above; it lands entirely on the primary")
Output (executed):
cross-country (US east-west): 41.2 ms of propagation delay alone
cross-continent (US-Europe): 58.8 ms of propagation delay alone
antipodal (US-Singapore): 150.0 ms of propagation delay alone
required combined edge+in-memory hit ratio >= 95%
miss traffic reaching replicas: 1,500 reads/sec, ~300 reads/sec per region across 5 regions
write target: 5,000 TPS (transactions per second), unaffected by any of the caching above; it lands entirely on the primary
Two things fall out of this. First, cross-continent and antipodal round trips already exceed 40 milliseconds on propagation delay alone, before a single database query runs, which is why no amount of origin-side tuning closes this gap for a global user base; the cache hit ratio has to do essentially all the work. Second, at a 97% hit ratio the replica tier only has to absorb about 300 reads per second per region, comfortably inside what a single well-indexed read replica handles, so the hard part of this design is the cache hit ratio, not replica capacity.
The companion write target matters here because it sets the freshness floor discussed above: 5,000 writes per second all land on the primary (caching and read replicas do nothing for writes), and the primary has to both sustain that rate and ship its write-ahead log (WAL, the durable record of every change) to every regional replica fast enough that replication lag stays inside whatever staleness window the cache tiers are promising. If write volume grows further, replication lag grows with it, and that is the point at which every cache TTL and every "how stale can a replica read be" assumption in this design needs to be re-validated, not treated as a one-time decision.
flowchart TD
U[Global user request] --> E["Edge / CDN PoP\ncache-only, long TTL"]
E -- hit --> RU[Response]
E -- miss --> R["Regional in-memory cache\nRedis, short TTL, SWR"]
R -- hit --> RU
R -- miss --> RR[Regional read replica]
RR -- fresh enough --> RU
RR -- lag over budget --> P[Primary]
P --> RU
W[Write / user action] --> P
P -- WAL replication --> RR
P -- fresh write --> OW["Read-your-own-write\nbypass cache, route to primary"]
Trade-offs & pitfalls
- Cache stampede on expiry. When a hot key's TTL lapses, many concurrent requests can all miss at once and hammer the layer below simultaneously. Stale-while-revalidate already helps (it keeps serving the old value during refresh), but pair it with request coalescing (only one in-flight refresh per key, other requests wait on that result) and TTL jitter (randomizing expiry slightly so many keys do not expire in the same instant).
- Over-caching mutable, high-stakes data. The temptation to cache everything for uniform simplicity is real; resist it for anything where staleness has a real cost (money, inventory at the moment of purchase, access control decisions). Naming the exceptions explicitly, as this design does, is safer than a blanket policy with silent edge cases.
- Edge caching does not help personalized data. If a page or response is unique per user, the edge tier will have as many cache keys as users and its hit ratio collapses; that traffic belongs in the regional in-memory tier, not forced into the edge tier for architectural tidiness.
- The 40ms budget math here is a simplified two-class model (hit versus miss), used to size the required hit ratio, not a promise about the real latency distribution. Validate the actual design against measured percentile telemetry once it is live; a real cache miss is not always exactly 120 milliseconds, and a real hit is not always a few milliseconds.
- Read-scaling work does not touch the write path. If the write target grows past 5,000 TPS, that is a primary-capacity and replication-fan-out problem, a different constraint than anything cache tiers or read replicas solve, and it should be sized and reasoned about on its own terms.
Unlock Full Question Bank
Get access to all 32 Database Performance Tuning and Scaling interview questions and detailed answers.
Sign in to ContinueJoin thousands of developers preparing for their dream job.