Cloud / Amazon Kinesis Interview questions
Last updated
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
running a nightly payroll batch job
ingesting clickstream events and reacting within seconds
storing archival backups for ten years
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
Kinesis Data Streams
Kinesis Video Streams
Amazon Data Firehose
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
1 hour
7 days
until the first consumer reads the record
24 hours
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
10 MiB/s or 10,000 records/s
1 MiB/s or 1,000 records/s
2 MiB/s or 500 records/s
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
random shards chosen per record
a separate stream for each key
the same shard
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
the producer
the consumer on first read
Kinesis Data Streams at ingestion
the KMS key
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
Kinesis Agent tailing a log file
a Lambda function triggered by the stream
a KPL process on a web server
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
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
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
staging it in S3 and running a COPY command
inserting each record over JDBC
streaming directly with no staging
Splunk
Amazon OpenSearch Service
Amazon S3
Amazon RDS for MySQL
10. What is Amazon Managed Service for Apache Flink?
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.
Take quiz
Apache Hive
Hadoop MapReduce
Apache Flink
Presto
windowed aggregations and stateful stream joins
storing camera footage
bulk loading S3 with no code
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
SMTP and IMAP
FTP only
SNMP
HLS and MPEG-DASH
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
PutRecordscalls. - 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
packing multiple user records into a single Kinesis record
compressing every record with GZIP
encrypting records with KMS
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
Amazon S3
Amazon DynamoDB
CloudWatch Logs
Amazon ElastiCache
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
the number of consumers only
GET requests only
shard-hours plus PUT payload units
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
100 records
1,000 records
10 records
500 records
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
the newest record only
the oldest record still in the shard
a timestamp you supply
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
to enable KMS encryption
to choose the target shard
so --data is accepted as raw text instead of base64
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
a shard iterator from get_shard_iterator
a consumer ARN
a partition key
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
inside the Kinesis service
on the source servers that produce the log files
only inside Lambda
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
MD5
SHA-256
CRC32
none, keys are matched as text
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
2 MiB/s total, shared by all consumers
2 MiB/s for each consumer
1 MiB/s in total
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.
- Write bandwidth: 10,000 x 2 KB is about 19.5 MiB/s, so 20 shards at 1 MiB/s each.
- Write records: 10,000 / 1,000 = 10 shards. Bandwidth is the bigger constraint.
- Read bandwidth: 3 consumers x 19.5 = 58.6 MiB/s, divided by 2 MiB/s gives 30 shards.
- 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
1
3
5
10
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
retry in a tight loop
drop all records
retry with exponential backoff and jitter
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
Kinesis Data Streams
Firehose
neither of them
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
SQS long polling
HTTP/2 push through SubscribeToShard
S3 event notifications
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
24 hours
10 seconds
up to 5 minutes, then it must be renewed
until the shard is merged
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
after the records have been processed successfully
before processing each batch
never, because KCL guesses the position
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
PutRecordsentries can land after newer ones, so usePutRecordwithSequenceNumberForOrderingwhen 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
across the entire stream
within a shard, which gives order per partition key
only when Firehose is used
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:
MaximumRetryAttemptsandMaximumRecordAgeInSecondsBisectBatchOnFunctionErrorto split a failing batch- An on-failure destination (SQS, SNS, or S3) for discarded batches
ReportBatchItemFailuresto return the failing sequence number
{"batchItemFailures": [{"itemIdentifier": "49590338271490256608559692538361571095921575989136588898"}]}
Take quiz
fully parallel per record
randomized across batches
sequential and in order
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
5
100
10
2
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
both thresholds are reached
either the buffer size or interval is reached first
only when the buffer interval ends
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
the Kinesis stream's metadata
an S3 bucket policy
an AWS Glue Data Catalog table
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
the record contents, using jq or Lambda
the shard ID only
the consumer's IAM role
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
an SQS standard queue
Kinesis Data Streams
neither of them
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
broker
topic
offset
partition
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
kms:GenerateDataKey on that key
kms:CreateKey
kms:ScheduleKeyDeletion
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
kinesis:DeleteStream
kinesis:PutRecord and kinesis:PutRecords on the stream ARN
kinesis:SplitShard
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
7 days
30 days
365 days
unlimited
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 write times out but succeeded, and the producer retries
the shard is encrypted
the stream has too many shards
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
flat, perfectly predictable high volume where cost is key
spiky, unpredictable traffic
a stream that must have exactly 7 shards
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
IncomingBytes
PutRecords.Success
GetRecords.IteratorAgeMilliseconds
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
PutMedia
GetShardIterator
GetHLSStreamingSessionURL
SubscribeToShard
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
deleted immediately with its data
closed but still readable until retention expires
still accepting writes
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
start with the child immediately for speed
delete the parent shard
finish reading the parent shard
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
one partition key alone exceeds a shard's limit
the keys are evenly spread
the stream uses KMS encryption
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.
- Check whether consumers are throttled (
ReadProvisionedThroughputExceeded) or just slow. - Measure processing time per batch and fix slow downstream calls.
- For Lambda, look for errors that retry the same batch, and raise the parallelization factor.
- For KCL, check the DynamoDB lease table for throttling or leases stuck on dead workers.
- If several consumers share throughput, move them to enhanced fan-out.
- 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
the shard splits automatically
unprocessed records expire and are lost
producers are blocked
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:
BisectBatchOnFunctionErrorsplits failing batches to isolate the bad record.MaximumRetryAttemptsandMaximumRecordAgeInSecondsstop endless retries.- An on-failure destination keeps the skipped records' metadata.
ReportBatchItemFailuresretries 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
it deletes the shard
it increases the shard count
ordered processing retries it and blocks later records in the shard
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.
- Give each logical event a unique ID at the producer, and reuse it on retries rather than generating a new one per attempt.
- In the consumer, write with a conditional or upsert operation keyed by that ID.
- Keep a dedupe table with a TTL longer than your retention period.
- 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
at the producer, once per event, and reused on retries
freshly on each retry attempt
randomly by the consumer
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.
- The producer calls
PutRecord(s)with data and a partition key. - Kinesis hashes the key with MD5 and routes the record to the owning shard.
- It assigns a sequence number, encrypts if SSE is on, and replicates across three AZs.
- The record is readable within roughly 70 ms (enhanced fan-out) or about 200 ms (shared).
- A consumer gets a shard iterator or subscription, reads, processes, and checkpoints.
- 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
it is sent to Firehose
the record is routed to the shard owning that hash range
the consumer is notified through SNS
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"]