Prev Next

Cloud / Amazon Kinesis Interview questions

Last updated

1. What is Amazon Kinesis? 2. What are the four services in the Amazon Kinesis family? 3. What is Amazon Kinesis Data Streams? 4. What is a shard in Kinesis Data Streams? 5. What is a partition key in Kinesis? 6. What is a data record in Kinesis Data Streams? 7. What are producers and consumers in Kinesis Data Streams? 8. What is Amazon Data Firehose? 9. What destinations does Amazon Data Firehose support? 10. What is Amazon Managed Service for Apache Flink? 11. What is Amazon Kinesis Video Streams? 12. What is the Kinesis Producer Library (KPL)? 13. What is the Kinesis Client Library (KCL)? 14. What are the capacity modes of Kinesis Data Streams? 15. What is the difference between PutRecord and PutRecords? 16. What are the shard iterator types in Kinesis? 17. How do you write a record to a stream using the AWS CLI? 18. How do you read records from a stream using boto3? 19. What is the Amazon Kinesis Agent? 20. How does Kinesis decide which shard receives a record? 21. What are the throughput limits of a single shard? 22. How do you calculate the number of shards required? 23. What happens when you exceed a shard's throughput limit? 24. What is the difference between Kinesis Data Streams and Firehose? 25. What is the difference between shared and enhanced fan-out consumers? 26. How does enhanced fan-out deliver records to consumers? 27. How does the KCL track progress using leases and checkpoints? 28. How does Kinesis guarantee record ordering? 29. How does Lambda consume records from a Kinesis stream? 30. What is the parallelization factor for Kinesis Lambda triggers? 31. How does Firehose buffering work? 32. How does Firehose convert records to Parquet? 33. How does Firehose dynamic partitioning work? 34. What is the difference between Kinesis Data Streams and SQS? 35. What is the difference between Kinesis Data Streams and Apache Kafka? 36. How do you encrypt data in Kinesis Data Streams? 37. How do you restrict access to a stream using IAM? 38. How do you replay old data from a Kinesis stream? 39. Why can Kinesis consumers receive duplicate records? 40. When would you choose on-demand over provisioned capacity mode? 41. How do you monitor a Kinesis data stream with CloudWatch? 42. How does Kinesis Video Streams ingest and play back video? 43. What happens when you split or merge a shard? 44. How do you preserve ordering during resharding? 45. How do you handle a hot shard? 46. How do you troubleshoot high iterator age? 47. How do you handle poison pill records in a Kinesis consumer? 48. How do you build idempotent processing on Kinesis? 49. Explain the lifecycle of a record from producer to consumer? 50. Explain the internal working of KPL aggregation and collection?

1. What is Amazon Kinesis?

Amazon Kinesis is the AWS family of managed services for collecting, processing, and analyzing streaming data in real time. Instead of waiting for a nightly batch job, you can react to clickstreams, logs, IoT telemetry, metrics, or video within seconds of the events happening.

AWS runs the infrastructure for you. There are no brokers to patch and no replication to configure, and data written to a stream is stored across three Availability Zones in the Region.

Typical uses include real-time dashboards, log and metric aggregation, fraud detection, IoT ingestion, and feeding data lakes or warehouses.

Take quiz
Which workload fits Kinesis best?
running a nightly payroll batch job
ingesting clickstream events and reacting within seconds
storing archival backups for ten years
Who operates the servers behind a Kinesis stream?
AWS, as a managed service
the customer, who patches the brokers
a third-party Kafka vendor

2. What are the four services in the Amazon Kinesis family?

The Kinesis family has four services, each solving a different part of the streaming problem.

Service What it does
Kinesis Data Streams Ingests and stores ordered records so custom consumers can read and replay them.
Amazon Data Firehose Loads streaming data into S3, Redshift, OpenSearch, Splunk and more without consumer code.
Managed Service for Apache Flink Runs stateful stream processing (windows, joins, aggregations) written in Java, Python, or SQL.
Kinesis Video Streams Ingests, stores, and plays back video, audio, and other time-serialized media.

Firehose was previously called Kinesis Data Firehose, and Managed Service for Apache Flink was called Kinesis Data Analytics.

Take quiz
Which service loads streaming data into S3 with no consumer code to write?
Kinesis Data Streams
Kinesis Video Streams
Amazon Data Firehose
Managed Service for Apache Flink was formerly named:
Kinesis Data Analytics
Kinesis Data Firehose
Kinesis Stream Processor
Kinesis Video Analytics

3. What is Amazon Kinesis Data Streams?

Kinesis Data Streams (KDS) is a serverless service for ingesting large volumes of streaming data and making it available to multiple consumers. Producers write records to a stream, and the stream keeps them in order within each shard.

Records stay available for 24 hours by default, and you can extend retention up to 365 days. Because records aren't deleted when read, several applications can consume the same data independently, and each can replay from an earlier point.

Capacity comes from shards, and you can run a stream in provisioned or on-demand mode.

Take quiz
What is the default retention period of a Kinesis data stream?
1 hour
7 days
until the first consumer reads the record
24 hours
Unlike a typical queue, a Kinesis stream lets consumers:
delete records once they have read them
read only if a single consumer is registered
re-read the same records independently within retention

4. What is a shard in Kinesis Data Streams?

A shard is the base unit of capacity in a stream: an ordered sequence of records with fixed throughput limits. Each shard accepts up to 1 MiB/s or 1,000 records/s of writes and supports 2 MiB/s of reads shared across standard consumers.

Every shard owns a slice of the 128-bit hash key space, and a record lands in the shard whose range contains the hash of its partition key. Shards are named like shardId-000000000000.

A stream's total capacity is simply the sum of its shards, so you scale by adding shards, either yourself or automatically in on-demand mode.

Take quiz
What is the write limit of one shard?
10 MiB/s or 10,000 records/s
1 MiB/s or 1,000 records/s
2 MiB/s or 500 records/s
The capacity of a stream is:
the sum of the capacity of its shards
fixed at 1 MiB/s regardless of shard count
determined by the number of consumers

5. What is a partition key in Kinesis?

A partition key is a Unicode string, up to 256 characters, that the producer attaches to every record. Kinesis runs it through MD5 to get a 128-bit value and uses that value to pick the shard.

Records with the same key go to the same shard, which is how you get per-key ordering, for example all events for one user_id in sequence.

Pick a key with high cardinality such as a device or customer ID. A key with few values, like a country code, funnels traffic into a handful of shards and causes hot shards and throttling.

Take quiz
Records that share a partition key are routed to:
random shards chosen per record
a separate stream for each key
the same shard
Which partition key gives the most even distribution?
the constant string 'events'
the AWS Region name
the stream name
a device ID with millions of distinct values

6. What is a data record in Kinesis Data Streams?

A data record is the unit of data stored in a stream. It has three parts that matter to you.

Field Set by Notes
Partition key Producer Decides the shard; up to 256 characters.
Sequence number Kinesis Unique per shard and increasing; assigned at ingestion.
Data blob Producer Opaque bytes, up to 1 MiB before base64 encoding.

Kinesis also records an approximate arrival timestamp. It never inspects the blob, so JSON, Avro, Protobuf, or plain text all work.

The partition key travels with every record to consumers, so it is a handy place to carry routing information such as a tenant ID. The 1 MiB limit applies to the data blob, so larger payloads usually go to S3 with a pointer stored in the record.

Take quiz
Who assigns the sequence number of a record?
the producer
the consumer on first read
Kinesis Data Streams at ingestion
the KMS key
What is the maximum size of a data blob?
1 MiB
256 KB
5 MiB
128 KB

7. What are producers and consumers in Kinesis Data Streams?

Producers put records into the stream and consumers read and process them. A single stream can have many of each.

Role Examples
Producers AWS SDK, Kinesis Producer Library, Kinesis Agent, IoT Core rules, CloudWatch Logs subscriptions
Consumers KCL applications, Lambda, Firehose, Managed Service for Apache Flink, AWS Glue streaming, Spark on EMR

Consumers track their own position in the stream, so a slow analytics job doesn't hold back a fast alerting job.

In the standard model Kinesis doesn't track consumers for you. Each consumer pulls records and remembers its own position, which is why the KCL keeps checkpoints in DynamoDB. One application can be both, such as a Flink job that reads one stream and writes results to another.

Take quiz
Which of these is a consumer of a stream?
Kinesis Agent tailing a log file
a Lambda function triggered by the stream
a KPL process on a web server
The Kinesis Agent is best described as:
a producer that ships log files to a stream or Firehose
a consumer that checkpoints to DynamoDB
a tool that splits shards

8. What is Amazon Data Firehose?

Amazon Data Firehose is a fully managed service that captures streaming data and delivers it to a destination such as S3, Redshift, or OpenSearch. There are no shards to size and no consumer code to run, and it scales automatically.

It buffers records by size or time, so delivery is near real time rather than per record. Along the way it can transform data with Lambda, convert JSON to Parquet or ORC, compress, and encrypt.

Sources include direct PUT calls, Kinesis Data Streams, MSK, and CloudWatch Logs. Firehose doesn't keep data for replay, so use Kinesis Data Streams when you need that.

Take quiz
How does Firehose deliver data?
once a day on a fixed schedule only
only when a consumer polls for it
in near real time, after buffering by size or interval
Do you manage shards in Firehose?
Yes, you provision shards per delivery stream
No, it scales automatically
Yes, but only for S3 destinations

9. What destinations does Amazon Data Firehose support?

Firehose can deliver to several AWS and third-party destinations.

Destination Notes
Amazon S3 Most common; supports compression, partitioning, and Parquet/ORC conversion.
Amazon Redshift Data is staged in S3 first, then loaded with a COPY command.
OpenSearch Service / Serverless Good for log search and dashboards.
Splunk Pushes events to a Splunk HTTP Event Collector.
Apache Iceberg tables, Snowflake Lakehouse and warehouse delivery.
HTTP endpoint Custom endpoints and partners such as Datadog or New Relic.

Records that fail delivery can be backed up to an S3 bucket for later inspection.

Take quiz
Firehose loads data into Redshift by:
staging it in S3 and running a COPY command
inserting each record over JDBC
streaming directly with no staging
Which of these is NOT a native Firehose destination?
Splunk
Amazon OpenSearch Service
Amazon S3
Amazon RDS for MySQL

It is a serverless service for running Apache Flink applications on streaming data without managing clusters. It used to be called Kinesis Data Analytics for Apache Flink.

You write Flink code in Java, Scala, Python, or SQL, or explore interactively in Studio notebooks. Typical sources are Kinesis Data Streams, MSK, and S3, and sinks include Kinesis, Firehose, S3, and databases.

Flink brings stateful processing, event-time windows, joins, and checkpointing, which is what you need for things like rolling averages or sessionization that a simple Lambda consumer struggles with.

11. What is Amazon Kinesis Video Streams?

Kinesis Video Streams ingests, stores, and plays back video, audio, and other time-serialized data from cameras and devices. Devices send media using the Producer SDK or a GStreamer plugin.

Media is stored as fragments indexed by timestamp. You can play it back with HLS or MPEG-DASH, fetch it with GetMedia, or use WebRTC for low-latency two-way streaming.

It integrates with services such as Rekognition Video and SageMaker for analysis, which makes it common in smart-camera and security applications.

Retention is set per stream in hours, and a value of zero means media is only available live with nothing archived. Longer retention lets you replay footage after an incident.

Take quiz
Which protocols does Kinesis Video Streams support for playback?
SMTP and IMAP
FTP only
SNMP
HLS and MPEG-DASH
Kinesis Video Streams is typically used for:
moving JSON clickstream events into S3
streaming camera and device video for storage and analysis
running SQL over application logs

12. What is the Kinesis Producer Library (KPL)?

The KPL is an AWS library that makes writing to Kinesis Data Streams more efficient and reliable. It runs a native process under a Java API and handles the details producers otherwise code by hand.

  • Aggregation: packs many user records into one Kinesis record.
  • Collection: batches records into PutRecords calls.
  • Retries: automatic retry with a configurable record TTL.
  • Metrics: publishes throughput and error metrics to CloudWatch.

The trade-off is a little added latency, because records wait in a buffer (RecordMaxBufferedTime, 100 ms by default) before being sent.

Take quiz
KPL aggregation means:
packing multiple user records into a single Kinesis record
compressing every record with GZIP
encrypting records with KMS
A trade-off of using the KPL is:
it cannot retry failed writes
records are limited to 100 bytes
a small amount of added latency from buffering

13. What is the Kinesis Client Library (KCL)?

The KCL is a library for building consumer applications on Kinesis Data Streams. You implement a record processor, and the KCL takes care of the plumbing: spreading shards across workers, checkpointing, failover, and reacting to resharding.

It coordinates workers through a DynamoDB lease table named after your application name, with one lease per shard. KCL 2.x supports enhanced fan-out, and KCL 3.0 rebalances shards based on worker CPU utilization.

It also deaggregates records written by the KPL automatically.

In KCL 2.x the record processor gets callbacks for initialization, each batch of records, lease loss, and shard end. It must checkpoint when a shard ends so the child shards can be processed.

Take quiz
KCL stores lease and checkpoint data in:
Amazon S3
Amazon DynamoDB
CloudWatch Logs
Amazon ElastiCache
The name of the KCL's DynamoDB table comes from:
the application name you configure
the AWS account ID
the IAM role of the producer

14. What are the capacity modes of Kinesis Data Streams?

A stream runs in one of two capacity modes. In provisioned mode you choose the shard count and pay per shard-hour plus PUT payload units. In on-demand mode Kinesis scales capacity for you and you pay per stream-hour and per GB written and read.

You can switch an existing stream between the two modes, but only twice in 24 hours. Both modes support enhanced fan-out and encryption.

Switching modes causes no downtime for producers or consumers. In provisioned mode you can call UpdateShardCount or split and merge shards, while in on-demand mode those calls aren't available because Kinesis manages the shards.

Take quiz
In provisioned mode, cost is mainly driven by:
the number of consumers only
GET requests only
shard-hours plus PUT payload units
In on-demand mode you:
must call UpdateShardCount every hour
don't choose a shard count; capacity scales automatically
cannot use enhanced fan-out

15. What is the difference between PutRecord and PutRecords?

PutRecord writes one record per API call, while PutRecords writes a batch in a single call, which is far more efficient for throughput.

PutRecord PutRecords
Records per call 1 Up to 500
Request size 1 MiB record Up to 5 MiB total, 1 MiB per record
Failure handling Whole call succeeds or fails Partial failures possible
Strict ordering SequenceNumberForOrdering supported Not available

With PutRecords, always check FailedRecordCount and retry only the failed entries, otherwise you create duplicates.

Prefer PutRecords for high-volume producers, since per-call overhead dominates when records are small. Both calls count against the same per-shard limits, so batching reduces the number of HTTP requests, not shard throughput.

Take quiz
A single PutRecords call can contain up to:
100 records
1,000 records
10 records
500 records
If FailedRecordCount is greater than zero, the producer should:
resend the entire batch
discard all records in the call
retry only the failed records

16. What are the shard iterator types in Kinesis?

A shard iterator tells GetRecords where to start reading. You get one by calling GetShardIterator with one of five types.

Type Starts reading
TRIM_HORIZON At the oldest record still in the shard
LATEST Just after the newest record
AT_TIMESTAMP At the first record at or after a given time
AT_SEQUENCE_NUMBER At the exact sequence number given
AFTER_SEQUENCE_NUMBER Right after the sequence number given

An iterator expires after 5 minutes, so keep using the NextShardIterator returned by each call.

In practice, use LATEST for a monitor that only cares about new data, TRIM_HORIZON for a first-time backfill, and AT_TIMESTAMP to start from a moment such as the start of an incident.

Take quiz
TRIM_HORIZON starts reading from:
the newest record only
the oldest record still in the shard
a timestamp you supply
A shard iterator expires after:
24 hours
60 seconds
5 minutes
it never expires

17. How do you write a record to a stream using the AWS CLI?

Use aws kinesis put-record with the stream name, a partition key, and the data.

aws kinesis put-record \
  --stream-name orders \
  --partition-key user-42 \
  --cli-binary-format raw-in-base64-out \
  --data '{"orderId":1001,"amount":49.90}'

AWS CLI v2 expects base64 input by default, so --cli-binary-format raw-in-base64-out lets you pass plain text. The response returns the ShardId and the SequenceNumber assigned to your record.

Create the stream first, and wait until it is ACTIVE before writing:

aws kinesis create-stream --stream-name orders --shard-count 2
aws kinesis describe-stream-summary --stream-name orders

To read it back from the CLI, use get-shard-iterator and then get-records.

Take quiz
Why add --cli-binary-format raw-in-base64-out in CLI v2?
to enable KMS encryption
to choose the target shard
so --data is accepted as raw text instead of base64
The put-record response contains:
the full record payload
the ShardId and SequenceNumber
the list of registered consumers

18. How do you read records from a stream using boto3?

Reading directly takes three steps: find the shard, get a shard iterator, then call get_records in a loop using the iterator each response hands back.

import boto3, json, time

kds = boto3.client("kinesis")
shard = kds.describe_stream(StreamName="orders")["StreamDescription"]["Shards"][0]["ShardId"]
it = kds.get_shard_iterator(StreamName="orders", ShardId=shard,
                            ShardIteratorType="TRIM_HORIZON")["ShardIterator"]
while it:
    resp = kds.get_records(ShardIterator=it, Limit=100)
    for r in resp["Records"]:
        print(r["SequenceNumber"], json.loads(r["Data"]))
    it = resp["NextShardIterator"]
    time.sleep(1)

The sleep keeps you under 5 GetRecords calls per second per shard. For production, prefer the KCL or Lambda, which handle shards, checkpoints, and resharding for you.

Take quiz
What do you need before calling get_records?
a shard iterator from get_shard_iterator
a consumer ARN
a partition key
Why sleep between get_records calls?
to avoid partition key collisions
to refresh the KMS key
to stay under 5 GetRecords calls per second per shard

19. What is the Amazon Kinesis Agent?

The Kinesis Agent is a standalone Java application you install on Linux servers. It watches log files, handles rotation, and sends new lines to Kinesis Data Streams or Firehose.

It keeps its own checkpoint, retries on failure, and can preprocess data before sending, for example turning a multi-line entry into one line, or converting CSV or log lines to JSON. Configuration lives in /etc/aws-kinesis/agent.json.

{
  "flows": [{
    "filePattern": "/var/log/app/*.log",
    "kinesisStream": "app-logs",
    "partitionKeyOption": "RANDOM"
  }]
}

Use the Agent when your source is plain log files. Its partitionKeyOption can be RANDOM to spread load evenly or DETERMINISTIC to hash a value. For application events, call the SDK or KPL directly instead.

Take quiz
Where does the Kinesis Agent run?
inside the Kinesis service
on the source servers that produce the log files
only inside Lambda
Which targets can the Agent send data to?
Kinesis Data Streams and Firehose
only Amazon S3
only SQS

20. How does Kinesis decide which shard receives a record?

Kinesis computes the MD5 hash of the partition key, which gives a 128-bit integer. Each shard owns a contiguous HashKeyRange within 0 to 2128-1, and the record goes to the shard whose range contains that integer.

A producer can bypass hashing by passing an ExplicitHashKey, which is used directly as the hash value. That is useful for targeting a specific shard, but you have to know the current ranges.

flowchart LR
  A["Partition key"] --> B["MD5 hash"]
  B --> C["128-bit integer"]
  C --> D{Which HashKeyRange contains it?}
  D --> E["Shard 1"]
  D --> F["Shard 2"]
  D --> G["Shard N"]
Take quiz
Which hash function does Kinesis apply to the partition key?
MD5
SHA-256
CRC32
none, keys are matched as text
ExplicitHashKey lets a producer:
change the stream's retention period
choose the target shard directly
skip encryption

21. What are the throughput limits of a single shard?

Each shard has separate write and read limits, and the read side depends on the consumer type.

Direction Limit per shard
Writes 1 MiB/s and 1,000 records/s
Reads, shared throughput 2 MiB/s total across all consumers, 5 GetRecords calls/s
Reads, enhanced fan-out 2 MiB/s per registered consumer
Single GetRecords call Up to 10 MiB or 10,000 records

Whichever write limit you hit first wins, so many small records can throttle you on record count before bandwidth.

These limits are per shard, so 10 shards give 10 MiB/s of writes only if keys spread evenly. Enhanced fan-out consumers have dedicated throughput, so adding one doesn't slow down existing shared consumers.

Take quiz
Shared-throughput reads per shard are:
2 MiB/s total, shared by all consumers
2 MiB/s for each consumer
1 MiB/s in total
How many GetRecords calls per second does a shard allow?
1
100
1,000
5

22. How do you calculate the number of shards required?

Work out the shard count for writes and for reads, then take the larger. Example: 10,000 records/s of 2 KB each, read by 3 applications on shared throughput.

  1. Write bandwidth: 10,000 x 2 KB is about 19.5 MiB/s, so 20 shards at 1 MiB/s each.
  2. Write records: 10,000 / 1,000 = 10 shards. Bandwidth is the bigger constraint.
  3. Read bandwidth: 3 consumers x 19.5 = 58.6 MiB/s, divided by 2 MiB/s gives 30 shards.
  4. Decision: 30 shards on shared throughput, or 20 shards if the consumers use enhanced fan-out.

Add headroom of 20-30% for bursts, or choose on-demand mode if the traffic is hard to predict.

Take quiz
Writes total 5 MiB/s at 800 records/s. Minimum shards for writes?
1
3
5
10
Writes total 6 MiB/s and 3 shared-throughput consumers each read everything. Minimum shards?
9
6
3
18

23. What happens when you exceed a shard's throughput limit?

Kinesis throttles the request and returns ProvisionedThroughputExceededException. For PutRecords, the call itself succeeds but the throttled entries come back as failed records.

CloudWatch shows it through WriteProvisionedThroughputExceeded and ReadProvisionedThroughputExceeded. Producers should retry with exponential backoff and jitter.

  • Spread load with a higher-cardinality partition key.
  • Add shards, or switch to on-demand mode.
  • Use enhanced fan-out if read limits are the problem.
  • Use the KPL to aggregate small records.

The exception message names the shard that was throttled, which helps you find hot shards. For PutRecords, inspect each entry's ErrorCode, since only some entries may have failed. The SDKs retry a few times by default, and the KPL keeps retrying until RecordTtl expires.

Take quiz
The best immediate producer response to throttling is:
retry in a tight loop
drop all records
retry with exponential backoff and jitter
Throttling on one shard while others sit idle usually means:
KMS key rotation
a hot shard caused by skewed partition keys
too many streams in the account

24. What is the difference between Kinesis Data Streams and Firehose?

Kinesis Data Streams is a storage-and-delivery layer you build consumers on, while Firehose is a managed pipe that loads data into a destination for you.

Data Streams Firehose
Capacity Shards or on-demand Fully automatic
Consumers Custom (KCL, Lambda, Flink) Managed destinations only
Latency Around 70-200 ms Near real time, based on buffering
Retention / replay 24 hours up to 365 days, replayable None
Ordering Per shard Not a design goal
Transformation In your consumer Built-in Lambda and format conversion

A common pattern uses both: producers write to a data stream, one consumer does real-time work, and Firehose reads the same stream to land data in S3.

Take quiz
Which service lets you replay data from an earlier point?
Kinesis Data Streams
Firehose
neither of them
You need data in S3 with minimal code and no capacity planning. Choose:
Kinesis Data Streams with a KCL app
Kinesis Video Streams
Amazon Data Firehose
a custom Flink sink

25. What is the difference between shared and enhanced fan-out consumers?

Shared consumers split a shard's 2 MiB/s read capacity between them, while each enhanced fan-out (EFO) consumer gets its own full 2 MiB/s per shard.

Shared throughput Enhanced fan-out
Read capacity 2 MiB/s per shard, shared 2 MiB/s per consumer per shard
Delivery Pull with GetRecords Push over HTTP/2 with SubscribeToShard
Typical latency About 200 ms About 70 ms
Cost Included Extra per consumer-shard-hour and per GB
Limit No registration 20 registered consumers per stream by default

Use EFO when several applications each need low latency and full throughput. For one or two light consumers, shared throughput is cheaper.

EFO costs more because it is billed per consumer-shard pair per hour plus per GB retrieved. Five EFO consumers on a 20-shard stream means 100 consumer-shard pairs, so it is worth checking before enabling it everywhere.

Take quiz
Enhanced fan-out delivers records using:
SQS long polling
HTTP/2 push through SubscribeToShard
S3 event notifications
Choose EFO when:
several consumers each need full throughput with low latency
a single low-volume consumer reads occasionally
you want the cheapest option for one nightly batch job

26. How does enhanced fan-out deliver records to consumers?

A consumer first registers with RegisterStreamConsumer and receives a consumer ARN. It then calls SubscribeToShard for each shard, which opens an HTTP/2 connection over which Kinesis pushes batches of records as events.

A subscription lasts at most 5 minutes, so the consumer must resubscribe, starting from the last sequence number it processed. KCL 2.x and Lambda do this for you.

sequenceDiagram
  participant C as Consumer
  participant K as Kinesis
  C->>K: RegisterStreamConsumer
  K-->>C: Consumer ARN
  C->>K: SubscribeToShard (HTTP/2)
  K-->>C: Push record batches
  Note over C,K: Subscription ends after 5 minutes
  C->>K: SubscribeToShard again

Push delivery avoids the 5-calls-per-second polling limit and the idle gaps between polls, which is where most of the latency gain comes from. Register consumers under a stable name so the consumer ARN stays the same across restarts.

Take quiz
How long can a single SubscribeToShard connection last?
24 hours
10 seconds
up to 5 minutes, then it must be renewed
until the shard is merged
What is the first step to use enhanced fan-out?
increase the retention period
register the consumer with the stream
enable KMS encryption

27. How does the KCL track progress using leases and checkpoints?

The KCL keeps a lease per shard in a DynamoDB table. A lease row holds the current owner, a counter, and the latest checkpoint. A worker may process a shard only while it holds that shard's lease.

Workers renew their leases regularly. If a worker dies, its leases expire and other workers take them over, and workers can also steal leases to balance load.

A checkpoint stores the sequence number you've fully processed. After a restart, reading resumes after it. Checkpoint only after processing succeeds, which gives at-least-once behavior; checkpointing earlier risks losing records.

Give the lease table enough DynamoDB capacity, since on-demand billing works well here. Failover speed depends on the lease duration (failoverTimeMillis, 10 seconds by default), so shorter values recover faster but renew more often.

Take quiz
To get at-least-once processing, checkpoint:
after the records have been processed successfully
before processing each batch
never, because KCL guesses the position
When a KCL worker crashes, its leases are:
deleted along with the shard data
frozen until someone reassigns them by hand
taken over by other workers after they expire

28. How does Kinesis guarantee record ordering?

Ordering is guaranteed within a shard: records are stored in sequence-number order and read back in that order. Because one partition key always maps to one shard, you get ordering per key, but not across the whole stream.

  • Producers: send records for the same key sequentially. Retried PutRecords entries can land after newer ones, so use PutRecord with SequenceNumberForOrdering when strict order matters.
  • Consumers: only one reader per shard at a time, which the KCL lease enforces.
  • Resharding: read parent shards before their children.

A common mistake is expecting a global order. If events from different entities must be ordered relative to each other, they need the same key, or the stream must be one shard, which caps throughput. Per-entity order is usually enough.

Take quiz
Kinesis guarantees ordering:
across the entire stream
within a shard, which gives order per partition key
only when Firehose is used
Retried PutRecords failures may:
arrive after later records and break per-key order
be deduplicated automatically by Kinesis
be sent to another Region

29. How does Lambda consume records from a Kinesis stream?

An event source mapping polls each shard (or reads through an enhanced fan-out consumer), collects records into a batch, and invokes your function synchronously. Batch size goes up to 10,000 records and the batch window up to 300 seconds. Record data arrives base64-encoded.

Within a shard, batches are processed in order. If the function fails, Lambda retries the same batch until it succeeds or the records expire, which can block the shard. Controls to limit that:

  • MaximumRetryAttempts and MaximumRecordAgeInSeconds
  • BisectBatchOnFunctionError to split a failing batch
  • An on-failure destination (SQS, SNS, or S3) for discarded batches
  • ReportBatchItemFailures to return the failing sequence number
{"batchItemFailures": [{"itemIdentifier": "49590338271490256608559692538361571095921575989136588898"}]}

Take quiz
Processing of one shard by Lambda is:
fully parallel per record
randomized across batches
sequential and in order
ReportBatchItemFailures lets your function:
skip every failed batch silently
report where it failed so only the remaining records are retried
increase the shard count

30. What is the parallelization factor for Kinesis Lambda triggers?

The parallelization factor sets how many batches Lambda can process concurrently per shard. The default is 1, and the maximum is 10.

Records with the same partition key are still processed in order, because Lambda keeps them on the same concurrent worker. Total concurrency can reach shards x factor, so a 4-shard stream with a factor of 5 can run up to 20 invocations at once.

Raise it when IteratorAgeMilliseconds is climbing because your function is slow, rather than adding shards. Remember that downstream databases will see the extra concurrency.

Higher concurrency means more simultaneous invocations, which can use up account concurrency or overload a database with connections. Reserve function concurrency and check downstream limits before raising the factor.

Take quiz
What is the maximum parallelization factor?
5
100
10
2
A 4-shard stream with a factor of 5 allows up to how many concurrent batches?
20
4
9
5

31. How does Firehose buffering work?

Firehose collects records until either the buffer size or the buffer interval is reached, whichever comes first, and then delivers a batch.

For S3 the size can be 1-128 MiB and the interval 0-900 seconds (default 300). Setting the interval to 0 turns on zero buffering for lower latency. Dynamic partitioning needs a larger buffer, at least 64 MiB.

Smaller buffers mean fresher data but many small S3 objects, which hurts Athena and Spark scans. Larger buffers produce fewer, bigger files at the cost of delay.

If a Lambda transformation is configured, Firehose invokes it on each buffered batch using its own smaller buffer. When delivery fails, Firehose retries for a limited period and can back failed records up to S3.

Take quiz
Firehose delivers a batch when:
both thresholds are reached
either the buffer size or interval is reached first
only when the buffer interval ends
A very small buffer interval typically gives:
lower latency but many small S3 objects
higher latency and larger files
no effect on file count

32. How does Firehose convert records to Parquet?

Enable record format conversion on the delivery stream and point it at a table in the AWS Glue Data Catalog. The table supplies the schema, and Firehose writes Parquet or ORC files to S3.

The incoming data must be JSON, read with the OpenX or Hive JSON SerDe. If your records aren't JSON, use a Lambda transformation first, and Firehose runs it before conversion.

Columnar files compress better and let Athena scan only the columns it needs, which reduces storage and query cost compared with raw JSON.

Records that fail conversion, such as malformed JSON or values that don't fit the schema, are written to an S3 error prefix instead of being dropped, so you can inspect and fix them. Parquet compression defaults to Snappy and can be changed.

Take quiz
Firehose gets the schema for conversion from:
the Kinesis stream's metadata
an S3 bucket policy
an AWS Glue Data Catalog table
Format conversion expects the input to be:
CSV
Avro
JSON
images

33. How does Firehose dynamic partitioning work?

Dynamic partitioning lets Firehose write S3 objects under prefixes built from the contents of each record, such as customer_id or event date. Keys are extracted with an inline jq expression or by a Lambda function.

customer_id=!{partitionKeyFromQuery:customer_id}/dt=!{timestamp:yyyy-MM-dd}/

It works for S3 destinations only and needs a buffer of at least 64 MiB. Each delivery stream has a default limit of 500 active partitions. The benefit is partition pruning in Athena and Redshift Spectrum, with no separate ETL step.

For example, an event for customer 77 on 2026-10-04 lands under customer_id=77/dt=2026-10-04/, so a query filtered by customer and day reads only that folder. Avoid keys with thousands of values, which create tiny files and approach the active-partition limit.

Take quiz
Dynamic partitioning takes keys from:
the record contents, using jq or Lambda
the shard ID only
the consumer's IAM role
Dynamic partitioning is supported for:
Splunk
any destination
OpenSearch Service
Amazon S3 destinations

34. What is the difference between Kinesis Data Streams and SQS?

Kinesis Data Streams is a replayable log, while SQS is a queue where each message is removed once it has been handled.

Kinesis Data Streams Amazon SQS
Consumption Many consumers read the same records Each message is processed by one consumer, then deleted
Replay Yes, within retention No
Ordering Per shard / partition key Only in FIFO queues, per message group
Retention 24 hours to 365 days 1 minute to 14 days
Scaling Shards or on-demand Automatic
Acknowledgement Consumer checkpoints a position Per-message delete

Pick Kinesis for ordered event streams with several readers, and SQS for simple work queues with independent retries per message.

Take quiz
Which one lets several independent applications read the same records?
an SQS standard queue
Kinesis Data Streams
neither of them
In SQS, a message is removed when:
the consumer deletes it after processing
the shard splits
the retention period is increased

35. What is the difference between Kinesis Data Streams and Apache Kafka?

Both are partitioned, replayable logs. Kinesis is an AWS-native service where you scale shards, while Kafka (or Amazon MSK) is open source and gives you more control and portability.

Kinesis Data Streams Apache Kafka / MSK
Unit of parallelism Shard Partition
Operations Serverless, no brokers Brokers to size, or MSK Serverless
Consumer coordination KCL leases in DynamoDB, or Lambda Consumer groups
Retention Up to 365 days Configurable, plus tiered storage
Ecosystem Deep AWS integrations Kafka Connect, Kafka Streams, many vendors
Portability AWS only Runs anywhere

Choose Kinesis for low operations effort inside AWS. Choose Kafka when you need its ecosystem, very high per-topic throughput tuning, or a multi-cloud option.

Take quiz
The Kafka equivalent of a Kinesis shard is a:
broker
topic
offset
partition
Which option is more portable across cloud providers?
Kinesis, because it is an open standard
Kafka, because it is open source
neither one

36. How do you encrypt data in Kinesis Data Streams?

For data at rest, enable server-side encryption with AWS KMS. Kinesis encrypts the data payload before storing it and decrypts it when consumers read.

aws kinesis start-stream-encryption --stream-name orders --encryption-type KMS --key-id alias/aws/kinesis

You can use the AWS managed key or your own customer managed key. With a customer managed key, producers need kms:GenerateDataKey and consumers need kms:Decrypt.

In transit, all calls use HTTPS/TLS. If you need end-to-end protection, you can also encrypt on the client before calling PutRecord, though then Kinesis can't help consumers like Firehose inspect the content.

Disabling or deleting a customer managed key makes reads and writes fail, so protect it with a key policy and don't schedule its deletion casually. Very busy streams should also watch KMS request quotas.

Take quiz
With SSE using a customer managed key, a producer needs:
kms:GenerateDataKey on that key
kms:CreateKey
kms:ScheduleKeyDeletion
Encryption in transit for Kinesis is provided by:
KMS alone
VPC peering alone
HTTPS/TLS on the service endpoints

37. How do you restrict access to a stream using IAM?

Grant the minimum actions on the specific stream ARN. A producer role usually needs only the put actions.

{
  "Effect": "Allow",
  "Action": ["kinesis:PutRecord", "kinesis:PutRecords"],
  "Resource": "arn:aws:kinesis:us-east-1:123456789012:stream/orders"
}

A consumer needs GetRecords, GetShardIterator, DescribeStreamSummary, and ListShards. For enhanced fan-out it also needs SubscribeToShard on the consumer ARN.

To keep traffic off the internet, use an interface VPC endpoint (AWS PrivateLink). Resource-based policies on the stream allow cross-account access without assuming a role.

You can tighten access further with the aws:SourceVpce condition so the stream only accepts calls arriving through your VPC endpoint. For Lambda consumers, the managed policy AWSLambdaKinesisExecutionRole covers the read permissions.

Take quiz
A producer role needs at least:
kinesis:DeleteStream
kinesis:PutRecord and kinesis:PutRecords on the stream ARN
kinesis:SplitShard
Private access from a VPC without using the internet uses:
an interface VPC endpoint (AWS PrivateLink)
an internet gateway only
VPC peering with S3

38. How do you replay old data from a Kinesis stream?

Replay is possible as long as the records are still inside the retention period. Retention is 24 hours by default, and you can raise it up to 365 days (8,760 hours) at extra cost.

aws kinesis increase-stream-retention-period --stream-name orders --retention-period-hours 168

To reread, start a consumer from TRIM_HORIZON or AT_TIMESTAMP. The KCL setting initialPositionInStream only applies when no checkpoint exists, so to reprocess you start the app under a new application name or delete its lease table. A Lambda mapping takes its StartingPosition when the mapping is created.

Retention beyond 24 hours is billed extra. For anything longer-term or very large, land a copy in S3 through Firehose rather than keeping everything in the stream.

Take quiz
The maximum retention period of a data stream is:
7 days
30 days
365 days
unlimited
To reprocess history with the KCL you can:
call GetRecords with a larger Limit
run with a new application name and TRIM_HORIZON as the starting position
split the shard

39. Why can Kinesis consumers receive duplicate records?

Kinesis offers at-least-once delivery, so duplicates can happen on both the producer and consumer side.

  • Producer retries: a write times out but actually succeeded, so the retry stores the record twice.
  • Consumer restarts: a worker fails after processing but before checkpointing, so the next worker rereads the batch.
  • Failover and resharding: lease changes can replay records since the last checkpoint.
  • Lambda retries: a failed batch is retried in full.

Treat duplicates as normal. Make processing idempotent instead of trying to prevent them.

Firehose can also deliver duplicates to S3 when a delivery is retried, so tables built on those files may need dedupe by a unique ID during query or ETL.

Take quiz
A producer-side duplicate typically happens when:
a write times out but succeeded, and the producer retries
the shard is encrypted
the stream has too many shards
The delivery guarantee of Kinesis Data Streams is:
exactly-once end to end
at-most-once
at-least-once

40. When would you choose on-demand over provisioned capacity mode?

Choose on-demand when traffic is spiky, unpredictable, or new, and you'd rather not manage shards. Choose provisioned when throughput is steady and predictable and you want to optimize cost.

Situation Better fit
New workload, unknown volume On-demand
Large daily swings or sudden bursts On-demand
Flat, high, predictable throughput Provisioned, usually cheaper
Need exact control of shard count and resharding Provisioned

On-demand starts at 4 MiB/s of write capacity and scales to up to double the previous peak seen in the last 30 days, with a default ceiling of 200 MiB/s for writes that you can raise. A single hot partition key can still hit the 1 MiB/s per-shard limit, so key design still matters.

Take quiz
On-demand mode is the best fit for:
flat, perfectly predictable high volume where cost is key
spiky, unpredictable traffic
a stream that must have exactly 7 shards
On-demand scaling is based on:
observed throughput, up to double the previous peak
a schedule you set in CloudWatch
the number of consumers

41. How do you monitor a Kinesis data stream with CloudWatch?

Stream-level metrics are published to CloudWatch for free. The ones worth alarming on are in the table.

Metric What it tells you
GetRecords.IteratorAgeMilliseconds How far consumers lag behind; rising means falling behind
WriteProvisionedThroughputExceeded Producers are being throttled
ReadProvisionedThroughputExceeded Consumers are being throttled
PutRecords.Success Share of successful writes
IncomingBytes / IncomingRecords Ingest volume versus capacity

To find a hot shard, turn on enhanced (shard-level) monitoring with EnableEnhancedMonitoring. It adds per-shard metrics and costs extra. KPL and KCL also publish their own metrics.

Sensible alarms are iterator age above a few minutes for five periods, any sustained write throttling, and IncomingBytes above about 80% of provisioned capacity. Metrics arrive at one-minute granularity, and shard-level metrics are billed like custom metrics.

Take quiz
Which metric shows how far consumers are behind?
IncomingBytes
PutRecords.Success
GetRecords.IteratorAgeMilliseconds
To spot a hot shard you enable:
KMS encryption
enhanced shard-level monitoring
extended retention

42. How does Kinesis Video Streams ingest and play back video?

A camera or device uses the Producer SDK (or a GStreamer plugin) to send media with PutMedia over HTTPS. Kinesis Video Streams stores it as fragments indexed by timestamp, encrypted with KMS, for the retention period you set.

For playback, GetHLSStreamingSessionURL and GetDASHStreamingSessionURL return URLs for live or archived video, and GetClip exports an MP4. Applications can read raw fragments with GetMedia and feed them to Rekognition Video or SageMaker.

WebRTC is a separate path using signaling channels for peer-to-peer, low-latency two-way video.

flowchart LR
  A["Camera with Producer SDK"] -->|PutMedia| B["Video stream"]
  B --> C["Fragments stored by timestamp"]
  C --> D["HLS or DASH playback"]
  C --> E["GetMedia consumers"]
  E --> F["Rekognition Video or SageMaker"]
Take quiz
Which API returns an HLS playback URL?
PutMedia
GetShardIterator
GetHLSStreamingSessionURL
SubscribeToShard
Video in Kinesis Video Streams is stored as:
shards selected by partition key
rows in a DynamoDB table
fragments indexed by timestamp

43. What happens when you split or merge a shard?

Both operations are resharding: the old shards are closed and new child shards take over their hash key ranges.

  • SplitShard: one parent becomes two children. You choose a NewStartingHashKey, often the midpoint, to divide the range.
  • MergeShards: two adjacent shards become one child covering both ranges.
  • UpdateShardCount: scales the whole stream uniformly, within limits such as doubling or halving per call.

Closed parent shards stop accepting writes but stay readable until their data ages out of retention. Children record their parent IDs, which is how consumers know the right reading order. Producers need no changes, because new writes route to the children automatically.

A reshard takes seconds to minutes, during which the stream status is UPDATING but reads and writes continue. Don't hard-code shard IDs because they change; use ListShards. Children hold only new writes, since parent data is not copied across. In on-demand mode these APIs aren't available.

Take quiz
After SplitShard, the parent shard is:
deleted immediately with its data
closed but still readable until retention expires
still accepting writes
MergeShards requires the two shards to be:
adjacent in hash key range
in different Regions
in different streams

44. How do you preserve ordering during resharding?

A consumer must finish reading the parent shard before reading its children. Records for a given partition key sat in the parent before the reshard and in the child afterwards, so draining the parent first keeps them in order.

The KCL does this automatically: it won't take a child lease until the parent is fully processed. Lambda handles it too. If you wrote a custom consumer, check each shard's ParentShardId and keep reading the parent until GetRecords returns a null NextShardIterator, which means the shard is closed and drained.

flowchart TD
  A["Parent shard closed"] --> B["Read parent until NextShardIterator is null"]
  B --> C["Parent fully drained"]
  C --> D["Start reading child shards"]

Example: key A wrote records 1-100 to shard 1, which is then split, and records 101 onward go to a child. A consumer that reads the child before finishing shard 1 could process record 101 before record 99. Merges follow the same rule: both parents must be drained before the child is read. Lambda, the KCL, and the Flink Kinesis connector all handle this for you.

Take quiz
Before reading a child shard, a consumer should:
start with the child immediately for speed
delete the parent shard
finish reading the parent shard
A closed, fully read shard signals the end by returning:
an empty partition key
a null NextShardIterator
a negative sequence number

45. How do you handle a hot shard?

First confirm it: enable shard-level metrics and look for one shard with far higher IncomingBytes or throttling than the rest. The usual cause is a skewed partition key.

  • Improve the key: use a higher-cardinality key.
  • Salt the key: append a small random suffix so one busy key spreads across shards.
  • Split the shard: works when the load comes from many different keys.
  • Aggregate first: batch with the KPL to cut record counts.
import random
key = f"{device_id}-{random.randint(0, 9)}"

Splitting can't help if a single key exceeds 1 MiB/s, since that key still maps to one shard. Salting fixes that but gives up strict ordering across the salts, and consumers must recombine the results.

The salt width matters. N suffixes spread one key over up to N shards, so choose N from the traffic the key needs divided by 1 MiB/s, and have consumers recombine results for that key. On-demand mode doesn't remove hot shards, since a single shard's limit still applies.

Take quiz
Splitting a shard will NOT help when:
one partition key alone exceeds a shard's limit
the keys are evenly spread
the stream uses KMS encryption
Salting a partition key trades away:
encryption at rest
the retention period
strict per-key ordering across the salted keys

46. How do you troubleshoot high iterator age?

GetRecords.IteratorAgeMilliseconds is the age of the last record read. If it keeps growing, consumers are falling behind, and if it passes the retention period, unread records are lost.

  1. Check whether consumers are throttled (ReadProvisionedThroughputExceeded) or just slow.
  2. Measure processing time per batch and fix slow downstream calls.
  3. For Lambda, look for errors that retry the same batch, and raise the parallelization factor.
  4. For KCL, check the DynamoDB lease table for throttling or leases stuck on dead workers.
  5. If several consumers share throughput, move them to enhanced fan-out.
  6. Add shards if the stream itself is undersized.

Alarm on iterator age well before it reaches about half your retention period.

The pattern helps diagnosis. Rising age on every shard suggests the consumer is too slow overall, rising age on one shard suggests a poison record or hot shard, and a jump right after a deployment points to a code regression.

Take quiz
If iterator age exceeds the retention period:
the shard splits automatically
unprocessed records expire and are lost
producers are blocked
A good first check for high iterator age is:
whether consumer processing is slow, failing, or throttled
the length of the partition key
the KMS key policy

47. How do you handle poison pill records in a Kinesis consumer?

A poison pill is a record your code can never process successfully. Because a shard is read in order and a failed batch is retried, one bad record can block everything behind it until it expires.

In your own consumer, wrap processing in a try/catch, send bad records to a dead-letter location such as SQS or S3, and checkpoint past them. For Lambda, use the built-in controls:

  • BisectBatchOnFunctionError splits failing batches to isolate the bad record.
  • MaximumRetryAttempts and MaximumRecordAgeInSeconds stop endless retries.
  • An on-failure destination keeps the skipped records' metadata.
  • ReportBatchItemFailures retries from the failing record only.

Validating schema at the producer reduces how many bad records get in at all.

Log the shard ID, sequence number, and error for each skipped record so it can be repaired later, and alert on dead-letter queue depth so skipped data doesn't go unnoticed. The stream still holds the record until retention expires, so a manual replay is possible.

Take quiz
Why is a poison pill dangerous on Kinesis?
it deletes the shard
it increases the shard count
ordered processing retries it and blocks later records in the shard
BisectBatchOnFunctionError helps by:
doubling the batch size
splitting a failing batch to isolate the bad record
encrypting the batch

48. How do you build idempotent processing on Kinesis?

Because delivery is at-least-once, design so that processing the same record twice has the same effect as once.

  1. Give each logical event a unique ID at the producer, and reuse it on retries rather than generating a new one per attempt.
  2. In the consumer, write with a conditional or upsert operation keyed by that ID.
  3. Keep a dedupe table with a TTL longer than your retention period.
  4. As a fallback key, use the shard ID plus sequence number, which is unique per record.
try:
    table.put_item(Item={"event_id": eid, "payload": body},
                   ConditionExpression="attribute_not_exists(event_id)")
except table.meta.client.exceptions.ConditionalCheckFailedException:
    pass  # already processed

Managed Service for Apache Flink can give exactly-once state updates through checkpoints, but external sinks still need idempotent writes.

For counters and aggregates, store the last processed sequence number per shard in the same atomic transaction as the result, for example with DynamoDB TransactWriteItems. A redelivered batch is then recognized and skipped.

Take quiz
Where should a dedupe ID be generated?
at the producer, once per event, and reused on retries
freshly on each retry attempt
randomly by the consumer
A DynamoDB put with attribute_not_exists(event_id) is used to:
increase read capacity
rotate the KMS key
reject a second write of the same event ID

49. Explain the lifecycle of a record from producer to consumer?

A record passes through a handful of stages between being written and being deleted.

  1. The producer calls PutRecord(s) with data and a partition key.
  2. Kinesis hashes the key with MD5 and routes the record to the owning shard.
  3. It assigns a sequence number, encrypts if SSE is on, and replicates across three AZs.
  4. The record is readable within roughly 70 ms (enhanced fan-out) or about 200 ms (shared).
  5. A consumer gets a shard iterator or subscription, reads, processes, and checkpoints.
  6. The record stays until retention expires, then is deleted whether read or not.
sequenceDiagram
  participant P as Producer
  participant S as Shard
  participant C as Consumer
  P->>S: PutRecord with partition key
  S-->>P: ShardId and SequenceNumber
  C->>S: GetRecords or SubscribeToShard
  S-->>C: Records in order
  C->>C: Process and checkpoint
  Note over S: Deleted after retention expires
Take quiz
What happens right after the partition key is hashed?
it is sent to Firehose
the record is routed to the shard owning that hash range
the consumer is notified through SNS
When are records removed from a stream?
when they age past the retention period
as soon as one consumer reads them
when the shard iterator expires

50. Explain the internal working of KPL aggregation and collection?

The KPL applies two stages before data reaches Kinesis. Aggregation packs many small user records into one larger Kinesis record. Collection then batches several of those records into a single PutRecords call.

A native daemon sends requests asynchronously and flushes a buffer when it reaches a size or count threshold or when RecordMaxBufferedTime (100 ms by default) expires. Failed records are retried until RecordTtl runs out.

Setting Default
RecordMaxBufferedTime 100 ms
AggregationMaxSize 51,200 bytes
CollectionMaxCount 500 records
CollectionMaxSize 5 MiB

Since the 1,000 records/s shard limit counts Kinesis records, aggregation lets you move far more user records per shard and lowers PUT payload cost. Consumers must deaggregate; the KCL does so automatically, while Lambda needs the KPL deaggregation library.

flowchart LR
  A["User records"] --> B["Aggregation into one Kinesis record"]
  B --> C["Collection into PutRecords batch"]
  C --> D["Async send with retries"]
  D --> E["Stream shard"]
Take quiz
Aggregation lets the KPL exceed 1,000 user records/s per shard because:
the KPL disables the shard limit
it skips partition key hashing
the limit counts Kinesis records, and many user records fit in one
The default RecordMaxBufferedTime is about:
100 ms
5 seconds
1 minute
0 ms
«
»

Comments & Discussions