Producers to a Kinesis Data Streams stream are getting ProvisionedThroughputExceededException even though the stream's total IncomingBytes sits far below shard count multiplied by 1 MB/s. What is happening, and how do you fix it?
answer
- limits bind per shard, not per stream
- turn on shard-level metrics first
- the key distribution is the real bug
- one key never exceeds one shard
basics
~20 sA hot shard. Kinesis enforces throughput per shard, and a low-cardinality or skewed partition key concentrates writes on one shard's hash-key range while the others idle. Fix the key distribution first; resharding alone only buys time.
solid answer
~50 sAggregate throughput is the wrong number to look at — the limits bind per shard. The partition key is MD5-hashed into a shard's hash-key range, so a skewed key sends most traffic to one shard, which throttles at 1 MB/s or 1,000 records/s while the rest of the stream idles. I would confirm it with `EnableEnhancedMonitoring` to get shard-level `IncomingBytes`, `IncomingRecords` and `WriteProvisionedThroughputExceeded`, and check `GetRecords.IteratorAgeMilliseconds` per shard — a single shard lagging is the same fingerprint. The durable fix is the key: raise cardinality, or salt the hot value with a bounded suffix and re-merge downstream, accepting that salting breaks strict per-original-key ordering. Resharding helps only if the hot range holds many distinct keys — `SplitShard` on that range, or `UpdateShardCount`. And the hard ceiling to state out loud: a single partition key can never exceed one shard's throughput, so no amount of resharding and no switch to on-demand rescues one genuinely hot key.
code
python · 27 linesimport random, time
import boto3
kinesis = boto3.client("kinesis")
# Salt a known-hot tenant across 32 hash values; leave others untouched.
HOT = {"acme"}
SALT_BUCKETS = 32
def partition_key(tenant_id: str) -> str:
if tenant_id in HOT:
return f"{tenant_id}#{random.randrange(SALT_BUCKETS)}"
return tenant_id
def put_with_retry(stream, records, attempts=5):
pending = records
for attempt in range(attempts):
resp = kinesis.put_records(StreamName=stream, Records=pending)
if resp["FailedRecordCount"] == 0:
return
pending = [
pending[i]
for i, r in enumerate(resp["Records"])
if r.get("ErrorCode")
]
time.sleep((2 ** attempt) * 0.1 + random.random() * 0.1)
raise RuntimeError(f"{len(pending)} records still failing")go deeper
Know that Kinesis throttling is per shard and that the partition key decides which shard a record lands on, so an uneven key means some shards are full while others are empty.
Explain how to prove it: enable shard-level metrics and compare IncomingBytes and WriteProvisionedThroughputExceeded across shards, and describe salting and resharding as the two families of fix.
Demonstrate the diagnosis-to-remediation path under production pressure — confirm skew with data, apply a safe short-term measure, and state explicitly what ordering guarantee the fix costs you.
Own the structural limit: one key cannot exceed one shard, so a genuinely hot entity is a data-modelling decision, not a capacity request. Be ready to argue for decomposing the entity or isolating it on its own stream.
## The symptom and the arithmetic `ProvisionedThroughputExceededException` on `PutRecord`, or `FailedRecordCount` climbing in `PutRecords` responses with the per-record `ErrorCode` set to the same exception, while the stream's aggregate `IncomingBytes` in CloudWatch looks comfortable. The contradiction is only apparent: `IncomingBytes` at stream level is a **sum**, and Kinesis enforces its limit on **each shard independently**. A 10-shard stream at 3 MB/s total is fine if the traffic is even and throttled solid if 2.5 MB/s of it lands on one shard. ## Why traffic concentrates Every record carries a partition key. Kinesis MD5-hashes that key to a 128-bit integer and stores the record in whichever shard owns the containing hash-key range. Two things concentrate traffic: - **Low cardinality.** A key like `event_type` or `region` has a handful of distinct values, so however many shards you add, only a handful of hash values exist and most shards receive nothing. - **Skewed cardinality.** The key has many values but the distribution is Zipfian — `tenant_id` where one enterprise customer is 80% of the volume, `device_id` where one buggy fleet retries in a loop. MD5 spreads *distinct* keys well. It does nothing about a key that is not distinct enough or not evenly used. ## Confirming it before you change anything Stream-level metrics cannot show skew, by construction. Turn on shard-level metrics: ```bash aws kinesis enable-enhanced-monitoring \ --stream-name events \ --shard-level-metrics IncomingBytes IncomingRecords \ WriteProvisionedThroughputExceeded \ IteratorAgeMilliseconds ``` Then compare `IncomingBytes` across `ShardId` dimensions. A hot shard is unmistakable — one series an order of magnitude above the rest. `WriteProvisionedThroughputExceeded` on that shard confirms the writes are being rejected there specifically. On the read side, `GetRecords.IteratorAgeMilliseconds` climbing for one shard while the others sit near zero is the same fingerprint seen from the consumer end. A second, cheaper check: log the partition key alongside throttled records in the producer and count the top keys. It takes minutes and usually names the culprit outright. ## The fixes, in the order I would consider them **1. Change the key.** If the key was chosen carelessly — `event_type` when nothing requires per-type ordering — replacing it with a high-cardinality key (`user_id`, `session_id`, or a random UUID when no ordering is required) removes the problem permanently and costs nothing at runtime. **2. Salt the hot value.** When the key is correct in principle but one value dominates, append a bounded random suffix: `acme#0` … `acme#31`. Writes now spread over up to 32 hash values. The price is explicit and must be stated: records for `acme` are no longer in a single shard, so **per-key ordering for that tenant is gone**. Salting is only acceptable when the consumer is order-insensitive, or when it can re-sequence downstream using an application-level timestamp or version. **3. Reshard.** `SplitShard` divides the hot shard's hash-key range at a chosen point; `UpdateShardCount` scales the whole stream uniformly. Both close the parent shard and open children, and consumers must drain the parent before the children — which the KCL handles for you, and which is exactly how per-key ordering survives a reshard. Resharding only helps if the hot range holds **many** distinct keys that a split can separate. Splitting a range whose traffic is one key just moves the problem to a child shard. **4. Switch to on-demand.** On-demand mode splits hot shards automatically, which absorbs organic growth and moderate skew without an operator in the loop. It does not repeal the physics: it reacts over minutes, and it cannot split a single key. **5. Reduce what you write.** If the hot key is a retry storm or a debug event, the correct fix may be upstream — deduplicate, sample, or batch. Also check the *record count* ceiling: a producer emitting many tiny records can hit 1,000 records/s at a fraction of the byte budget, and KPL aggregation solves that specific case without touching keys. ## The ceiling to say out loud A partition key hashes to exactly one hash-key range, therefore to exactly one shard at any moment. **One partition key can never carry more than one shard's worth of throughput.** No shard count, no capacity mode, no resharding strategy changes that. If a single logical entity genuinely produces more than 1 MB/s, the entity must be decomposed in the data model — salted, split by sub-entity, or moved to a stream of its own. Interviewers ask this question specifically to see whether you know the limit is structural rather than a tuning matter. ## While you fix it Producers must not drop throttled records. `PutRecords` is a partial-failure API: the call returns HTTP 200 with a non-zero `FailedRecordCount`, and only the failed entries carry an `ErrorCode`. Retry exactly those, with exponential backoff and jitter — and note that a retried record is appended after later successful ones, so retries themselves perturb order within the shard.
- You salt the hot tenant's key. What did you just break, and how do you compensate?Per-tenant ordering. Records for that tenant now spread across up to N shards with no relative order between them. Compensate by carrying an application-level sequence or event timestamp in the payload and re-ordering downstream, by making the consumer commutative or idempotent so order stops mattering, or by salting only the subset of event types that genuinely do not need ordering.
- How does the KCL keep per-key ordering correct across a SplitShard?A split closes the parent shard and creates children covering halves of its hash-key range. The KCL records the parent-child relationship in its lease table and will not begin a child shard until the parent has been read to its end and checkpointed. Since any given key lived in the parent before the split and in exactly one child after it, that ordering rule preserves per-key order across the reshard.
- Would switching the stream to on-demand mode have resolved this?Only partially. On-demand splits hot shards automatically, so it absorbs organic growth and mild skew without an operator. But it reacts over minutes rather than instantly, and it cannot split a single partition key across shards. If one key is the hot spot, on-demand throttles exactly as provisioned mode does.
saying these in an interview costs you the question
- Reads aggregate IncomingBytes and concludes there is headroom
- Adds shards without inspecting the key distribution
- Claims on-demand mode eliminates hot-shard throttling
- Thinks a single partition key can span multiple shards
- Treats PutRecords' HTTP 200 as proof every record was stored