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.
Design an automated failover system for primary-replica database pairs that minimizes split-brain risk. Cover health checks, fencing mechanisms, promotion safety checks, and how you would validate and audit automatic promotions. Also explain what data can be lost during the failover and how RPO and RTO differ between manual and automatic failover.
Sample Answer
Direct answer
An automated failover system for primary-replica pairs needs four pieces working together to minimize split-brain risk: a consensus-based health check (a majority, not a single observer, must agree the primary is actually down), fencing that guarantees the old primary cannot keep accepting writes once it loses leadership, a promotion-safety check that refuses to promote a replica too far behind, and an audit trail proving what happened and when. This is the same design pattern Patroni implements for PostgreSQL clusters, and it generalizes to any primary-replica system.
Architecture
flowchart TD
A[Health check: primary missed N heartbeats] --> B{DCS quorum reachable?}
B -->|no quorum| C[Do not promote: minority side, refuse writes]
B -->|quorum reached| D[Elect candidate: most caught-up in-sync replica]
D --> E[Fence old primary via watchdog/STONITH]
E --> F{Fencing confirmed?}
F -->|no| G[Abort promotion, alert on-call: cannot guarantee safety]
F -->|yes| H[Promote candidate to primary]
H --> I[Update DNS/service discovery to new primary]
I --> J[Audit log: who/what promoted, at what LSN, data-loss window]
J --> K[Old primary rejoins later as a follower once fencing is lifted]
Health checks. A single node observing "the primary didn't respond" is not enough evidence, since that node itself could be the one that is network-partitioned. Route the decision through a distributed configuration store (DCS: a small, separately-replicated coordination service such as etcd, Consul, or ZooKeeper) that requires a majority, or quorum, of its own members to agree before any leadership change proceeds; this is the same mechanism Patroni uses for PostgreSQL. If the surviving side of a partition cannot reach a majority of the DCS, it must not promote, even if it genuinely believes the primary is down, because it cannot distinguish "the primary is dead" from "I am the one who is isolated."
Fencing mechanisms. Before any promotion, the old primary must be guaranteed unable to keep accepting writes. Two layers, used together: revoke its ability to renew its leadership lease in the DCS (a soft fence: it should demote itself once it notices the lease expired), and a hardware or OS-level watchdog that forcibly resets the host if the soft fence has not been confirmed within a bounded time (a hard fence, sometimes called STONITH, "shoot the other node in the head"). The hard fence exists specifically for the case where the old primary is unresponsive rather than well-behaved, since a well-behaved node would have already demoted itself.
Promotion safety checks. Never promote a replica that is more than a defined, bounded amount behind the primary's last known position (measured by comparing log sequence numbers, LSNs). If the most caught-up available replica is still too far behind, the safer action is to alert and wait (accepting downtime) rather than promote and silently accept a larger-than-intended data loss; that threshold (how far behind is "too far") is a deliberate, pre-agreed trade-off between availability and data loss, not a default left unexamined.
Validating and auditing automatic promotions. Every automatic promotion writes an audit record: which node was promoted, what LSN it was at, what the old primary's last known LSN was (which bounds the data-loss window), and what triggered the decision. This is what turns "the system did something" into something a human can verify after the fact, and it is what a chaos-testing exercise (deliberately triggering failures in a controlled way, including in production during low-traffic windows, specifically to verify this pipeline works under realistic conditions rather than only in a lab) should be validating end-to-end, including checking that dependent services handling the failover (connection pools reconnecting, in-flight requests retried safely) do not cascade the failure outward.
Manual versus automatic failover: what data can be lost, and the RPO/RTO difference
| Automatic failover | Manual failover | |
|---|---|---|
| Recovery time objective (RTO) | Seconds to low tens of seconds: bounded by health-check interval, DCS quorum round trip, and fencing confirmation time | Minutes: bounded by human detection, human judgment time, and manual execution of the same steps |
| Recovery point objective (RPO) | Determined mechanically by whichever replica was most caught-up at decision time, which may not be the theoretical best choice if the promotion-safety threshold accepted a replica with some lag | Potentially better: a human can inspect multiple replicas' exact positions, attempt to recover unshipped writes from the old primary's log if it is still reachable in a degraded state, and choose the option with the least data loss, at the cost of the extra time that inspection takes |
| Failure mode if wrong | Fast, but a bug in the automation (an incorrect promotion-safety check, a fencing race condition) can promote unsafely at machine speed, before a human notices | Slow, but a human in the loop can catch an obviously wrong situation (e.g., "the primary is actually fine, this is a monitoring false positive") before acting |
| What can be lost | Any write acknowledged by the old primary but not yet replicated to the promoted replica at the moment of the fence; bounded by the promotion-safety threshold | The same category of loss, but potentially smaller, since a human can wait slightly longer to let a lagging replica catch up, or attempt log recovery from the old primary, trading RTO for a better RPO |
Automatic failover is the right default for a well-tested, generic primary-down scenario, because RTO matters more than shaving the last few seconds of RPO for most services. Manual (or manual-confirmed) failover is worth keeping as an option, or even the default, for a scenario where the failure mode is ambiguous (a regional network partition isolating a meaningful fraction of followers, rather than a clean single-node crash), where an automated system cannot reliably distinguish "promote now" from "wait, this might resolve itself" and a wrong automatic decision (over-eager promotion during a transient blip) is worse than the extra minute a human takes to confirm.
Trade-offs and pitfalls
The most common design flaw is building the health check and the fencing mechanism as two independently-tuned systems that can disagree: a health check that decides to promote faster than the fencing mechanism can guarantee the old primary is actually stopped creates exactly the split-brain window the whole design exists to prevent. The second common flaw is never testing the failover path against a real, partial-partition scenario (only testing against a clean node crash), since a partition that isolates a meaningful fraction, not all, of the followers behaves very differently (some followers see a different quorum than others) than a simple full-node failure, and is the scenario most likely to expose a subtle bug in the promotion-safety logic.
Explain database federation as an alternative to sharding: querying across multiple independently-owned databases without redistributing their data. Describe one scenario where federation is preferable to sharding, and one where federation introduces unacceptable complexity.
Sample Answer
Direct answer
Database federation queries across multiple databases that each remain independently owned, operated, and never redistributed, typically via a federation layer that fans a query out to the relevant source databases and stitches the results together. Sharding, by contrast, deliberately redistributes ONE logical dataset across multiple nodes that a single system controls end to end. The difference that matters in practice: sharding is something you design and own, choosing the shard key and rebalancing strategy yourself; federation is something you do TO data you often do not own or cannot move, because it already lives in separately operated systems (a different team's service, an acquired company's legacy database, a third-party system you only have read access to).
When federation is the right call
A company acquires another company whose customer data lives in its own separate, actively-used production database, still owned and operated by that team, with its own release cadence and its own operational constraints. Migrating that data into a shared, sharded schema is a large, risky project with no urgent deadline forcing it. Federation lets a unified "search across all customers" feature query both databases live and merge the results, without either team giving up ownership of their own system or taking on a risky migration before there is a clear business reason to. This is the right trade when the data's owners, schemas, and operational lifecycles are genuinely going to stay independent for the foreseeable future, and the query patterns across them are occasional cross-source lookups rather than the bulk of the workload.
When federation introduces unacceptable complexity
A query that needs to join large amounts of data across two of the federated sources (a report needing every acquired company's transactions joined against the main company's customer table, filtered and sorted across the combined set) turns into pulling large result sets out of both systems over the network and joining them in application code or a federation middleware layer, since neither underlying database can push down a join against data it does not have. This scales badly and gets slower exactly as both datasets grow, with no way to add an index that spans both sides of the join. At that point, either the query pattern needs to change (accept eventually-consistent denormalized copies, replicated via change-data-capture (CDC: a technique that watches a source database's transaction log for every row-level insert, update, and delete, then streams those changes out as events, instead of querying the source live for each one) into one place purpose-built for the cross-source query, rather than querying live), or the underlying decision to keep the data federated instead of migrated needs to be revisited, since federation was never designed to make heavy cross-source joins fast.
Trade-offs and pitfalls
- Federation trades query performance and consistency for organizational and operational independence. It is the right tool when that independence is the actual constraint (you cannot or should not move the data), not a substitute for sharding when you simply have not gotten around to a proper data architecture.
- A federation layer becomes a hidden single point of failure and a hidden performance bottleneck if every cross-source query has to flow through it, even though no single underlying database is overloaded.
- Schema drift across the federated sources is a real ongoing cost: each source can and will change its own schema independently, on its own schedule, and the federation layer has to keep working across those changes without any coordinated migration process to keep everything in sync.
Your distributed database cluster is showing growing WAL (write-ahead log) shipping and apply lag on replicas during peak traffic, and downstream consumers are seeing stale reads. Walk through a remediation plan that addresses network bandwidth, replica IO throughput, and replica apply configuration, and states what level of read staleness you would accept for which consumers. Include both short-term mitigations and longer-term architectural fixes.
Sample Answer
Direct answer
Before writing a remediation plan, triage which of the three plausible causes is actually responsible: network bandwidth between primary and replica, replica disk IO throughput, or a single-threaded (or otherwise under-parallelized) apply worker that cannot keep up with the primary's write-ahead log (WAL, the durable, ordered record of every change the primary commits) even though both the network and the disk have headroom to spare. Treating all three the same wastes the incident's first minutes on the wrong lever: throwing more network bandwidth or a faster disk at a single-threaded apply bottleneck does nothing, since neither was ever the constraint. Once the bottleneck is identified, apply a short-term mitigation scoped to that specific cause, then follow with a longer-term architectural fix, and set an explicit staleness budget per consumer rather than treating "zero lag" as the only acceptable target.
Triage: compare the same three numbers every time
Pull three measurements: the primary's current WAL generation rate, the actual available network throughput between primary and replica, and the replica's actual sustained apply throughput (not its raw disk throughput, its APPLY throughput, since apply is usually bottlenecked on per-transaction commit and fsync overhead (fsync: the system call that forces the operating system to physically flush written data to durable storage before a transaction counts as committed, which is comparatively slow per call), not on how fast the disk can stream bytes). Whichever of the three is closest to (or below) the WAL generation rate is the bottleneck; the other two are very likely fine and not worth touching first.
# Triage arithmetic: which of network, replica IO, or single-threaded apply is the
# bottleneck, given measured/assumed throughput numbers (stated explicitly, not
# fabricated benchmarks -- this is exactly the kind of number an operator would pull
# from real monitoring during the incident).
WAL_GENERATION_RATE_MBPS = 45.0 # measured: bytes of WAL the primary produces per second at peak
NETWORK_LINK_MBPS = 1000 # measured: primary -> replica link capacity, megabits/s
NETWORK_THROUGHPUT_MB_S = NETWORK_LINK_MBPS / 8 # megabits -> megabytes/s
REPLICA_DISK_SEQUENTIAL_MB_S = 400.0 # measured: raw sequential write throughput available
REPLICA_APPLY_MB_S = 20.0 # measured: ACTUAL apply throughput achieved, single apply worker
# (each transaction commits individually; apply is bottlenecked
# on fsync/commit overhead per transaction, not on raw disk speed)
print(f"WAL generation rate: {WAL_GENERATION_RATE_MBPS:.1f} MB/s")
print(f"network capacity: {NETWORK_THROUGHPUT_MB_S:.1f} MB/s "
f"({'NOT the bottleneck' if NETWORK_THROUGHPUT_MB_S > WAL_GENERATION_RATE_MBPS else 'BOTTLENECK'})")
print(f"replica raw disk throughput: {REPLICA_DISK_SEQUENTIAL_MB_S:.1f} MB/s "
f"({'NOT the bottleneck' if REPLICA_DISK_SEQUENTIAL_MB_S > WAL_GENERATION_RATE_MBPS else 'BOTTLENECK'})")
print(f"replica ACTUAL apply throughput (single worker): {REPLICA_APPLY_MB_S:.1f} MB/s "
f"({'BOTTLENECK' if REPLICA_APPLY_MB_S < WAL_GENERATION_RATE_MBPS else 'not the bottleneck'})")
print()
deficit_mb_s = WAL_GENERATION_RATE_MBPS - REPLICA_APPLY_MB_S
print(f"apply deficit: {deficit_mb_s:.1f} MB/s of WAL accumulating faster than the replica applies it")
for minutes in (5, 30, 60):
backlog_mb = deficit_mb_s * minutes * 60
print(f" after {minutes} min at this rate: {backlog_mb:,.0f} MB of unapplied WAL backlog")
# how many parallel apply workers would close the gap, holding per-worker throughput
# constant at the same per-transaction commit overhead observed on the single worker
import math
workers_needed = math.ceil(WAL_GENERATION_RATE_MBPS / REPLICA_APPLY_MB_S)
print()
print(f"parallel apply workers needed to match generation rate (naive, assumes linear scaling "
f"of independent-transaction throughput): {workers_needed}")
print("in practice provision headroom above this minimum, since linear scaling from 1 worker "
"rarely holds exactly once workers start contending for the same hot rows/pages")
Output:
WAL generation rate: 45.0 MB/s
network capacity: 125.0 MB/s (NOT the bottleneck)
replica raw disk throughput: 400.0 MB/s (NOT the bottleneck)
replica ACTUAL apply throughput (single worker): 20.0 MB/s (BOTTLENECK)
apply deficit: 25.0 MB/s of WAL accumulating faster than the replica applies it
after 5 min at this rate: 7,500 MB of unapplied WAL backlog
after 30 min at this rate: 45,000 MB of unapplied WAL backlog
after 60 min at this rate: 90,000 MB of unapplied WAL backlog
parallel apply workers needed to match generation rate (naive, assumes linear scaling of independent-transaction throughput): 3
in practice provision headroom above this minimum, since linear scaling from 1 worker rarely holds exactly once workers start contending for the same hot rows/pages
In this worked scenario, both the network (125 MB/s available) and the replica's raw disk (400 MB/s) have plenty of headroom above the 45 MB/s WAL generation rate. The actual apply throughput, 20 MB/s on a single apply worker, is the real bottleneck, and the deficit compounds: at 25 MB/s of accumulating backlog, an hour of sustained peak traffic without intervention produces about 90 GB of unapplied WAL, and downstream consumers reading from that replica see staleness growing roughly in step with it.
Short-term mitigations
- If the triage points at the apply worker (as in the worked example): increase apply parallelism if the engine supports it (parallel apply workers applying independent, non-conflicting transactions concurrently), which the naive linear-scaling estimate above suggests would need about 3 workers to match generation rate, with real headroom above that minimum since workers contending for the same hot rows will not scale perfectly linearly.
- If the triage points at network bandwidth: enable WAL compression on the replication stream if supported, and check for anything else competing for the same link (a backup job, a bulk export) that could be paused or rescheduled during the incident.
- If the triage points at replica disk IO: check for a competing workload on the same replica (an ad hoc analytical query, a backup snapshot in progress) stealing IO from the apply process, and move or pause it.
- Across all three causes: communicate the current staleness level to consumers immediately, so anyone with a tight freshness requirement can fail over to a different data source or degrade gracefully, rather than silently serving stale reads while the fix is in progress.
Longer-term architectural fixes
- Parallelize apply structurally, not just as an incident response: if the workload regularly produces peak WAL rates that a single apply worker cannot sustain, that is a standing capacity gap, not a one-time incident, and needs permanent parallel-apply configuration plus monitoring that pages before the backlog reaches a consumer-visible threshold, not after.
- Decouple read-heavy consumers from replication capacity entirely with change-data-capture (CDC) into a message bus. Instead of every downstream reader hitting a physical replica that also has to keep up with WAL apply, stream committed changes into a durable message bus (Kafka or an equivalent) once, and let read-heavy consumers subscribe to that stream or read from a materialized view (a query's result set computed once and stored as its own table-like object, then refreshed on a schedule or trigger, instead of being recomputed from scratch on every read) built from it. This means a spike in downstream read demand no longer adds any load to the replication path at all, and it decouples "how fast can the replica apply WAL" from "how many consumers need this data," which is exactly the two problems this incident conflated.
- Right-size replica capacity to the actual peak WAL generation rate, not the average, since a replica sized for average load is exactly the kind that falls behind during the peaks that matter most.
Setting a staleness budget per consumer
Not every consumer needs the same freshness. A fraud-detection system reading recent transactions might need sub-second staleness and should read from the primary or a tightly-bounded replica, not a general-purpose read replica at all. A daily analytics rollup can tolerate minutes of staleness without anyone noticing. Write these budgets down explicitly per consumer rather than implicitly assuming "replica lag should always be near zero," since that assumption is what turns an acceptable, monitored staleness level into a false-alarm incident, and what leaves an actually-too-stale consumer with no defined threshold to alert on in the first place.
Trade-offs and pitfalls
- Skipping triage and applying all three mitigations at once wastes the most valuable early minutes of an incident and makes it harder to know afterward which change actually fixed it, information you need for the longer-term capacity plan.
- Parallel apply is not free: transactions that touch the same rows still have to serialize (each one waits for a lock, the mechanism a database uses to stop two transactions from modifying the same row at the same time, to be released) to preserve correctness, so parallelizing apply helps most when the workload has many independent transactions, and helps little for a workload dominated by lock contention on a small hot set of rows.
- A CDC-to-message-bus architecture adds its own new failure modes (consumer lag on the bus itself, schema evolution of the event stream) in exchange for removing the original coupling; it is a real architectural change, not a free upgrade, and needs its own monitoring once in place.
- A staleness budget only helps if consumers actually know about it and build against it; publishing the number without making it visible in the API or dashboard consumers actually look at just moves the assumption problem one level down instead of solving it.
How would you detect and mitigate silent data corruption or a split-brain scenario in a replicated database? Propose detection mechanisms, automated mitigation steps, and offline repair procedures that preserve data correctness.
Sample Answer
Direct answer
Split-brain (two nodes both believing they are the authoritative primary at the same time, typically after a network partition) and silent data corruption both stem from the same root cause: something wrote data that later turns out to conflict with, or be inconsistent with, the rest of the system, and nobody noticed at write time. Detection has to be active (continuously comparing state) rather than passive (waiting for a user to report wrong data), automated mitigation has to stop the bleeding fast (fence the wrong writer, redirect traffic), and offline repair has to reconcile divergent data without guessing, using whichever side has a verifiable, complete record of what actually happened.
Detection mechanisms
sequenceDiagram
participant N1 as Node A (was primary)
participant N2 as Node B (promoted primary)
participant Detector as Reconciliation job
Note over N1,N2: Network partition isolates Node A
N2->>N2: Promoted after fencing timeout, begins accepting writes
Note over N1: Node A incorrectly believes it is still primary (fencing failed to reach it in time)
N1->>N1: Also accepts writes during the partition window
Note over N1,N2: Partition heals
Detector->>N1: Compare checksums / row versions
Detector->>N2: Compare checksums / row versions
Detector->>Detector: Detect divergent rows written on both sides
Detector->>Detector: Quarantine divergent rows for offline reconciliation
- Consensus-layer detection (best, catches it before it happens). A properly implemented leader-election system (using a distributed configuration store like etcd or ZooKeeper, requiring a majority quorum to grant and renew a leader lease) makes split-brain structurally hard, not just detected after the fact: a node cannot legitimately believe it holds the leader lease unless it can prove it against a majority. This is a design choice made before the incident, not a detection technique applied during one, but it is the most effective mitigation available.
- Fencing-token/generation-number checks. Every write carries a monotonically increasing "generation" or "epoch" number tied to the current leadership term. Any downstream consumer (storage layer, replicas, or an external system) rejects a write carrying an older generation number than one it has already seen, which catches a zombie old-primary's writes even if the primary itself has not yet realized it lost leadership.
- Continuous checksum/version reconciliation. A background job periodically compares row-level checksums or version numbers between what should be identical copies (a shard and its replica, or two nodes that briefly diverged), flagging mismatches for investigation rather than waiting for a user-visible symptom.
- Application-level invariant checks. Domain-specific sanity checks (an account balance that went negative when the business rule says it never should, a row updated with a timestamp older than its own last-modified timestamp) catch corruption that passes structural checks (the row is syntactically valid) but violates business meaning.
Automated mitigation steps
The moment split-brain is suspected (two nodes both claiming leadership, or a generation-number mismatch detected), the automated response should: immediately fence the node with the lower/older generation number (revoke its ability to accept writes, via a network-level block, a credential revocation, or a hardware watchdog reset), stop routing any traffic to it, and freeze (do not auto-merge) any data written during the suspected divergence window rather than guessing which side is "right." Auto-merging without a clear, safe resolution rule is how a detection system turns one incident into a second, worse one.
Offline repair, preserving correctness
Once the systems are stable, reconciliation is a deliberate, evidence-based process, not an automated guess: identify the exact window of potential divergence using the generation-number or timestamp evidence, extract the rows written on each side during that window, and resolve each conflict using domain rules the team defines in advance (for genuinely append-only or idempotent operations, keep both if they do not actually conflict; for a true same-key conflicting write, prefer whichever side has independently verifiable evidence of being the legitimate majority-quorum leader during that window, and treat the other side's writes as candidates for compensating transactions rather than silent overwrites). Anything that cannot be confidently reconciled gets flagged for manual review rather than resolved by a rule that might be wrong, and the incident is not closed until every flagged row has an explicit resolution, not just "most of them look fine."
Trade-offs and pitfalls
The most dangerous mistake is auto-resolving conflicts with a blanket rule like last-writer-wins by wall-clock timestamp, without checking whether clocks were actually synchronized across the two nodes during the incident; clock skew during exactly the kind of network trouble that causes split-brain is common, which means the "last" writer by timestamp may not be the actually-later write. The second mistake is treating detection as a one-time incident-response activity rather than continuous background reconciliation; silent corruption that never triggers an obvious symptom (a slowly-diverging count, not a crash) can persist for a long time before anyone notices without an active, scheduled check.
A single PostgreSQL table receives millions of inserts per minute and experiences write hotspots on a monotonically increasing primary key (serial). This causes contention and slows ingestion. As a data engineer, propose architectural and database-level changes to eliminate hotspots and increase throughput while preserving insert order where necessary. Consider alternatives like partitioning, UUIDs, batching, or using a distributed write service; explain pros and cons.
Sample Answer
Direct answer
The root cause is that a serial/auto-incrementing primary key always assigns the next value at the "end" of the key space, so every insert contends for the same rightmost page of the primary key's B-tree index (the balanced tree structure the database uses to keep the index sorted and quickly searchable) and, on some engines, the same physical page on disk. I would eliminate the hotspot by decoupling the key's value from insert order: switch to a key that spreads inserts across the keyspace (a UUID variant, or a prefixed/hashed id), while preserving whatever ordering guarantee is actually needed (usually created_at, not the primary key itself) as a separate indexed column.
Structured elaboration
Why a serial PK specifically causes contention. Every insert has to acquire the same rightmost leaf page of the primary key's B-tree to add its new, always-higher value. Under high concurrent insert rates, many sessions are contending for a lock on that same page (and the same underlying disk block) simultaneously, which serializes work that should otherwise be parallel across many independent index pages.
The realistic options, and their trade-offs:
- Switch to a random or time-ordered UUID. A fully random UUID (v4) spreads inserts uniformly across the keyspace, eliminating the rightmost-page contention entirely, but destroys any locality: rows inserted around the same time end up scattered across the whole index and table, which hurts range scans over recent data and bloats the index's working set (random insert points mean more index pages touched per insert, and worse cache locality for "recent rows" queries). A time-ordered UUID variant (one that encodes a timestamp in its high bits, sorting close to insertion order while still having enough randomness in the low bits to avoid a single hot page) keeps most of the write-distribution benefit while preserving much of the locality a fully random UUID gives up; this is usually the better default unless there is a specific reason to fully randomize.
- Partition the table (e.g., by a hash of a different, non-monotonic column, or by time range) so inserts spread across several partitions' indexes instead of one table's single index. This does not change the primary key itself, but each partition gets its own B-tree, and if the partition key is chosen so writes distribute across partitions (not, again, monotonically), the rightmost-page contention is divided N ways.
- Batching: instead of many small, independent inserts each contending for the tail of the index, buffer incoming events for a short window and insert them as fewer, larger batches. This reduces the number of separate lock acquisitions per unit of throughput (one batch insert instead of many single-row inserts), though the underlying keyspace contention pattern for a serial key is unchanged if the key is still monotonic; batching is a mitigation for insert overhead generally, not a fix for the key-choice root cause.
- A distributed write service / multiple insert paths, each responsible for a distinct sub-range or shard, so no single serial sequence is the sole source of new keys. This is the most architecturally involved option and usually only justified once a single instance's insert throughput, not just this specific contention pattern, is the actual ceiling.
Preserving insert order where it's actually needed. The common mistake is assuming the primary key itself has to carry ordering information. In practice, "preserve order" almost always means "be able to query recent rows in order" or "know which row came first for tie-breaking," both of which are satisfied by keeping (or adding) an indexed created_at (or a monotonic sequence number used only for tie-breaking, not as the primary storage key) column, while the primary key itself is free to be chosen purely for write distribution.
Worked example
The contention mechanism is a property of any B-tree-backed monotonic key, but its symptom, index bloat and page-split pressure concentrated at one end of the tree, is directly observable. Using the same instrumentation as a table already partitioned to avoid this (from the partitioning discussion above): a BIGSERIAL primary key on a single, unpartitioned table means every one of, say, 10,000 concurrent inserts targets the same logical position in the index; splitting the same insert stream across 4 hash-partitioned sub-tables (by a non-monotonic column, e.g. a hash of the row's event_type plus a random shard suffix) means each partition's own B-tree only sees roughly 2,500 of those inserts, each targeting that partition's own rightmost page, a quarter of the original contention on any single page, with the same total throughput.
total_inserts_per_window = 10_000
n_shards = 4
print(f"unpartitioned serial PK: {total_inserts_per_window} inserts contend for one page")
print(f"sharded by a non-monotonic key across {n_shards}: "
f"~{total_inserts_per_window // n_shards} inserts/shard contend for that shard's own rightmost page")
unpartitioned serial PK: 10000 inserts contend for one page
sharded by a non-monotonic key across 4: ~2500 inserts/shard contend for that shard's own rightmost page
Trade-offs & pitfalls
- Switching to a fully random UUID without checking the downstream cost. It solves the write hotspot completely but can turn what used to be a cheap "recent rows" range scan (because recent rows were physically adjacent under the serial key) into a scan that touches pages scattered across the whole table; measure that cost against the write-side win before committing to it.
- Batching alone is not a fix if the underlying key stays monotonic. It reduces the number of separate insert operations, which helps overall throughput, but does not change the fact that a monotonic key still concentrates each batch's inserts at the same end of the index; it is a complement to a key-choice fix, not a substitute for one.
- Losing sight of what "preserve insert order where necessary" actually requires. The natural instinct is to keep the primary key monotonic because "we need order somewhere," but conflating "the storage key" with "the field the application queries order by" is exactly what causes this problem in the first place; keep them as two separate columns with two separate jobs.
Unlock Full Question Bank
Get access to all Replication, Partitioning, and Sharding interview questions and detailed answers.
Sign in to ContinueJoin thousands of developers preparing for their dream job.