Skip to content
BytePatterns

DynamoDB Hot Partitions: Why Writes Throttle and How to Shard

9 min readBytePatterns

Why one busy partition key throttles a DynamoDB table with spare capacity, how per-partition limits work, and how random or calculated suffixes spread the load.

A DynamoDB table can have plenty of capacity and still throttle writes. The cause is almost never the table; it is one partition key taking most of the traffic. Understanding why takes two facts about how DynamoDB places data, and fixing it takes one design pattern, write sharding. Every behaviour and number below comes from the DynamoDB Developer Guide pages listed at the end, as of September 2026; limits change, so check them before you rely on them.

The problem it solves

Picture a flash sale. Every order write uses the partition key sale#26, because "all orders for this sale" is the obvious query. Traffic spikes, the table's total capacity is far from used, and yet writes start failing with throttling errors. Adding capacity to the table does not help, because the bottleneck is not the table.

The intuition

Fact one: the partition key picks the partition. DynamoDB passes the partition key value through an internal hash function, and the output decides which physical partition stores the item. Items with the same partition key always land in the same partition. With a composite key, items that share a partition key are stored together, ordered by the sort key, which is what makes "this customer's orders, newest first" a single Query.

Fact two: each partition has its own ceiling. Every partition is designed to deliver at most 3,000 read units and 1,000 write units per second. One read unit is one strongly consistent read per second, or two eventually consistent reads, for an item up to 4 KB; one write unit is one write per second for an item up to 1 KB. Larger items cost proportionally more units.

Put the two together and the flash sale explains itself. All writes for sale#26 hash to one partition, so they all compete for that partition's 1,000 write units per second, no matter how much capacity the rest of the table has. The documentation's advice is to design for uniform activity across all partition keys.

Item size matters on the read side too. The guide's own example: with 20 KB items, one strongly consistent read costs 5 read units, so a single item can serve at most 600 such reads per second before its partition hits the limit.

The fix is to widen the key space. Write sharding appends a suffix to the partition key, so one logical key becomes many physical ones:

  • Random suffix. Pick a number from 1 to N for each write. Writes spread evenly, but to read one specific item you do not know which suffix it got.
  • Calculated suffix. Derive the number from something you query by, such as the order ID. Writes still spread, and a single item can be fetched with GetItem, because its suffix can be recomputed.

Either way, reading everything for the logical key now means one Query per suffix, merged in your application.

Watch it run

The animation starts with three customers whose keys hash to different partitions, and shows a sort key turning "c#42's orders, newest first" into one Query. Then every flash-sale write uses sale#26: one partition turns hot and throttles while the other two sit idle. Write sharding appends #1, #2, #3 to the key, and the writes fan out across partitions. The last frames show the price: reading the sale back is one query per suffix, merged in your code.

RDS vs DynamoDB

Step 1 of 10

Every DynamoDB item has a partition key. Three customers are about to write orders.

The same interactive animation as the lesson — step through it with the controls.

The code

A toy model, not DynamoDB: the hash below stands in for the real, internal one, and the table has four partitions. The capacity arithmetic follows the documented unit definitions:

import hashlib
import math
from collections import Counter

PARTITION_WCU = 1000          # write units per second, per partition (docs)
PARTITION_RCU = 3000          # read units per second, per partition (docs)

def partition_of(key, partitions=4):
    """Toy placement: DynamoDB's real hash function is internal."""
    return int(hashlib.md5(key.encode()).hexdigest(), 16) % partitions

def write_units(item_kb):
    return math.ceil(item_kb / 1)                  # 1 unit per 1 KB written

def read_units(item_kb, strongly_consistent=True):
    units = math.ceil(item_kb / 4)                 # 1 unit per 4 KB read
    return units if strongly_consistent else units / 2

print(write_units(0.5), write_units(2.5))          # 1 3
print(read_units(20), read_units(20, False))       # 5 2.5
print(PARTITION_RCU // read_units(20))             # 600

The flash sale at 3,000 writes per second of 1 KB items. On one key, 2,000 writes a second are over the partition's limit. Split over ten suffixes, the toy hash places them unevenly, 4 of the 10 on one partition, so a little throttling remains:

def throttled(writes_per_key, item_kb=1):
    """writes_per_key: key -> writes per second. Returns throttled writes/s."""
    load = Counter()
    for key, rate in writes_per_key.items():
        load[partition_of(key)] += rate * write_units(item_kb)
    return sum(max(0, units - PARTITION_WCU) for units in load.values())

hot = {"sale#26": 3000}                            # every write, one key
print(throttled(hot))                              # 2000

shards = 10
spread = {f"sale#26#{i}": 3000 / shards for i in range(1, shards + 1)}
print(sorted(Counter(partition_of(k) for k in spread).values()))
print(throttled(spread))
# [2, 2, 2, 4]
# 200.0

A calculated suffix. The same order always maps to the same shard, so one item can be read directly, while reading the whole sale visits every suffix:

SHARDS = 10

def sharded_key(sale_id, order_id):
    """Calculated suffix: the same order always lands on the same shard."""
    digest = hashlib.sha256(order_id.encode()).digest()
    return f"{sale_id}#{digest[0] % SHARDS + 1}"

table = {}                                         # partition key -> items
for n in range(25):
    order = f"o-{n}"
    table.setdefault(sharded_key("sale#26", order), []).append(order)

print(sharded_key("sale#26", "o-7") in table and "o-7" in table[sharded_key("sale#26", "o-7")])
everything = [o for i in range(1, SHARDS + 1) for o in table.get(f"sale#26#{i}", [])]
print(len(everything), len(table))                 # one Query per suffix, merged
# True
# 25 10

The model's invariants on 2,000 random cases: unit counts match a brute-force search for the smallest sufficient number of units, sharding a key never throttles more than leaving it whole, a calculated suffix is repeatable, and querying every suffix finds every item:

import random

random.seed(26)
ok = True
for _ in range(2000):
    kb = random.uniform(0.1, 50)
    ok &= read_units(kb) == min(u for u in range(1, 100) if 4 * u >= kb)
    ok &= write_units(kb) == min(u for u in range(1, 100) if u >= kb)

    total = random.randint(0, 6000)
    n = random.randint(1, 30)
    single = throttled({"k": total})
    sharded = throttled({f"k#{i}": total / n for i in range(1, n + 1)})
    ok &= single == max(0, total - PARTITION_WCU)
    ok &= sharded <= single + 1e-9                 # sharding never makes it worse

    orders = {f"o-{random.randint(0, 10**6)}" for _ in range(random.randint(0, 40))}
    t = {}
    for o in orders:
        key = sharded_key("s", o)
        ok &= key == sharded_key("s", o)           # calculated: repeatable
        t.setdefault(key, []).append(o)
    merged = [o for i in range(1, SHARDS + 1) for o in t.get(f"s#{i}", [])]
    ok &= sorted(merged) == sorted(orders)         # read-back finds everything
print(ok)                                          # True

The complexity

The trade is write throughput against read cost:

  • Writes to one logical key can use up to N partitions' worth of write capacity instead of one, if the suffixes hash to different partitions.
  • Reading one item stays a single GetItem with a calculated suffix; with a random suffix you may have to try every shard.
  • Reading the whole key becomes N queries plus a merge. The guide's example uses dates with suffixes 1 to 200, which means 200 queries to read one day.

Pick N from the peak write rate and the per-partition limit, with headroom, because hashing does not spread a small number of keys perfectly evenly, as the toy run shows.

Where it goes wrong

  • Adding table capacity to fix a hot key. The table's total was never the limit; the partition's was.
  • A low-cardinality partition key. Status values, today's date or a single tenant ID concentrate traffic by design.
  • Random suffixes when you need point reads. You lose the ability to fetch one item without trying every shard; use a calculated suffix instead.

How to say it in an interview

"DynamoDB hashes the partition key to pick a partition, and each partition tops out at 3,000 read and 1,000 write units a second. So a single popular key can throttle while the table has spare capacity. I'd design for uniform access across keys, and for a key that must be hot, like a flash sale, I'd shard it by appending a suffix. A suffix calculated from the order ID keeps single-item reads as GetItem; reading the whole sale is one query per suffix, merged in the application. If access patterns are unknown or need ad-hoc queries, that's a point for a relational database instead."

The wider comparison, including Multi-AZ and read replicas, is in the RDS vs DynamoDB lesson, and the same "hash a key to pick a node" idea appears in design a distributed cache.

Sources