skip to content

A team maintains an index table that maps order status — one of PENDING, SHIPPED, DELIVERED, or CANCELLED — to order IDs, so they can quickly list all orders in a given status. Under heavy load, writes to this index table start throttling even though the base orders table stays healthy. What's the likely root cause, and how would you redesign the index table's key to fix it?

level: seniorimportance: should knowfreq 45%

answer

  1. low-cardinality partition key -> hot partition
  2. per-partition throughput ceiling regardless of table-level provisioned capacity
  3. write sharding / salting: append hash(id) mod N to the key
  4. read fan-out cost: N parallel queries + merge
  5. shard only the hot values, not uniformly

basics

~20 s

Using status as the partition key means all orders with the same status pile into one 'bucket,' so writes for that status compete for the same limited capacity even though the rest of the system is fine. Spreading the writes across more buckets — for example by adding some kind of extra prefix or suffix — fixes the bottleneck.

solid answer

~50 s

The root cause is a hot partition: with only four possible values for the partition key (PENDING, SHIPPED, DELIVERED, CANCELLED), every write for a given status lands on the same physical partition, and most stores cap the throughput a single partition can sustain regardless of how much overall capacity is provisioned for the table. Since most orders sit in PENDING or SHIPPED at any given time, those two partitions absorb almost all the write traffic while the other two sit idle, and the table as a whole throttles well below its aggregate provisioned capacity. The fix is to widen the effective key space: shard the partition key by appending a calculated suffix (e.g., status + a hash of the order ID mod N, giving PENDING-0 through PENDING-N), so writes for the same status spread across N physical partitions instead of one, and reads for 'all PENDING orders' fan out N queries and merge the results instead of one query.

go deeper

for a junior

Should recognize in plain terms that putting everything under a few labels creates a bottleneck, without needing partition-ceiling numbers.

for a middle

Should identify low cardinality of the partition key as the cause and know 'spread the writes out somehow' as the general direction of the fix.

for a senior

Should name the write-sharding/salting technique concretely (hash/random suffix, modulo N) and articulate the read-side fan-out trade-off it introduces.

for a principal

Should reason about sizing N from observed load and per-partition limits, and know to shard selectively (only the hot values) rather than uniformly, weighing read complexity against write-throughput headroom.

## Why the key choice bites This is a hot-partition problem, and it's a direct consequence of how index tables are usually keyed: the partition key is the indexed attribute's value, chosen because it's exactly what a query needs to look up by. That choice is correct for query simplicity but can be badly wrong for write distribution when the indexed attribute has low cardinality — few distinct values — because most horizontally-scaled key-value stores (DynamoDB, Azure Table Storage, Cassandra, Cosmos DB): - achieve their aggregate throughput by spreading data and traffic across many physical partitions keyed by hash ranges of the partition key; - but every single partition key value always maps to exactly one physical partition (or a small, fixed replica set), and that one partition has a hard throughput ceiling regardless of how much capacity is provisioned for the table overall. DynamoDB, for example, has historically capped a single partition around 1,000 write capacity units per second (and around 3,000 RCU for reads) no matter how many thousands of WCU you've provisioned in aggregate for the table. ## What four values do to the table In this scenario, the indexed attribute — order status — has exactly four distinct values. That means the index table has, at most, four physical partitions doing any work at all, no matter how the store's autoscaling or partition-splitting logic tries to help; you cannot split a partition further than one distinct key value can go. And the traffic isn't even evenly spread across those four: - a healthy e-commerce system has most of its live orders sitting in `PENDING` (just created, being fulfilled) or transitioning through `SHIPPED`; - while `DELIVERED` and `CANCELLED` are comparatively rare events per unit time — an order becomes `DELIVERED` once and then typically isn't written to again under that key, whereas `PENDING` receives a write for every single new order created. So the write load concentrates overwhelmingly on one or two of the four partitions, and those specific partitions throttle — returning capacity-exceeded errors — while the table's aggregate provisioned capacity, spread in theory across four partitions, looks nowhere near exhausted in the account-level or table-level metrics, which is exactly the confusing symptom described: the base table (keyed by the high-cardinality order ID) is healthy because every order gets its own partition-ish slice, but this index table (keyed by four-valued status) is starving. ## The fix — widen the key space The fix is to break the natural one-partition-per-status structure by artificially widening the key space — a technique generally called write sharding or salting. Instead of using the raw status value as the partition key, compute a composite key that combines the status with some evenly-distributed extra bits, such as a hash of the order ID (or a random number, or a round-robin counter) modulo some shard count N — for example, `PENDING-0`, `PENDING-1`, ... `PENDING-7` for `N=8` shards. Each individual order's index-table write now lands on one of the N sub-partitions for its status, chosen by the hash, spreading what used to be one hot partition's worth of traffic across N physical partitions, each handling roughly 1/N of the load. The trade-off is on the read side: 'list all `PENDING` orders' can no longer be answered by a single partition query — it now requires fanning out N parallel queries (one per shard) and merging the results client-side, which adds complexity and some latency (bounded by the slowest of the N parallel reads) in exchange for removing the write bottleneck. ## Choosing N Choosing N is itself a judgment call: - too small and the hot-partition problem persists in diminished form; - too large and you pay unnecessary read fan-out cost and per-shard overhead for a table whose absolute load might not have needed that much spreading. Teams typically size N based on the peak observed write rate for the hottest status divided by the known single-partition throughput ceiling, with headroom for growth, and may choose to shard only the specific status values that are actually hot (leave `DELIVERED` and `CANCELLED` unsharded, shard `PENDING` and `SHIPPED` into, say, 8 sub-partitions each) rather than applying uniform sharding to a key space that doesn't need it everywhere. This exact class of fix — sharding a low-cardinality partition key with an appended hash or random suffix to avoid a hot partition — is explicitly documented in AWS's own DynamoDB best-practices guidance as the standard remedy for this failure mode.

  • Why doesn't provisioning more total capacity for the index table fix a hot partition?
    Because most stores enforce a per-partition throughput ceiling independent of the table's aggregate provisioned capacity — a single partition can only absorb so many writes per second no matter how much capacity sits unused on other partitions. Adding capacity increases the total pie, but the hot partition can't eat a bigger slice of it if all the traffic keeps landing on that one key.
  • What does the read path look like after sharding the status key into N sub-partitions?
    A query for 'all PENDING orders' becomes N separate queries, one against each PENDING-0 through PENDING-(N-1) sub-partition, run in parallel, with the results merged (and possibly re-sorted, if ordering matters) by the caller. This adds fan-out complexity and ties overall latency to the slowest of the N parallel reads, which is the cost paid for removing the write bottleneck.
  • Would sharding help if the problem were the reverse — a single order ID being written to extremely frequently?
    No — sharding by adding a hash suffix to a key only helps when many distinct entities share the same low-cardinality key value; if the hot spot is a single specific order ID being updated at very high frequency, that's a hot key/hot item problem on one logical entity, which sharding by hash doesn't fix since all writes for that one order still need to converge somewhere to stay correct.

It's like a supermarket with only four checkout lanes, one per first letter group of the customer's last name. If most customers happen to have last names starting with the same letter group, that one lane backs up into a huge line while the other three sit empty — even though the store has plenty of total checkout capacity, just badly distributed.

saying these in an interview costs you the question

  • Suggests only 'add more capacity/provisioned throughput' without addressing the per-partition ceiling
  • Doesn't recognize that low cardinality (few distinct values) in the partition key is the root cause
  • Proposes sharding but forgets to mention the read-side fan-out cost it introduces
  • Assumes the base table and index table should automatically have the same throughput characteristics
  • Can't explain why PENDING/SHIPPED would be hotter than DELIVERED/CANCELLED in this domain

context