Menu

AWS course · Lesson 8 of 12

Amazon Kinesis: Data Streams, Firehose and Stream Processing

Amazon Kinesis for Data Engineers: Data Streams, shard throughput, KCL and KPL, enhanced fan-out, Amazon Data Firehose, Managed Service for Apache Flink, and Kafka.

  • Intermediate
  • 16 min read
  • Updated Oct 2026
On this page
  1. Sample code
  2. Kinesis Data Streams
  3. Shards and throughput
  4. KCL and KPL
  5. Enhanced fan-out
  6. Amazon Data Firehose
  7. Kinesis Data Analytics (now Managed Service for Apache Flink)
  8. Kinesis vs Kafka
  9. Practice questions
  10. Key takeaways

Amazon Kinesis is the family of AWS services for streaming data. Kinesis Data Streams is a durable, ordered, replayable log, similar in spirit to Kafka. Amazon Data Firehose (formerly Kinesis Data Firehose) delivers streams into S3, Redshift, OpenSearch and other destinations with no code. Amazon Managed Service for Apache Flink (formerly Kinesis Data Analytics) runs stateful stream processing. This lesson covers how they work, how to size them, and how to choose between Kinesis and Kafka.

Sample code

The Python models reproduce the arithmetic behind Kinesis design decisions (shard sizing, how partition keys map to shards, and Firehose buffering) and run locally. The producer and consumer API calls need an AWS account and are not executed.

Kinesis Data Streams

What it is. A data stream is an ordered log of records, split into shards. A record is a data blob plus a partition key chosen by the producer. Kinesis assigns each record a sequence number, unique and increasing within its shard. Consumers read records in order per shard and can re-read them until they expire.

How it works.

  • Retention is 24 hours by default and can be raised up to 365 days (8,760 hours), at extra cost above 24 hours. Records cannot be deleted individually before they expire.
  • Producers call PutRecord or PutRecords (batches of up to 500 records), or use the KPL, agents, or AWS services that write to streams.
  • Consumers read with GetShardIterator and GetRecords, or more usually through Lambda event source mappings, the KCL, Firehose, or Flink.
  • Capacity modes: in provisioned mode you choose the shard count and reshard (split or merge shards) as load changes; in on-demand mode Kinesis manages shards for you, starting from a default write capacity of 4 MB/s (4,000 records/s) and scaling with observed traffic.
import json
import boto3

kinesis = boto3.client("kinesis")
events = [{"user_id": "u-17", "action": "add_to_cart", "sku": "B-204"},
          {"user_id": "u-42", "action": "checkout", "sku": "A-001"}]
resp = kinesis.put_records(
    StreamName="clickstream",
    Records=[{"Data": json.dumps(e).encode(), "PartitionKey": e["user_id"]} for e in events],
)
if resp["FailedRecordCount"]:
    failed = [events[i] for i, r in enumerate(resp["Records"]) if "ErrorCode" in r]
    # retry only the failed records, with backoff

Pitfalls.

  • PutRecords is not all-or-nothing: some records can fail (throttling) while others succeed. Check FailedRecordCount and retry the failed ones.
  • Retries by the producer can write a record twice, so Kinesis delivery is at least once. Consumers must de-duplicate or be idempotent, for example using an event ID in the payload.
  • Ordering is per shard only. Records with different partition keys can be processed in any relative order.

In interviews. Define the stream, shard, partition key and sequence number, state the default and maximum retention, and say “ordered per shard, at-least-once, replayable within retention”.

Shards and throughput

What it is. A shard is the unit of capacity. Each shard accepts writes of up to 1 MB per second or 1,000 records per second, and supports reads of up to 2 MB per second shared by all standard consumers (and up to five GetRecords calls per second). Enhanced fan-out consumers each get their own 2 MB/s per shard.

How it works. In provisioned mode you size from these limits. The busiest constraint wins: bytes in, records in, or bytes out across all shared consumers:

import math

def shards_needed(records_per_sec, avg_record_kb, consumers_shared=1):
    """Provisioned-mode sizing from the per-shard limits: writes 1 MB/s or 1,000 records/s,
    shared reads 2 MB/s split across all standard consumers."""
    in_mb = records_per_sec * avg_record_kb / 1024
    by_bytes_in = math.ceil(in_mb / 1.0)
    by_records_in = math.ceil(records_per_sec / 1000)
    by_reads = math.ceil(in_mb * consumers_shared / 2.0)
    return max(by_bytes_in, by_records_in, by_reads), (by_bytes_in, by_records_in, by_reads)

for rps, kb, consumers in [(5_000, 0.5, 1), (5_000, 3.0, 1), (5_000, 3.0, 3)]:
    n, parts = shards_needed(rps, kb, consumers)
    print(f"{rps} rec/s x {kb} KB, {consumers} shared consumer(s) -> {n} shards (bytes in, records in, reads) = {parts}")
5000 rec/s x 0.5 KB, 1 shared consumer(s) -> 5 shards (bytes in, records in, reads) = (3, 5, 2)
5000 rec/s x 3.0 KB, 1 shared consumer(s) -> 15 shards (bytes in, records in, reads) = (15, 5, 8)
5000 rec/s x 3.0 KB, 3 shared consumer(s) -> 22 shards (bytes in, records in, reads) = (15, 5, 22)

Small records are limited by the record count, large ones by bytes, and several standard consumers by the shared read limit (the case for enhanced fan-out). Add headroom for peaks.

Partition keys decide the shard. Kinesis hashes each partition key with MD5 into a 128-bit number and sends the record to the shard whose hash key range contains it. A key that is far busier than the others makes a hot shard, which throttles even when the stream as a whole has spare capacity:

import hashlib
from collections import Counter

def shard_for(partition_key, shard_count):
    """Kinesis hashes the partition key with MD5 into a 128-bit key space split into shard ranges."""
    h = int(hashlib.md5(partition_key.encode()).hexdigest(), 16)
    return h * shard_count // 2**128            # equal-sized ranges, like a freshly created stream

events = [f"user-{i % 2000}" for i in range(20_000)] + ["user-bot"] * 6_000   # one very chatty key
by_user = Counter(shard_for(k, 4) for k in events)
print("partition key = user id :", [by_user[s] for s in range(4)])
import random
random.seed(7)
by_random = Counter(shard_for(str(random.random()), 4) for _ in events)
print("partition key = random  :", [by_random[s] for s in range(4)])
partition key = user id : [10970, 5200, 4980, 4850]
partition key = random  : [6462, 6521, 6444, 6573]

Random keys spread load evenly but give up per-entity ordering; entity keys keep each user’s events in order but can be skewed. A common compromise is a composite key (for example user_id plus a small bucket number) for the few known heavy keys.

Pitfalls.

  • ProvisionedThroughputExceededException on writes: a hot key or too few shards. Look at per-shard metrics, not stream totals.
  • Resharding a provisioned stream changes the hash ranges. Parent shards must be fully read before child shards to keep order; the KCL handles this.
  • More shards than consumers can process in parallel waste money; fewer than the load needs throttle producers.

In interviews. Quote the limits exactly (1 MB/s or 1,000 records/s in, 2 MB/s out per shard), size a stream from a given rate, and explain hot shards and how partition key choice trades ordering for balance.

KCL and KPL

What it is. The Kinesis Client Library (KCL) is a consumer library that handles the hard parts of reading a stream with many workers. The Kinesis Producer Library (KPL) is a producer library that maximises write throughput.

How the KCL works.

  • Each shard is a lease stored in a DynamoDB table. Workers in an application share the leases, so each shard is processed by exactly one worker at a time, similar to partitions in a Kafka consumer group.
  • Workers checkpoint the last processed sequence number per shard in DynamoDB. After a restart or failover, processing resumes from the checkpoint, so records after it are processed again: at-least-once.
  • It follows resharding (finishing parent shards before children) and balances leases when workers join or leave. KCL 3.x adds load-based lease balancing; use the current major version for new consumers and check the documentation for the support status of older ones.

How the KPL works. It batches records, aggregates many small user records into one Kinesis record (so the 1,000 records/s limit stops being the bottleneck), retries failures and publishes metrics. Consumers must de-aggregate (the KCL and Lambda’s Kinesis integration with the de-aggregation library do this).

Pitfalls.

  • Checkpointing before processing finishes (losing records on failure) or very rarely (reprocessing a lot after a restart).
  • Two applications sharing one KCL application name, and therefore one lease table, so they steal each other’s leases.
  • The lease table’s DynamoDB capacity becoming a bottleneck for streams with many shards.
  • KPL buffering adds latency (it waits to fill batches); tune the buffering time for latency-sensitive paths.

In interviews. Explain leases and checkpoints in DynamoDB as Kinesis’s equivalent of consumer groups and committed offsets, and KPL aggregation as the way around the per-shard record limit.

Enhanced fan-out

What it is. By default all consumers of a shard share its 2 MB/s read throughput and poll with GetRecords. Enhanced fan-out (EFO) gives each registered consumer its own 2 MB/s per shard, with records pushed over HTTP/2 (SubscribeToShard) instead of polled, which also lowers latency.

How it works. You register a consumer on the stream; the KCL, Lambda (with an EFO consumer ARN) and Flink can use it. Each EFO consumer costs extra per consumer-shard hour and per GB retrieved, and the number of registered consumers per stream is limited (see the quotas page).

Shared throughput Enhanced fan-out
Read throughput 2 MB/s per shard, shared by all consumers 2 MB/s per shard per consumer
Delivery Polling GetRecords (5 calls/s per shard) Push over HTTP/2
Latency Higher, grows with the number of consumers Lower and consistent
Cost Included Extra per consumer-shard hour and data retrieved

Pitfalls. Adding EFO for a single consumer that is not limited by read throughput: it adds cost without benefit.

In interviews. Use EFO when several independent applications read the same stream, or when latency matters; otherwise shared throughput is fine.

Amazon Data Firehose

What it is. Amazon Data Firehose (named Kinesis Data Firehose until early 2024) is a fully managed delivery service. Producers send records (directly, or from a Kinesis data stream or MSK topic) and Firehose buffers, optionally transforms, and delivers them to S3, Redshift, OpenSearch, Splunk, Snowflake, Apache Iceberg tables, HTTP endpoints and other destinations. There are no shards to manage.

How it works.

  • Buffering hints: Firehose delivers when the buffer reaches a size or a time interval, whichever comes first. If you do not set them, the S3 defaults are 5 MB or 300 seconds (5 minutes).
  • Record format conversion turns JSON into Parquet or ORC using a schema from the Glue Data Catalog.
  • Transformation with a Lambda function (enrich, filter, reshape records).
  • Dynamic partitioning writes to S3 prefixes built from fields in the record (for example dt= and customer_id=), extracted with jq expressions or a Lambda.
  • Failed records go to an error prefix in S3 so nothing is silently dropped.

This model shows the “size or interval, whichever first” rule for a busy and a quiet stream:

def firehose_flushes(arrivals, buffer_mb=5, buffer_seconds=300):
    """Firehose delivers when the buffer reaches its size or its interval, whichever comes first.
    arrivals: list of (second, megabytes)."""
    flushes, size, opened = [], 0.0, None
    for t, mb in arrivals:
        if opened is not None and t - opened >= buffer_seconds:
            flushes.append((opened + buffer_seconds, round(size, 1), "interval"))
            size, opened = 0.0, None
        if opened is None:
            opened = t
        size += mb
        if size >= buffer_mb:
            flushes.append((t, round(size, 1), "size"))
            size, opened = 0.0, None
    return flushes

busy = [(t, 0.5) for t in range(0, 60, 2)]          # 0.25 MB/s for a minute
quiet = [(t, 0.01) for t in range(0, 900, 30)]       # a trickle for 15 minutes
print("busy :", firehose_flushes(busy)[:3])
print("quiet:", firehose_flushes(quiet))
busy : [(18, 5.0, 'size'), (38, 5.0, 'size'), (58, 5.0, 'size')]
quiet: [(300, 0.1, 'interval'), (600, 0.1, 'interval')]

(The quiet stream’s third buffer would flush at 900 seconds, after the simulated window.) A quiet stream produces tiny files on every interval, which is the small files problem; raise the buffer interval, or compact downstream.

aws firehose create-delivery-stream --delivery-stream-name clicks-to-lake \
  --delivery-stream-type KinesisStreamAsSource \
  --kinesis-stream-source-configuration KinesisStreamARN=arn:aws:kinesis:eu-west-1:111122223333:stream/clickstream,RoleARN=arn:aws:iam::111122223333:role/firehose-read \
  --extended-s3-destination-configuration file://s3-destination.json

Pitfalls.

  • Expecting exactly-once delivery: Firehose can deliver duplicates on retries. Deduplicate downstream (or use an Iceberg destination with keys).
  • Dynamic partitioning on a high-cardinality field creates many small files and hits the limit on active partitions per stream.
  • Firehose is for delivery, not for joins, windows or stateful logic. Use Flink for those.

In interviews. Firehose is the answer to “land this stream in S3 as Parquet, partitioned by date, with no code”. Mention buffering hints, format conversion, dynamic partitioning and the error prefix.

What it is. Kinesis Data Analytics was renamed Amazon Managed Service for Apache Flink on 30 August 2023. It runs Apache Flink applications (Java, Scala, Python, or SQL through Flink Studio notebooks) on managed infrastructure for stateful stream processing: event-time windows, joins between streams, aggregations, pattern detection, and exactly-once state.

How it works. You upload a Flink application (a JAR or Python package), configure parallelism in Kinesis Processing Units (KPUs), and the service runs it with checkpoints and snapshots for recovery. Sources and sinks include Kinesis Data Streams, MSK, S3, DynamoDB and others. The rename did not change APIs, endpoints or IAM actions, so the CLI is still aws kinesisanalyticsv2.

The legacy SQL product, Kinesis Data Analytics for SQL Applications, has been discontinued: no new applications from 15 October 2025, and applications deleted from 27 January 2026. Its replacement is Managed Service for Apache Flink or Flink Studio.

aws kinesisanalyticsv2 create-application --application-name sessionise-clicks \
  --runtime-environment FLINK-1_20 \
  --service-execution-role arn:aws:iam::111122223333:role/flink-app \
  --application-configuration file://flink-app-config.json

Check the documentation for the Flink runtime versions currently supported before choosing one.

Pitfalls.

  • Under-provisioned parallelism makes back-pressure and growing consumer lag; watch the application’s lag and checkpoint metrics.
  • Exactly-once inside Flink does not make external sinks exactly-once; use transactional or idempotent sinks.

In interviews. Name the service correctly (Managed Service for Apache Flink), state that the SQL-based Kinesis Data Analytics has been discontinued, and describe when you need Flink: stateful, event-time processing that Lambda or Firehose cannot do.

Kinesis vs Kafka

What it is. Both are partitioned, replayable logs. The choice is mostly about operations, ecosystem and integration.

Kinesis Data Streams Apache Kafka (self-managed or Amazon MSK)
Unit of parallelism Shard Partition
Throughput per unit Fixed: 1 MB/s or 1,000 records/s in, 2 MB/s out Depends on brokers and hardware
Retention 24 hours default, up to 365 days Configurable, unlimited with tiered storage
Ordering Per shard Per partition
Consumer coordination KCL leases in DynamoDB; Lambda ESM Consumer groups and committed offsets
Exactly-once Not built in; idempotent consumers Idempotent producers and transactions within Kafka
Operations Fully managed, on-demand mode available MSK manages brokers (or MSK Serverless); you still size clusters and topics
Ecosystem Tight AWS integration (Lambda, Firehose, Flink) Kafka Connect, Kafka Streams, Schema Registry, multi-cloud
Pricing model Shard-hours or on-demand throughput Broker hours and storage (MSK), or your own servers

When to choose which. Kinesis for AWS-native teams that want minimal operations and direct Lambda and Firehose integration. Kafka (on MSK) when the organisation already uses Kafka tooling, needs Kafka Connect or Kafka Streams, needs longer retention or transactions, or wants portability.

In interviews. Map the concepts one to one (shard and partition, lease and offset), then argue the choice from operations, ecosystem and cost rather than raw speed.

Practice questions

You must ingest 8,000 events per second of about 2 KB each, read by two independent applications. How many shards, and would you use enhanced fan-out?

Inbound is about 15.6 MB/s, so at least 16 shards by bytes (8 by record count). Two shared consumers read about 31 MB/s in total, which needs 16 shards at 2 MB/s each, so 16 shards cover it with no headroom; add margin for peaks, or use on-demand mode. Enhanced fan-out would give each application its own 2 MB/s per shard and lower latency, which is worth it if the consumers are latency-sensitive or more consumers will be added.

One shard is throttling while the others are idle. Why, and what do you do?

A hot partition key: most records hash to that shard’s range. Spread the heavy keys (add a suffix bucket to known hot keys, or use a more granular key), accept losing strict per-entity ordering where possible, or split that shard. Monitor per-shard metrics to confirm.

How does the KCL guarantee each shard is processed by one worker, and what happens on failure?

Each shard is a lease in a DynamoDB table; a worker holds a lease and renews it. If the worker dies, its lease expires and another worker takes it, resuming from the last checkpointed sequence number. Records processed after that checkpoint are processed again, so processing is at least once and must be idempotent.

What is the difference between Kinesis Data Streams and Amazon Data Firehose?

Data Streams is a durable, replayable log that you read with your own consumers; you manage capacity (shards or on-demand). Firehose is a managed delivery pipe: it buffers records and writes them to destinations such as S3 or Redshift, with optional Lambda transformation, Parquet conversion and dynamic partitioning, and it does not let you replay or read records yourself.

What happened to Kinesis Data Analytics?

It was renamed Amazon Managed Service for Apache Flink in August 2023, with no API changes. The older SQL-based product, Kinesis Data Analytics for SQL Applications, was discontinued: new applications could not be created from 15 October 2025 and applications were deleted from 27 January 2026. New stream processing should use Flink applications or Flink Studio.

Key takeaways

  • Kinesis Data Streams is an ordered, at-least-once, replayable log of shards; retention is 24 hours by default and up to 365 days.
  • Each shard takes 1 MB/s or 1,000 records/s in and gives 2 MB/s out; partition keys hash with MD5 to shards, so hot keys make hot shards.
  • The KCL coordinates consumers with DynamoDB leases and checkpoints; the KPL batches and aggregates writes.
  • Enhanced fan-out gives each consumer its own 2 MB/s per shard with push delivery.
  • Amazon Data Firehose buffers (by default 5 MB or 300 seconds for S3) and delivers streams with no code.
  • Kinesis Data Analytics is now Managed Service for Apache Flink; its SQL predecessor is discontinued.

By Data Career Hub Editorial · Last reviewed Oct 2026 · Limits and service names checked against the Kinesis Data Streams, Amazon Data Firehose and Managed Service for Apache Flink documentation in October 2026. The shard-sizing, partition-hashing and Firehose buffering models run locally on Python 3.11. AWS CLI and boto3 calls were written from the documentation and not executed.

Progress is saved in this browser only. No account needed.

Search
Filter by type