Replication, Partitioning, and Sharding Questions
Scaling and distributing data across nodes: primary-replica and multi-primary replication, read-replica scaling, horizontal partitioning, and sharding strategies with their key-selection and rebalancing challenges. Covers replication lag, failover and split-brain handling, cross-shard operations such as joins, distributed transactions, and global secondary indexes, and the operational cost of a partitioned topology. Key to designing databases that scale horizontally.
Write an operational playbook for recovering from a single-shard failure that causes read/write errors for a subset of users. Include detection, immediate mitigation steps, data integrity checks, and post-recovery validation steps relevant to a sharded SQL database.
Sample Answer
Direct answer
A single-shard failure in a sharded SQL database only affects the subset of users whose data lives on that shard, which is both the good news (the blast radius is bounded) and the trap (it is easy to under-react to an outage that "only" affects a fraction of users but is, for those users, a full outage). The playbook is: detect which shard and confirm the primary is actually unreachable (not just slow), fence it so it cannot accept writes if it comes back unexpectedly, promote the most caught-up in-sync replica, repoint the routing layer's shard map, then run data-integrity checks before declaring the incident over.
The playbook
flowchart TD
A[Alert: elevated errors on shard 7] --> B{Is shard 7 primary reachable?}
B -->|no| C[Fence old primary: revoke credentials, block writes]
C --> D[Promote most-caught-up in-sync replica to primary for shard 7]
D --> E[Repoint routing layer: shard map entry for 7 -> new primary]
B -->|yes, but degraded| F[Diagnose: disk, connections, replication lag on shard 7]
F --> G{Root cause fixable in place?}
G -->|yes| H[Mitigate in place, no promotion]
G -->|no| C
E --> I[Run data-integrity checks: row counts, checksums vs other shards' expectations]
H --> I
I --> J[Validate application-level reads/writes on shard 7 succeed]
J --> K[Post-recovery: reconcile any writes lost in the gap, document timeline]
1. Detection. The alert should identify the specific shard, not just "database errors," which means shard-aware monitoring (per-shard error rate, per-shard connection health, per-shard replication lag, meaning how far behind the primary this shard's replica has fallen) has to exist before the incident, not be built during it. Confirm scope immediately: query the shard-to-tenant or shard-to-user-range mapping to know exactly which users are affected, since this drives both the urgency (how many users, and are any high-priority accounts among them) and the communication (what to tell support or status pages).
2. Immediate mitigation. First check whether the primary is truly unreachable versus merely degraded (high load, a stuck query, disk pressure); a degraded-but-reachable primary may be fixable in place (kill a runaway query, add I/O capacity) without the risk of a promotion. If it is genuinely unreachable, fence it (revoke its credentials or otherwise ensure it cannot accept writes if it unexpectedly comes back mid-recovery, which prevents a split-brain where two nodes for the same shard both think they are primary), then promote the most caught-up in-sync replica for that shard specifically, and update the routing layer's shard map so application traffic for that shard's key range goes to the new primary.
3. Data-integrity checks. Before declaring the shard healthy, verify: row counts and, where feasible, checksums on recently-written tables against expectations (comparing against a recent backup or against another shard's equivalent tables if the schema allows structural comparison); confirm the promoted replica's applied position (its last-replayed log sequence number) to bound how many, if any, in-flight writes from the old primary may not have made it across, which tells you the actual recovery point objective (RPO) hit for this specific incident, not a theoretical one; and run a small set of representative application-level read/write operations end-to-end against the shard, not just database-level health checks, since a shard can be database-healthy while still broken from the application's perspective (wrong connection string, stale cached routing entry).
4. Post-recovery validation. Reconcile any writes that were in flight but unconfirmed at the moment of failure (check application-level idempotency logs or an outbox table if one exists, to identify anything that needs replaying or compensating). Document the actual timeline (detection time, fencing time, promotion time, full recovery time) against the shard's specific RTO target, and capture the RPO actually incurred (how many seconds, or how many specific rows, of writes were lost, if any) for the post-mortem, since "we recovered" and "we recovered with zero data loss" are different claims that need different evidence.
Trade-offs and pitfalls
The most common mistake is skipping the fencing step under time pressure ("it's clearly dead, let's just promote"), which is exactly the scenario that produces split-brain if the old primary was actually just network-partitioned rather than dead, and comes back mid-incident still believing it is primary. The second common mistake is declaring victory once the promoted replica is serving traffic without running the data-integrity checks, which defers discovering data loss or corruption from "during the incident, with full attention on it" to "days later, when a user reports a missing record," at which point root-causing it is far harder.
Explain how database-level replication lag can affect read-after-write consistency in systems using read replicas. As a data engineer, describe patterns to ensure users see their own recent writes (for example, sticky sessions, session routing, read-your-writes guarantees) and trade-offs involved.
Sample Answer
Direct answer
Replication lag breaks read-after-write consistency (the guarantee that a user who just wrote something sees that write on their very next read) whenever the read happens to land on a follower that has not yet applied it. There are four standard patterns to guarantee a user sees their own recent writes despite this: sticky sessions (route a user's reads to the same node their write went to, for a window of time), session-scoped read-your-writes tracking (remember the position a user's last write reached, such as an LSN, log sequence number: a monotonically increasing marker of a database's write history, and only serve their reads from a replica that has caught up to at least that position), forcing reads-after-writes to the primary for a short window, and, less commonly, bounded-staleness reads that block briefly until a replica catches up rather than either redirecting or serving stale data outright.
The four patterns, and their trade-offs
| Pattern | How it works | Trade-off |
|---|---|---|
| Sticky sessions | Load balancer or client routes a given session/user consistently to the same node (often the primary, or a designated "home" replica) for some time window | Simple to implement, but concentrates that user's reads on one node (losing the read-scaling benefit of replicas for that user) and needs a fallback if the sticky node becomes unavailable |
| Read-your-writes token / session tracking | After a write, the client stores the LSN or a similar log position it reached; on the next read, only route to a replica whose applied position is at or past that token | Precise (only redirects exactly the reads that need it) but requires plumbing the token through the client or session layer and requires each replica to expose its current apply position for comparison |
| Force-read-from-primary flag | Explicitly mark reads that must be fresh (e.g., "show me my own just-submitted order") to bypass replicas and hit the primary directly | Zero staleness for the marked reads, but adds primary load proportional to how many reads get flagged, so it is only safe when a small, well-identified subset of reads actually needs it |
| Bounded-staleness / wait-for-replica reads | The replica delays answering a read until its applied position reaches the requested minimum, instead of redirecting elsewhere | Avoids extra primary load and avoids sticky routing, but turns a routing problem into a latency problem: the read blocks for however long the replica is behind, which is unbounded unless capped with a timeout-and-fallback |
Worked example: a social-network timeline
A user posts a comment, then immediately reloads their own timeline. If the reload hits a replica still 800 milliseconds behind (a fully normal amount of replication lag under moderate load), the comment appears to have vanished, which reads as a bug even though nothing is actually wrong. A concrete implementation: on write, the API server captures the WAL position (or an equivalent monotonic write marker) the write reached and returns it to the client as an opaque read-your-writes token embedded in the response (client-side technique); the client attaches that token to its next read request; the read-routing layer (server-side technique) checks each candidate replica's current applied position and only sends the request to one that has caught up, falling back to the primary if no replica has caught up within a short timeout (for example 50 milliseconds), rather than blocking indefinitely. This combination avoids sticky-session's node-affinity cost and avoids sending every read to the primary, while still guaranteeing the user sees their own comment.
Implementation details worth naming explicitly
A connection pooler alone does not solve this: PgBouncer, the common PostgreSQL pooler, is a pure TCP-level connection pooler with no query inspection and no awareness of a replica's replication position, so it cannot itself implement token-aware read/write splitting, and pointing it at a random or round-robin pool of replicas will pick one that has not caught up. The token-aware routing logic has to live in a layer that actually understands queries and replica state: application-level routing that picks a connection target per request based on the token, or a proxy built for query routing rather than plain pooling (for example Pgpool-II's load-balancing mode, which does inspect queries to split reads from writes) sitting in front of, or instead of, a plain pooler like PgBouncer, which continues to do only what it is built for: pooling connections to whichever destination it is told to use. A retry-on-lag pattern (catch the case where no replica is caught up, retry against the primary once, rather than retrying the same lagging replica) keeps tail latency bounded. Session-affinity state (which token belongs to which in-flight session) needs to live somewhere the load balancer or API layer can cheaply access it on every read, typically in the request context itself (the token travels with the request) rather than in a separate lookup service, to avoid adding a new dependency to every read's critical path.
Trade-offs and pitfalls
Sticky sessions are the easiest pattern to ship first, but teams often leave them in place permanently rather than migrating to a token-based approach, which silently caps how much read-scaling benefit replicas actually provide as traffic grows, since sticky users never spread across the replica fleet. Forcing all writes' subsequent reads to the primary "to be safe" is the opposite mistake: it looks correct but reintroduces exactly the primary-load problem replicas exist to solve, for reads that, in most product flows, do not actually need same-millisecond freshness. The right default is to reserve the primary-forcing pattern for the specific, small set of user actions where staleness is visibly wrong (their own just-written content), and use the token-based approach for everything else.
You receive reports of stale reads for a small percentage of users. Outline a step-by-step troubleshooting plan to identify root cause: what logs, traces, metrics, and checks would you run; how to reproduce; and how to validate whether the issue stems from replication lag, caching, DNS, or client behavior.
Sample Answer
Direct answer
Stale reads for a small percentage of users almost always trace back to one of four cause categories: replication lag on a specific replica, a caching layer serving an outdated value, a DNS or load-balancer routing quirk sending a subset of users to a stale node, or client-side behavior (a browser cache, a retried request, a session stuck on an old connection). The troubleshooting plan is to check them in that order, cheapest and most likely first, and to reproduce the issue against a specific replica before guessing at a fix.
Step-by-step plan, ordered by cause category
- Reproduce and scope it. Identify whether "a small percentage of users" correlates with a specific replica, a specific region, a specific client version, or a specific time window (for example, only right after a deploy or only during a nightly batch job). Pull the affected users' request logs and check which backend node/replica actually served each stale read; if a distributed trace (a request-scoped record of every service and backend a single request touched, including which specific database replica it hit) exists for these requests, it is the fastest way to pin down the exact replica and timing without guessing from aggregate logs. If neither the logs nor traces currently capture which replica served a given read (many setups only log "database", not "which replica"), fix that first, because without it every other step is guesswork.
- Network/replication category: Check the lag metrics (for example PostgreSQL's
pg_stat_replication.replay_lag, or MySQL'sSeconds_Behind_Source) for the specific replica(s) identified in step 1, at the specific timestamps of the reported staleness, not just current lag. A replica that spiked to 15 seconds of lag for two minutes during a burst and recovered will look perfectly healthy by the time you check it live, which is why historical, not current, lag data matters here. - Config category: Confirm the read-routing configuration actually matches intent: is a connection pooler or load balancer sending reads to a replica that was supposed to be excluded (e.g., a replica dedicated to backups, or one flagged unhealthy but not actually removed from rotation)? Check for a recently changed routing rule or a replica added to the pool without its lag-based health check wired up.
- I/O category: On the suspect replica, check disk I/O saturation and apply-thread backlog around the incident window; a replica falling behind because its own disk cannot keep up with the apply rate (as opposed to a network-caused lag) needs a different fix (more I/O capacity, or removing competing load like backups from that node) than a network issue does.
- Caching category: If the read path includes an application cache or CDN layer, check whether the cache key for the affected records was invalidated at write time; a stale cache entry produces identical symptoms to replication lag but requires cache-invalidation logic fixes, not database fixes, so ruling this in or out early avoids chasing the wrong system.
- DNS/client category: Rule out a stale DNS entry still pointing at a decommissioned or unhealthy node (check DNS time-to-live and recent record changes), and rule out client-side caching or a client stuck on a long-lived connection to a node that was since removed from the healthy pool.
Worked example: BI reporting service-level agreement (SLA)
A business-intelligence dashboard shows a metric that occasionally contradicts a more recent number end users already saw elsewhere in the product. Reproducing it: correlate reported instances against the specific reporting replica's historical lag graph, and find recurring lag spikes to 40 to 90 seconds every night at 02:00, matching a nightly ETL (extract-transform-load) batch job hitting that same replica. Root cause: the reporting replica is shared between BI reads and the nightly batch job, and the batch job's write-heavy phase (if it writes anywhere near this replica's apply path) or its heavy read load starves the replica's apply thread of I/O. Remediation options: move BI reads off the shared replica onto a replica (or dedicated change-data-capture pipeline, CDC: a stream of row-level changes consumed to build a separate, purpose-built analytics store) that is isolated from the nightly job's resource contention, and add a per-replica lag SLA specifically for the reporting path.
Trade-offs and pitfalls
The single biggest time-waster in this kind of investigation is jumping straight to "it's probably replication lag" and tuning replication settings before confirming, with actual per-replica historical data, that lag was even elevated at the reported times; teams have spent days tuning replication only to find the real cause was a stale cache key or a misconfigured routing rule. The second pitfall is checking only current-state metrics instead of the historical window around each report, since transient lag spikes that already recovered are invisible to a live check but are exactly what caused the user-visible symptom.
Design a monitoring and alerting strategy for a sharded database cluster. List per-shard and cluster-wide metrics to collect (latency percentiles, QPS, replication lag, disk utilization), alert thresholds and severity, and what actions/automation or runbooks should be triggered for common alerts.
Sample Answer
Direct answer
Monitor at two levels that answer two different questions: per-shard metrics tell you which specific shard is unhealthy, and cluster-wide aggregates tell you whether the cluster as a whole is meeting its service-level objective (SLO, the target level of service the team has committed to). Alert on the aggregate for paging (nobody should be paged because one shard out of two hundred has a blip), and keep the per-shard view for diagnosis once a page fires. Thresholds should be set relative to each shard's own healthy baseline, not one global number, since shards can have legitimately different steady-state load.
Metrics to collect
Per-shard:
- Query latency, p50/p95/p99 (95th and 99th percentile response time, the tail behavior that matters more than the average for user-facing impact)
- Queries per second (QPS) and the read/write split
- Replication lag to that shard's replicas (for example Postgres's
pg_stat_replication.replay_lag, the delay between a write committing on the primary and becoming visible on a given replica) - Disk utilization (both space used and IO throughput/IOPS, input/output operations per second, since a shard can be IO-saturated well before it runs out of space)
- Connection count against the configured limit
- Lock wait time and deadlock rate (how often two transactions end up blocking each other in a cycle, forcing the database to kill one so the other can proceed)
Cluster-wide:
- Aggregate QPS and error rate across all shards
- Number of shards currently unhealthy (below their own baseline) as a single number, so a dashboard shows "3 of 200 shards degraded" instead of 200 individual panels
- Rebalance/resharding operations in flight and their progress
- Cross-shard query fan-out latency (the scatter-gather pattern from a query that has to touch multiple shards is only as fast as its slowest shard, so this is a distinct signal from any single shard's own latency)
Thresholds, severity, and what should page
| Signal | Warning | Critical (pages on-call) | Rationale |
|---|---|---|---|
| Per-shard p99 latency | 2x that shard's 7-day rolling baseline | 5x baseline, sustained 5 minutes | A short spike is normal; a sustained multiple of baseline is not, and using the shard's OWN baseline avoids a single global number being wrong for every shard with different traffic |
| Replication lag | 30 seconds | 5 minutes, or growing rather than stable | Below 30s is invisible to most consumers; past a few minutes, stale-read risk and failover data-loss risk (if the primary fails before a lagging replica has caught up, whatever had not yet replicated is gone once that replica gets promoted) both become real |
| Disk utilization | 75% | 90%, or projected to fill within 24 hours at current growth rate | Running out of disk on a primary is a hard outage, not a degradation, so this needs lead time to act, not just a threshold at the edge |
| Connection count | 70% of configured limit | 90% | Hitting the connection limit fails new requests outright; this should never be a surprise |
| Unhealthy shard count (cluster-wide) | 1 percent of shards | 5 percent of shards, or any single shard fully unreachable | A rising unhealthy-shard count often indicates a systemic cause (a bad deploy, a shared dependency failing) rather than 200 independent unlucky shards |
What should happen automatically versus what pages a human
- Automatic, no page: a single shard's connection pool (its fixed set of reusable, already-open database connections) exhausting temporarily under a burst (if the application has retry-with-backoff, this self-heals); routine rebalance operations proceeding on schedule.
- Page, with a runbook attached, not just an alert: replication lag critical (runbook: check for a long-running transaction blocking apply, check network throughput between primary and replica, decide whether to fail over or wait); disk critical (runbook: confirm whether growth is organic or a leak/bug, trigger an emergency resize or archive-and-purge); unhealthy-shard-count critical (runbook: check for a common cause across the affected shards before treating them as independent incidents).
- A runbook is not optional documentation, it is the alert's payoff. An alert that pages someone with no linked, concrete next action just moves the "what do I do now" research into the middle of the incident instead of before it.
Trade-offs and pitfalls
- Alerting on raw per-shard thresholds instead of each shard's own baseline produces both false pages (a shard that's always a bit hotter than average) and missed pages (a shard whose baseline is unusually low, where a real problem never crosses a one-size-fits-all number).
- Too many per-shard alerts wired directly to paging is how on-call burns out. Fan noisy per-shard signals into the aggregate unhealthy-shard-count metric for paging, and keep the per-shard dashboards for triage after the page, not as 200 separate alert rules.
- A metric with no owner decays. Review thresholds periodically against actual incident history: a threshold that never fires might be too loose to catch real problems, and one that fires constantly and gets silenced is worse than no alert at all, since it trains people to ignore it.
Tail latency spikes increased after increasing shard count. Outline how you would diagnose root causes (network issues, coordination overhead, compaction/GC, metadata lookups), which tools and traces you'd use, and what optimizations you would try to reduce tail latency in a sharded datastore.
Sample Answer
Direct answer
Diagnose a post-shard-count-increase tail-latency regression by checking, in order, whether queries are actually using the shard key (an un-keyed query now fans out to more shards than before, which alone degrades p99 through pure order-statistics, no per-shard slowdown required), then network and coordination overhead (more shards means more connections, more routing hops, more chances one participant in a fan-out is briefly slow), then compaction/GC pauses and metadata-lookup cost on individual shards. The mechanism that most directly explains "latency got worse specifically because shard count went up" is scatter-gather fan-out width itself, demonstrated below.
Structured elaboration
Root-cause checklist, in likely-impact order
- Un-keyed queries forcing full fan-out: if a query's
WHEREclause does not include the shard key (the column value the router uses to know which single shard should hold a given row), the router cannot route it to one shard, it must scatter-gather to ALL of them. As shard count grows, that same query fans out wider, and p99 latency, the max across the fan-out, degrades purely from that width even if every individual shard's own performance is completely unchanged (demonstrated below). - Full-table scans on individual shards: check the query plan on each shard (
EXPLAIN ANALYZE, the command that prints the actual steps and timing the database used to run a query) for a sequential scan (reading every row in the table, in order, to find matches) where an index scan (using an index to jump straight to the matching rows instead of reading the rest of the table) was expected; this often surfaces alongside symptom 1 above, because a query written assuming it always hits one small, well-indexed shard can silently degrade into a full scan per shard once it starts fanning out to many, especially if the shard-routing key was dropped from an index or from the query itself during some earlier change. - Network and coordination overhead: more shards means more open connections from the router/proxy tier, more routing-table lookups per request, and more individual RPCs to wait on for any multi-shard query; check connection-pool saturation and per-hop latency in a distributed trace, not just aggregate query time.
- Compaction and garbage collection pauses: on LSM-tree-backed storage engines (engines like RocksDB or Cassandra's storage layer that buffer writes in memory and periodically merge them into sorted files on disk instead of updating rows in place) or garbage-collected ones, adding shards (if it meant redistributing existing write load onto more, smaller instances, or if total write volume grew alongside the resharding) can shift the compaction schedule (the background job that merges and rewrites those sorted files to keep reads fast, briefly competing with foreground traffic for CPU and IO while it runs) or GC schedule on individual nodes; check per-node compaction/GC pause metrics for a correlation in time with the latency regression.
- Metadata/routing lookups: confirm the routing-table cache hit rate has not dropped; more shards means a larger routing table, and if any tier is doing an uncached or poorly-cached lookup per request, that cost grows with the size of the table it is searching.
Tools and traces
- Distributed tracing with one span per shard hop for a fan-out query, so you can see the actual per-shard latency distribution behind an aggregate p99, not just the aggregate number.
- Per-shard
EXPLAIN ANALYZEon the specific query pattern that regressed, to catch symptom 2 directly. - Connection-pool and routing-cache metrics (pool wait time, cache hit rate) to catch symptom 3 and 5.
- Storage-engine internal metrics (compaction queue depth, GC pause histograms) correlated by timestamp against the latency regression, to catch symptom 4.
Optimizations to try
- Push the shard key into every query that can legitimately use it, and add monitoring that flags un-keyed queries against a sharded table as a class of technical debt, since they get strictly more expensive as the cluster grows, not just today.
- For queries that genuinely need to touch many shards, narrow the fan-out where possible (route only to shards known to hold relevant data, if any partitioning scheme allows that) rather than always scattering to every shard.
- Keep-alive and pool connections between routers and shards rather than establishing a fresh connection per request, to remove per-hop connection-setup cost from the tail.
- Cache routing metadata aggressively at every router/proxy instance (a local, periodically-refreshed copy of the routing table, not a lookup per request) so metadata lookup cost does not scale with routing-table size on the hot path.
Worked example
The fan-out-width mechanism, isolated: this is an illustration of pure order statistics using synthetic, unitless latency samples (not a measurement of any real system, deliberately unitless to avoid implying otherwise), holding each shard's own latency distribution completely fixed while only the fan-out width N changes, to show that p99 degrades from width alone.
import random
random.seed(42)
def per_shard_latency_au():
return max(0.1, random.gauss(5, 1.5)) # arbitrary units (a.u.), unchanged as N grows
def simulate_fanout(n_shards, n_queries=8000):
samples = sorted(max(per_shard_latency_au() for _ in range(n_shards)) for _ in range(n_queries))
p99 = samples[int(len(samples) * 0.99)]
return p99, sum(samples) / len(samples)
single_shard_samples = sorted(per_shard_latency_au() for _ in range(8000))
single_p99 = single_shard_samples[int(len(single_shard_samples) * 0.99)]
print(f"single-shard p99 (n=1, no fan-out): {single_p99:.2f} a.u. <- unchanged baseline throughout")
for n in (8, 32, 128, 512):
p99, mean = simulate_fanout(n)
print(f"fan-out width N={n:>4}: p99={p99:6.2f} a.u. ({p99/single_p99:.2f}x baseline), mean={mean:5.2f} a.u.")
Output:
single-shard p99 (n=1, no fan-out): 8.48 a.u. <- unchanged baseline throughout
fan-out width N= 8: p99= 9.62 a.u. (1.13x baseline), mean= 7.15 a.u.
fan-out width N= 32: p99= 10.06 a.u. (1.19x baseline), mean= 8.10 a.u.
fan-out width N= 128: p99= 10.66 a.u. (1.26x baseline), mean= 8.88 a.u.
fan-out width N= 512: p99= 11.25 a.u. (1.33x baseline), mean= 9.56 a.u.
p99 climbs from 1.00x to 1.33x baseline purely from widening the fan-out from 8 to 512 shards, with the per-shard distribution never changed. This is the concrete argument for checklist item 1 above: a p99 regression that lines up in time with a shard-count increase, on a query pattern that fans out, does not require ANY individual shard to have gotten slower to be fully explained; it can be the arithmetic of "the tail is the max of more independent draws," and checking for that first, before chasing a per-shard regression that may not exist, saves real investigation time.
Trade-offs and pitfalls
- Jumping straight to compaction/GC tuning (item 4) because it is a familiar lever is a common wrong turn; if the actual cause is fan-out width (item 1), no amount of per-shard tuning will fix a p99 that is fundamentally an order-statistics effect of querying more shards per request.
- A full-table-scan symptom on an individual shard can be a genuine, independent regression OR a downstream consequence of item 1 (a query that used to be a fast, single-shard indexed lookup silently became a per-shard full scan once it started fanning out without its shard-key filter surviving the rewrite); check both, in that order, rather than assuming they are unrelated.
- Reducing fan-out by routing only to "likely relevant" shards is a real optimization but only works if your partitioning scheme actually correlates with the query's filter; applying it blindly to a genuinely unpredictable access pattern just reintroduces missed-data bugs to chase a latency number.
- Treating this as purely a network or infrastructure problem, without first confirming whether the regressed queries are even using the shard key, risks spending real engineering effort on connection pooling and tracing infrastructure when the actual fix is adding a
WHEREclause.
Unlock Full Question Bank
Get access to all 16 Replication, Partitioning, and Sharding interview questions and detailed answers.
Sign in to ContinueJoin thousands of developers preparing for their dream job.