skip to content

Kinesis Data Streams

You will learn AWS's managed record stream: shards as the unit of throughput and ordering, partition keys that decide placement, replayable retention, and consumers that lease shards and checkpoint. Interviewers ask when a stream beats a queue, and expect shard math and hot-key awareness in the answer.

part ofAWSoverview, primer and where to startread it →
on this pageshow

questions

6

In Amazon Kinesis Data Streams, what is a shard, and what throughput and ordering does a single shard give you?

level: juniorimportance: must knowfreq 78%

answer

  1. capacity is per shard, not per stream
  2. two write ceilings, bytes and count
  3. partition key hashed into a range
  4. sequence numbers order one shard only

basics

~20 s

A shard is Kinesis Data Streams' unit of capacity and ordering: roughly 1 MB/s or 1,000 records per second in, 2 MB/s of shared reads out. Records are placed by partition-key hash and are ordered only within one shard.

solid answer

~50 s

A Kinesis Data Streams stream is a set of shards, and every quota is per shard, not per stream. As of 2025 one shard accepts up to 1 MB/s **or** 1,000 records/s of writes — whichever ceiling you hit first — and serves 2 MB/s of reads shared across all standard consumers, with at most 5 `GetRecords` calls per second. A single record's data blob can be up to 1 MB. Placement is by partition key: Kinesis MD5-hashes the key into a 128-bit value and routes the record to the shard whose hash-key range contains it, so the same key always lands on the same shard while the shard map is stable. Within a shard, records carry increasing sequence numbers and are read in write order; across shards there is no ordering at all. So you size shard count from peak MB/s and peak records/s, and you pick the partition key from what must stay ordered together.

code

python · 18 lines
python
import boto3

kinesis = boto3.client("kinesis")

records = [
    {"Data": b'{"event":"click"}', "PartitionKey": f"user-{i}"}
    for i in range(500)
]

resp = kinesis.put_records(StreamName="clickstream", Records=records)

# PutRecords succeeds partially: failed entries carry an ErrorCode.
failed = [
    records[i]
    for i, r in enumerate(resp["Records"])
    if r.get("ErrorCode")
]
print(resp["FailedRecordCount"], "records need a retry")

go deeper

for a junior

Be able to say plainly that a stream is made of shards, that a shard is the unit of throughput and ordering, and that the partition key you pass on every write decides which shard a record lands in.

for a middle

Explain both write ceilings — bytes and record count — and show the arithmetic that turns a peak workload into a shard count. Know that reads are 2 MB/s shared per shard and that sequence numbers order records inside one shard only.

for a senior

Demonstrate that you size for the skewed peak rather than the average, and connect key choice to both ordering guarantees and hot-shard risk. Be ready to say what you would monitor to know the sizing was wrong.

for a principal

Own the tradeoff between per-key ordering and horizontal spread as a system-level decision: which entity really needs ordering, whether the consumer can be made order-insensitive instead, and what that buys you in shard count and cost.

## The shard is the stream A Kinesis Data Streams stream is not one pipe. It is a collection of **shards**, and essentially every limit you care about is expressed per shard. In provisioned capacity mode you choose the shard count explicitly; in on-demand mode AWS manages it for you, but the per-shard physics underneath are the same. A shard is an ordered, append-only sequence of records with its own throughput allowance and its own slice of the key space. ## The numbers that matter (as of 2025) Per shard: - **Ingest:** 1 MB per second **or** 1,000 records per second — whichever ceiling you hit first. - **Egress, shared fan-out:** 2 MB per second, and at most 5 `GetRecords` calls per second, each returning up to 10 MB or 10,000 records. - **Record size:** a single record's data blob may be up to 1 MB. The two-sided ingest limit is the one juniors miss. A firehose of tiny 200-byte events hits the 1,000 records/s ceiling at about 200 KB/s — a fifth of the byte budget. That is exactly why the Kinesis Producer Library (KPL) *aggregates* many application records into one Kinesis record, and why the KCL deaggregates them on the way out. Blow the write limit and the API returns `ProvisionedThroughputExceededException`; blow the read limit and `GetRecords` is throttled, visible in the CloudWatch metrics `WriteProvisionedThroughputExceeded` and `ReadProvisionedThroughputExceeded`. ## Partition keys decide placement Every write carries a **partition key**: a UTF-8 string of up to 256 bytes that you supply. Kinesis MD5-hashes it into a 128-bit integer. Each shard owns a contiguous *hash-key range*, and the record is stored in whichever shard's range contains that hash. `PutRecord` and `PutRecords` also accept an `ExplicitHashKey`, which bypasses the hash and targets a range directly — useful when you want to steer traffic at a specific shard. Because placement is a hash of a key you choose, the key is the only lever you have over distribution. A key with few distinct values concentrates traffic on few shards; a key with high cardinality spreads it. ## Ordering, precisely Within a shard, records are stored in arrival order and each gets a **sequence number** that increases monotonically for that shard. A consumer reading that shard sees them in that order. Across shards there is **no** ordering relationship whatsoever — two records written a millisecond apart to different shards can be processed hours apart. The practical rule: everything that must be ordered relative to each other must share a partition key. Order events by `customer_id` and one customer's history is strictly ordered; order them by a random UUID and you get maximum spread and no ordering at all. That tension — ordering versus spread — is the whole design conversation. ## Reading a shard A standard consumer calls `GetShardIterator` to get a position, then loops on `GetRecords`. The iterator type chooses where to start: `TRIM_HORIZON` (oldest retained record), `LATEST` (only new records), `AT_TIMESTAMP`, `AT_SEQUENCE_NUMBER`, or `AFTER_SEQUENCE_NUMBER`. Each `GetRecords` response includes a `NextShardIterator` to continue with and `MillisBehindLatest`, which tells you how far behind the tip you are. ```python it = kinesis.get_shard_iterator( StreamName="clickstream", ShardId="shardId-000000000000", ShardIteratorType="TRIM_HORIZON", )["ShardIterator"] resp = kinesis.get_records(ShardIterator=it, Limit=500) print(resp["MillisBehindLatest"], len(resp["Records"])) ``` Reading does not consume: records stay for the retention period (24 hours by default) regardless of how many consumers have read them, which is why a second consumer application can be added later and replay history. ## Sizing math Take peak write throughput and compute both ceilings, then take the larger: ``` shards = max( ceil(peak_MB_per_sec / 1), ceil(peak_records_per_sec / 1000) ) ``` Then sanity-check reads: with standard consumers, all of them together share 2 MB/s per shard, so three consumer apps on a stream running at 1 MB/s of ingest are already at the edge. Finally add headroom, because real partition keys are never perfectly uniform and the limits bind per shard, not on the average. ## Where candidates go wrong The classic mistake is quoting the limits as if they applied to the stream. A 10-shard stream is not "a 10 MB/s stream" if 80% of records carry the same partition key — it is a 1 MB/s stream with nine idle shards. The second mistake is assuming reads are elastic because the service is managed; the 2 MB/s shared read budget is the reason enhanced fan-out exists.

  • A producer sends 200-byte events at 900 KB/s to a one-shard stream and starts getting throttled. Why?
    Because 200-byte events at 900 KB/s is about 4,500 records per second, and a shard caps at 1,000 records/s as well as 1 MB/s. The record-count ceiling binds first. The usual fix is aggregation — batch many application events into one Kinesis record (the KPL does this automatically, and the KCL deaggregates on the consumer side) — or add shards.
  • If two separate consumer applications read the same stream, does each get its own 2 MB/s per shard?
    Not with the default standard consumer. The 2 MB/s per shard is shared across every standard consumer on that shard, so two applications each average about 1 MB/s and can throttle each other's `GetRecords` calls. Registering a consumer for enhanced fan-out gives it a dedicated 2 MB/s per shard instead, at extra cost.
  • Does reading a record remove it from the shard?
    No. Kinesis Data Streams is replayable storage: records remain until the retention period expires, independent of who has read them. That is what lets you add a new consumer application later and start it at `TRIM_HORIZON` to reprocess history, and what makes checkpointing the consumer's own responsibility rather than the service's.

saying these in an interview costs you the question

  • Quotes 1 MB/s as a per-stream limit rather than per shard
  • Claims Kinesis guarantees global ordering across the whole stream
  • Forgets the 1,000 records/s ceiling and sizes only on bytes
  • Thinks adding consumers adds read throughput to a shard
  • Believes reading a record deletes it from the stream

context

open as a page

When would you run a Kinesis Data Streams stream in on-demand capacity mode instead of provisioned, and what do you give up by doing so?

level: middleimportance: should knowfreq 52%

basics

~10 s

On-demand suits unpredictable or new workloads: AWS manages shard count and you pay per stream-hour plus data volume. Provisioned is cheaper at steady, well-understood throughput but makes you size and reshard the stream yourself.

open as a page

In Kinesis Data Streams, how does an enhanced fan-out consumer differ from the default shared-throughput consumer, and when is the extra cost justified?

level: middleimportance: should knowfreq 45%

basics

~20 s

A standard consumer polls GetRecords and shares one shard's 2 MB/s read budget with every other standard consumer. An enhanced fan-out consumer registers with the stream and gets its own 2 MB/s per shard, pushed over HTTP/2, at extra cost.

open as a page

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?

level: seniorimportance: should knowfreq 50%

basics

~20 s

A 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.

open as a page

How does the Kinesis Client Library distribute a stream's shards across worker instances and remember progress, and what operational problems come from that design?

level: seniorimportance: nice to knowfreq 36%

basics

~20 s

The KCL keeps a DynamoDB lease table, one item per shard, holding the owning worker and the last checkpointed sequence number. Workers heartbeat to hold leases and steal stale ones, which makes processing at-least-once and puts a second billable table in the critical path.

open as a page

A Kinesis Data Streams stream retains 24 hours of records by default. How do you decide between extending that retention and archiving the records elsewhere so you can replay them later?

level: principalimportance: nice to knowfreq 30%

basics

~20 s

Extended retention keeps replay on the same code path and preserves per-shard ordering, but you pay stream prices for cold data and replay is capped by shard read throughput. Archiving to object storage is far cheaper and unbounded, at the price of a second read path.

open as a page