skip to content

A high-throughput ingest table is clustered on a monotonically increasing key, and insert throughput plateaus well below what the hardware should allow even though there is no I/O bottleneck. How would you reason about the cause and the options?

level: principalimportance: nice to knowfreq 28%

answer

  1. Idle disks + flat throughput = serialization, not I/O
  2. Monotonic clustering key → one hot right-edge page
  3. Fix = k insert points: leading discriminator or partitioning
  4. Read cost: one range becomes k ranges — keep k small
  5. Random keys trade contention for bloat and cache misses

basics

~20 s

Every insert targets the same rightmost leaf page, so concurrent writers serialize on that one page's latch. Options: partition or shard the key space so there are several insert points, add a leading discriminator to spread writes, or accept it and scale out by table or node.

solid answer

~60 s

With a monotonically increasing clustering key, all inserts belong at the right edge of the tree, so every concurrent writer contends for the **same leaf page** (and often its parent during splits). That is a latch/contention ceiling, not an I/O one — the symptom is CPU spinning or waits on a single page while disks are idle. The options trade one problem for another: 1. **Multiple insert points.** Prefix the clustering key with a small discriminator — a hash bucket, shard id, or a natural grouping column like tenant — so writes land on k right edges instead of one. Cost: the key widens, and reads that scanned one contiguous range now scan k ranges. 2. **Partition the table** by that discriminator, giving each partition its own tree and right edge. Cleaner isolation, more objects to manage. 3. **Spread to separate tables or nodes** and merge at read time. 4. **Do nothing.** Random keys "fix" contention but reintroduce scattered writes, poor fill and cache misses — usually a worse trade. Decide by measuring: confirm the wait is on a page, then choose the discriminator that costs the least on the dominant read pattern.

code

text · 5 lines
text
clustering key = (id)                 -> all writers -> leaf[rightmost]   (1 hot page)
clustering key = (bucket, id)         -> writers     -> k right edges
clustering key = (tenant_id, id)      -> writers     -> per-tenant edges + better read locality

read cost: 1 contiguous range  ->  k ranges merged

go deeper

for a junior

Recognize that all inserts landing at the end of an ordered structure means writers queue for the same page.

for a middle

Explain the mechanism and name spreading the insert point via a leading discriminator, with the read-side cost.

for a senior

Insist on confirming the wait is page contention, weigh composite key versus partitioning, and try batching and index pruning first.

for a principal

Present it as a tradeoff space with no single answer — contention versus locality versus read complexity — size the discriminator to measured concurrency, and state when sharding out is the right end state.

## Recognize the shape of the problem A plateau with idle storage points at serialization, not throughput limits. In a clustered table with an ever-increasing key, the structural suspect is obvious: every insert must go where the key belongs, and the key always belongs at the far right of the tree. So all writers converge on: - the **current rightmost leaf page**, which each must latch exclusively to modify; - during a split, its **parent** and possibly higher levels; - and often the same few log/metadata structures. This is the mirror image of the random-key problem. Random keys scatter writes and destroy locality; monotonic keys concentrate writes and destroy parallelism. Both are consequences of the same fact — in clustered storage, key order *is* physical placement. ## Confirm before you redesign The reasoning only earns its keep if the diagnosis is verified. You want evidence that the wait is on a specific page or index rather than storage or log flush: high concurrency correlating with rising wait time on a single object, throughput that does not improve with more writer threads (or degrades), idle device queues, and a per-thread profile dominated by waiting rather than working. If instead the wait is on durability (log flush per commit), the fix is batching or group commit, and none of the key redesign below helps. Also check the trivial explanations first: an insert path doing one round trip per row, a trigger, an excessive number of secondary indexes each needing maintenance, or a lock on a uniqueness check against another table. ## Option 1 — introduce several insert points Make the clustering key lead with a low-cardinality discriminator so there are k right edges: - **A natural grouping column** — tenant id, region, device group — if one exists and the read patterns already filter by it. This is the best case: contention drops and read locality *improves*, because a tenant's rows become contiguous. - **A synthetic bucket** — for example a hash of some column modulo k, or a per-writer id. Purely mechanical: it buys parallelism but adds nothing for reads. The cost is real. A query that used to scan one contiguous range now scans k ranges and merges them, so scan cost and plan complexity rise with k. Choose k to be just large enough — a small multiple of the number of concurrent writers, not an arbitrary large number — because every increment permanently taxes the read side. ## Option 2 — partition the table Partitioning on the same discriminator gives each partition its own B+tree with its own right edge. The contention isolation is cleaner than a composite key (separate structures rather than separate subtrees), and partitions bring independent benefits: bulk removal of old data by dropping a partition, smaller per-partition indexes, and parallel maintenance. The costs are operational — more objects, and queries that do not filter on the partitioning column must touch every partition. If the ingest table is time-series-shaped, partitioning by time is often already justified for retention reasons; note though that time partitioning alone does *not* fix right-edge contention, since all current writes still land in the newest partition. You need a discriminator orthogonal to time for that. ## Option 3 — spread across tables or nodes At the extreme, give each writer or shard its own table or its own database node, and merge at read time (a view, a union, or application-side fan-out). This scales furthest but pushes complexity into every reader and into operations. It is a reasonable end state for genuinely enormous ingest rates, and an over-reaction for a single hot page. ## Option 4 — reconsider the key entirely The naive fix is to make the key random so writes scatter. Resist it. That converts a contention problem — bounded, measurable, and fixable with a small discriminator — into an I/O and footprint problem: read-before-write on uncached leaves, mid-page splits with poor fill, table bloat, and a working set equal to the whole index. Random keys are usually a worse trade unless the table is small enough to stay entirely in memory. A time-ordered identifier is not a fix either, since it is monotonic by design and lands in the same place. ## Also consider the non-structural levers - **Batch inserts.** Fewer, larger statements reduce the number of times the hot page is latched per row and amortize log flushes. Often the cheapest large win. - **Reduce per-row index maintenance.** Each secondary index is additional work per insert; dropping ones that no longer earn their keep raises ingest throughput directly. - **Deferred or bulk load paths.** Staging into an unindexed or heap-structured landing table and merging in batches decouples ingest rate from the clustered structure entirely — at the cost of read latency for the newest data. ## How to present the answer There is no single right answer here, and the interviewer knows it. What they are listening for is: (1) you identify concentration at the right edge as the mechanism; (2) you insist on confirming it before redesigning; (3) you propose spreading insert points and can state the read-side cost of doing so; (4) you explicitly reject the random-key "fix" and can say why; (5) you mention cheap non-structural levers like batching before schema surgery; and (6) you size the fix to the measured concurrency rather than reaching for the largest hammer.

  • Why not simply switch to random keys to eliminate the hot page?
    Because it replaces a bounded contention problem with an unbounded I/O one. Random insertion scatters writes across the whole tree, so once the index exceeds memory each insert becomes a random read before the write, pages split mid-way leaving poor fill, the table bloats, and the entire index becomes the working set. A small leading discriminator gives you the parallelism you actually need while keeping inserts concentrated in a few hot, cached regions.
  • How would you choose the number of buckets in a bucketed clustering key?
    Size it to the measured write concurrency, not to a round number — a small multiple of the number of simultaneous writers is usually enough to move the bottleneck off the page latch. Every additional bucket permanently multiplies the number of ranges a scan must merge, so the read side pays forever for parallelism the write side may not need. Measure, pick the smallest value that removes the wait, and revisit only if concurrency grows substantially.
  • Does partitioning a time-series ingest table by time solve right-edge contention?
    No. All current inserts still carry recent timestamps and therefore land in the newest partition, whose tree has exactly one right edge — the hot page simply moves with time. Time partitioning is still worth doing for retention and maintenance, since dropping an old partition is far cheaper than deleting rows. To spread inserts you need a discriminator orthogonal to time, such as a tenant or hash bucket, either in the key or as a sub-partitioning dimension.
  • What cheaper measures would you try before changing the key or partitioning scheme?
    Batch the inserts so each statement covers many rows, which amortizes both the page latching and the commit-time log flush. Audit the secondary indexes, since every one of them adds per-row maintenance work on the hot path. And confirm the wait really is on the index page rather than on log flush or a foreign-key check, because those have entirely different remedies. Schema surgery should follow measurement, not precede it.

One checkout lane at a supermarket: the queue is not limited by how fast items scan but by there being a single lane — you open a few more, not scatter shoppers randomly across the store.

saying these in an interview costs you the question

  • Jumping to random keys as the fix for right-edge contention
  • Diagnosing contention without evidence that the wait is on a page rather than on log flush or a lock
  • Adding a large number of hash buckets without accounting for the permanent read-side scan cost
  • Claiming time-based partitioning removes right-edge contention
  • Ignoring batching and index count before proposing a schema redesign

context