AWS 303: Kinesis Data Streams, Firehose and Managed Flink
Why this lesson matters
“Real time” does not define an architecture. A payment alert might need ordered events for each card within one second, replay for seven days, duplicate-safe processing, and a durable audit copy. A log archive may accept five-minute batches and need no custom consumer. A fraud model may need event-time windows, state, late-event handling, and a recoverable checkpoint. These are three different problems.
Amazon Kinesis Data Streams is durable streaming storage for custom producers and consumers. Amazon Data Firehose is managed buffered delivery to supported destinations. AWS renamed Kinesis Data Firehose to Amazon Data Firehose in 2024, so older diagrams and APIs may still contain “firehose” or Kinesis-era terminology. Amazon Managed Service for Apache Flink runs stateful Apache Flink applications. It was previously called Kinesis Data Analytics for Apache Flink, which explains the kinesisanalyticsv2 API name.
This workshop teaches event and delivery semantics before service menus. You will design a clickstream and fraud pipeline, calculate capacity, test partition-key distribution, select consumer modes, define Firehose buffering and failure retention, model Flink event time and state, plan replay and recovery, diagnose supplied failures, and prove cleanup. The required track makes no cloud changes.
What you will be able to do
By the end, you can:
- distinguish event time, ingestion time, processing time, latency, throughput, lag, and freshness;
- state ordering, duplication, replay, retention, and loss requirements precisely;
- choose provisioned or on-demand Kinesis Data Streams capacity from workload evidence;
- design partition keys that preserve required order without creating hot shards;
- compare shared-throughput polling, enhanced fan-out, KCL, Lambda, Firehose, and Flink consumers;
- explain producer retries, partial batch failure, consumer checkpoints, and duplicate handling;
- configure Firehose source, buffering, transformation, dynamic partitioning, backup, and destination retry behavior;
- model Flink operators, keyed state, event-time watermarks, windows, checkpoints, snapshots, parallelism, and KPUs;
- separate Flink exactly-once state recovery from end-to-end sink semantics;
- map IAM, KMS, VPC, endpoint, destination, and cross-account controls;
- diagnose throttling, hot keys, iterator age, late data, backpressure, failed delivery, and checkpoint stalls; and
- produce a capacity, cost, recovery, and operational dossier.
Before you start
- Complete AWS300 through AWS302 for S3, Catalog, Lake Formation, and distributed processing foundations.
- Use synthetic records only. Event streams commonly contain identifiers and behavioral data.
- The example Region is
ap-south-1. Verify exact service, destination, connector, and runtime availability. - The optional lab needs an approved budget, one short-lived stream, a KMS choice, log destination, and named cleanup owner.
- Do not send test records to an unknown stream or delivery stream.
- Never assume a successful producer response proves final destination delivery.
- Record current prices and quotas immediately before approving production capacity.
Create a local evidence directory:
mkdir -p "$HOME/aws303-evidence"
cd "$HOME/aws303-evidence"
export AWS_DEFAULT_REGION="ap-south-1"
aws sts get-caller-identity --query Arn --output text
aws configure list
Redact account-specific identity before sharing evidence.
1. Define streaming semantics before selecting services
business event
|
+--> event time: when the business action happened
+--> ingestion time: when the stream accepted it
+--> processing time: when an operator handled it
|
v
partitioned durable stream --> one or more independent consumers
| |
| +--> state, windows, alerts, APIs
v
buffered delivery --> durable analytical destination and failure copy
| Requirement | Question that must be measurable |
|---|---|
| Throughput | Records and bytes per second at average, peak, and burst |
| Latency | Which percentile from event creation to which outcome? |
| Ordering | Global, per entity, per partition key, or none? |
| Durability | How long must an accepted event remain replayable? |
| Duplication | Can producers, consumers, delivery, or retries duplicate it? |
| Loss | Which failure boundary can lose data, and how is loss detected? |
| Replay | Who may reset position, how far, and into what safe sink? |
| Late data | How late can events arrive and still change a window? |
| Backpressure | What happens when processing is slower than ingestion? |
| Recovery | What RPO/RTO applies to state and final destinations? |
Kinesis Data Streams ordering is scoped to records mapped to a shard, commonly through the same partition key. Multiple producers, resharding, retries, and downstream concurrency must still be considered. No AWS service can invent a meaningful event ID or business sequence if the producer did not create one.
Use an immutable envelope:
{
"event_id": "01J-example-unique-id",
"event_type": "page_viewed",
"event_version": 1,
"entity_id": "session-7f2",
"event_time": "2026-09-22T06:00:00Z",
"producer": "web-edge",
"trace_id": "redacted-example",
"payload": {"page": "/training/aws"}
}
Consumers should deduplicate by event_id or an equivalent business key for a bounded period. Sequence numbers are stream positions, not universal business IDs.
2. Understand Kinesis Data Streams storage and capacity
A data stream contains shards. A producer supplies a partition key; Kinesis hashes it into a shard's hash-key range. Each record receives a sequence number. Retention provides a replay window, not permanent archival.
Provisioned mode exposes shard count and resharding. A shard currently supports up to 1 MiB/s or 1,000 records/s for writes and 2 MiB/s for shared reads, subject to documented API limits. On-demand mode manages stream capacity, but a single partition key is still constrained by a shard's per-key path. On-demand is not immunity from hot keys, sudden beyond-profile growth, service quotas, or poor consumer design.
Capacity worksheet:
write shards >= max(
ceiling(peak records per second / 1000),
ceiling(peak bytes per second / 1 MiB)
)
shared read demand =
retained bytes per second x number of independently polling consumers
Then add growth, uneven distribution, retry bursts, aggregation effects, and reshard lead time. Do not use an average day to size a launch peak.
| Mode | Strength | Risk | Evidence |
|---|---|---|---|
| On-demand | Less capacity administration; adapts to traffic | Cost and scaling assumptions can be vague | observed peaks, hot-key test, quota and price |
| Provisioned | Explicit capacity and predictable shard model | Reshard operations and forecasting | shard math, alarms, scale runbook |
Retention can be extended when replay value justifies cost. Archive independently to S3 when legal, analytical, or long-term recovery requirements exceed stream retention.
3. Design partition keys with real distributions
A partition key should align the smallest ordering domain with enough cardinality for load distribution. customer_id may preserve customer order but fail if one customer dominates. A random UUID distributes well but destroys natural per-session ordering. Time-only keys create extreme hot spots.
Build a local test file:
printf '%s\n' \
'tenant-a,session-1,3500' \
'tenant-a,session-2,120' \
'tenant-b,session-3,90' \
'tenant-c,session-4,70' \
'tenant-d,session-5,40' \
> producer-profile.csv
awk -F, '{sum += $3; if ($3 > max) {max=$3; key=$1}} END {
printf "total_rps=%d hottest_key=%s hottest_rps=%d hottest_share=%.2f%%\n",
sum,key,max,(100*max/sum)
}' producer-profile.csv | tee hot-key-evidence.txt
For each candidate key, calculate hottest-key records and bytes per second, required ordering, consumer grouping, and rebalance behavior. If salting a hot entity into entity#bucket, document how the consumer reconstructs order or why cross-bucket order is unnecessary.
KPL aggregation can pack multiple user records into one Kinesis record to reduce request overhead. Consumers must deaggregate correctly. Aggregation changes request economics, not the logical need for event IDs and duplicate safety.
4. Build duplicate-safe producers
PutRecord writes one record. PutRecords batches records, but the response can contain partial failures. Retrying the entire batch can duplicate successes. Retry only failed entries with exponential backoff and jitter, while preserving original event IDs.
create immutable event ID
|
serialize and validate size/schema
|
choose partition key from ordering contract
|
PutRecord or PutRecords
|
inspect every result entry
| |
success failure
record evidence classify, bounded retry, dead-letter evidence
Producer acceptance proves the stream accepted the record. It does not prove a consumer processed it, Firehose delivered it, Flink checkpointed it, or a sink committed it. Carry correlation IDs through every layer.
Protect producers against:
- records larger than the current maximum;
- unbounded retry queues exhausting memory;
- credentials expiring during backlog drain;
- schema-breaking deployment;
- clock drift corrupting event-time logic;
- a single hot key;
- retry storms during throttling; and
- local buffering loss during process termination.
5. Choose consumer type and checkpoint semantics
Shared-throughput consumers poll shards through GetRecords and share each shard's read capacity. Enhanced fan-out registers consumers and provides dedicated read throughput with HTTP/2 push-style delivery, at additional cost. Use it when multiple low-latency consumers or read contention justify it.
The Kinesis Client Library coordinates workers, shard leases, checkpoints, failures, and resharding, commonly using DynamoDB for lease state. A checkpoint means “resume after this sequence position” for that consumer application. It does not prove the external side effect committed atomically.
Lambda event source mappings poll streams, batch records, retry failures, and expose controls such as batch size, age, bisect behavior, partial batch response, parallelization, and failure destination. Increasing per-shard parallelization can weaken assumptions about processing order. Design idempotent handlers.
| Consumer | Use when | Main concern |
|---|---|---|
| KCL application | Custom long-running logic and lease control | checkpoint versus side-effect atomicity |
| Lambda | Short event-driven handling | retries, poison records, concurrency, downstream limits |
| Firehose | Managed destination delivery | buffering, retry duration, failed-data handling |
| Managed Flink | Stateful event-time computation | state, checkpoint, parallelism, sink semantics |
| Custom API consumer | Specialized protocol/control | must own resharding, checkpoints, scaling |
Monitor iterator age or equivalent lag, read throughput, throttles, expired iterators, lease churn, processing errors, and downstream latency. Lag is an inventory of unfinished work and must have a recovery-time estimate.
6. Use Amazon Data Firehose for managed delivery
Firehose is not a replayable message bus. It accepts direct writes or uses supported sources, buffers records, optionally invokes transformations or format conversion, and delivers batches to a configured destination. It manages scaling, but destination throttling and service quotas still matter.
direct PUT or Kinesis source
|
v
optional Lambda transform and metadata extraction
|
v
buffer by size/time and optional dynamic partition
|
+--> primary destination
|
+--> failed/backup S3 prefix and CloudWatch evidence
Buffer size and interval are hints, not exact delivery deadlines. Low traffic can wait for the time threshold; dynamic partitioning uses multi-stage buffering and can add more delay. Zero buffering is not available with dynamic partitioning. Many partition values can create many active buffers, small objects, quota pressure, and expensive downstream analytics.
For S3, define:
- raw, transformed, backup, and error prefixes;
- UTC time and dynamic business partitions;
- compression and Parquet/ORC conversion;
- Glue schema ownership and compatibility;
- KMS key policy;
- bucket ownership and cross-account delivery;
- lifecycle and incomplete-delivery monitoring; and
- record reconciliation from source acceptance to S3 objects.
Transformation Lambda receives batches. It must return one result per input with the correct record identifier and status. Bound payload size, timeout, concurrency, logging, malformed-data behavior, and retries. Never silently mark an invalid record as delivered.
Destination retry behavior differs by destination. When retries expire, Firehose can place undelivered data in an S3 backup/error path where configured. “Delivery stream ACTIVE” does not prove current records reached the destination. Monitor freshness, bytes, failed records, throttling, transformation errors, and backup growth.
7. Build the Flink state and time model
Flink applications form a directed graph of sources, operators, and sinks. Each operator can have parallel subtasks. Keying a stream partitions records so all events for one key reach the same logical keyed-state owner.
Processing-time windows use machine time and are simple but change results when lag changes. Event-time windows use timestamps from records and watermarks that estimate how far event time has progressed. A watermark is not proof that no older event can arrive. The allowed-lateness and update/retraction policy determine what happens next.
events: 10:01, 10:04, 10:02-late, 10:06
|
timestamp extraction and watermark strategy
|
keyBy(session_id)
|
five-minute event-time window + keyed state
|
on-time result, allowed late update, or late side output
Document:
- timestamp field and invalid-clock handling;
- bounded out-of-orderness;
- idle partitions that could hold back watermarks;
- window type, trigger, accumulation, and late-data policy;
- state retention and timer growth;
- side output or quarantine;
- sink update semantics; and
- how corrected results reach consumers.
Without these choices, “five-minute session window” is incomplete.
8. Separate checkpoints, snapshots, and end-to-end guarantees
Flink checkpoints periodically capture consistent application state and source positions for automatic recovery. Managed Service for Apache Flink snapshots persist application state for stop, update, rollback, or planned recovery. Connector compatibility, state schema, operator identity, maximum parallelism, and application version affect restorability.
Exactly-once state recovery is not automatically exactly-once business delivery. The source must support replay, the checkpoint must include source position, and the sink must support transactional or idempotent commit coordinated with checkpoints. An external HTTP call or ordinary database insert can duplicate after recovery.
Create a failure matrix:
| Failure point | Recovered from | Duplicate risk | Data-loss risk | Test |
|---|---|---|---|---|
| task restart | latest completed checkpoint | sink dependent | after checkpoint if source not replayable | kill task |
| application update | compatible snapshot | sink dependent | incompatible state blocks restore | rollback version |
| source lag beyond retention | no stream data remains | possible during manual repair | high | capacity/retention game day |
| sink outage | backpressure and retry/state | connector dependent | if retry/retention exhausted | deny sink |
| bad event | validation/side output | low if quarantined | high if silently dropped | poison record |
Keep checkpointing enabled for production fault tolerance. Prove snapshot creation and restore before relying on it. An unhealthy application may fail to create a clean stop snapshot.
9. Size Flink parallelism and control backpressure
Managed Service for Apache Flink allocates Kinesis Processing Units. Current documentation describes one KPU as one vCPU, 4 GB memory, and 50 GB disk, with task slots controlled by ParallelismPerKPU. Application parallelism, operator-specific parallelism, max parallelism, source partitions, and sink capacity interact.
Automatic scaling is enabled by default and responds to sustained CPU thresholds. Scaling can cause application downtime and doubles parallelism on scale-up under the documented behavior. CPU-only reaction can be too late for sudden lag or can miss a blocked sink with low CPU. Use source lag, busy/backpressured time, checkpoint duration/failure, records in/out, sink throttling, state size, and watermark delay.
Backpressure is flow control: a slow sink prevents upstream operators from emitting freely. It protects memory only if the whole pipeline and source retention can absorb the delay. Calculate:
backlog bytes = incoming bytes per second x backlog seconds
catch-up time = backlog records / (sustainable processing rps - incoming rps)
If sustainable processing does not exceed incoming rate after recovery, the system never catches up.
10. Map security and network paths
Separate producer identities, consumer identities, Firehose delivery role, transformation role, Flink execution role, human deployer, and CI role. Scope permissions to exact streams, applications, buckets/prefixes, Glue tables, KMS keys, logs, secrets, and destinations.
Draw:
producer -> Kinesis endpoint and KMS
consumer/Flink -> source, checkpoint/state service, sink, logs
Firehose role -> source, Lambda, Glue, KMS, destination, backup
private application -> subnets, SGs, DNS, endpoints/NAT, databases
KMS encryption protects service-managed data at rest, but key policy, grants, rotation, cross-account use, and disabled-key response remain operational dependencies. TLS protects in transit. Private connectivity requires endpoint policy and DNS validation. Never place secret destination credentials in application code or plaintext configuration.
11. Diagnose streaming failures
| Symptom | Evidence | Likely cause | Safe response |
|---|---|---|---|
| Write throttling with low total traffic | per-key metrics/profile | hot partition key | redesign key or salt with order plan |
Partial PutRecords errors | entry-level response | shard/key pressure or transient error | retry failed entries only |
| Iterator age rises | consumer rate, errors, downstream | consumer slower than input | fix bottleneck, then controlled scale |
| Duplicate sink records | event IDs, retries, checkpoints | at-least-once replay | deduplicate/idempotent commit |
| Firehose freshness rises | destination and transform metrics | buffer, throttling, Lambda, destination | inspect exact failed stage |
| S3 has tiny files | buffer/partition cardinality | low volume or excessive dynamic keys | consolidate partitions/buffer design |
| Flink watermark stalls | per-partition activity | idle source partition or timestamp issue | configure idleness/fix timestamps |
| Checkpoints time out | duration, state, backpressure, storage | large state or blocked pipeline | repair bottleneck and tune safely |
| Autoscaling but lag worsens | source/sink limits and backpressure | sink bottleneck or source partition ceiling | scale correct layer |
| Restore rejects snapshot | release, operator IDs, max parallelism | incompatible state evolution | follow tested migration/rollback |
Never reset consumer position, delete a snapshot, increase retention, and scale every component simultaneously. Preserve the evidence and change one layer with a rollback.
12. Design the P18 streaming architecture
Scenario: a retailer emits 8,000 events/s normally and 40,000 events/s during campaigns. Average record size is 1.2 KiB; one celebrity campaign can generate 18 percent of events for one product. Requirements:
- order events ordered per order;
- page views do not require product-global order;
- fraud alert p99 under two seconds;
- five-minute event-time session windows with 10-minute late allowance;
- seven-day replay for incident recovery;
- S3 Parquet archive available within 10 minutes;
- independently deployable fraud, personalization, and archive consumers;
- no silent malformed-record loss; and
- Region failure strategy with declared RPO/RTO.
Produce:
- separate keys for order and page-view event families;
- peak records/bytes and hot-key calculations;
- on-demand versus provisioned KDS ADR;
- shared versus enhanced fan-out per consumer;
- retention and long-term S3 archive;
- Firehose buffer, transform, partition, error, and reconciliation design;
- Flink event-time, watermark, state, checkpoint, snapshot, parallelism, and sink design;
- duplicate-safe event envelope and sink contract;
- backpressure and catch-up model;
- identity, KMS, network, and cross-account paths;
- twelve failure injections; and
- full cost, quota, monitoring, and cleanup register.
Reject Firehose as the fraud processor because buffered delivery is not custom stateful subsecond computation. Reject Flink for a simple untransformed archive if Firehose alone meets the requirement. Consider MSK when Kafka protocol/ecosystem, longer broker semantics, or existing Kafka operations are requirements; AWS304 covers that comparison.
13. Cost, quotas, and cleanup
Kinesis Data Streams cost varies by capacity mode, data in/out, shard hours where applicable, retention, enhanced fan-out, and connected features. Firehose pricing can include ingested bytes, format conversion, dynamic partitioning, VPC delivery, transformation Lambda, destination, backup, and S3. Managed Flink uses KPUs, running application time, snapshots/state, logs, source/sink services, and network paths.
Also include KMS, CloudWatch, Lambda, DynamoDB KCL leases, S3 requests/storage, Glue, NAT, transfer, downstream write capacity, replay, failed retries, non-production environments, and operator labor.
Read-only inventory:
aws kinesis list-streams --output json | tee kds-streams.json
aws firehose list-delivery-streams --output json | tee firehose-streams.json
aws kinesisanalyticsv2 list-applications --output json | tee flink-applications.json
Cleanup proof for owned labs includes deleted Kinesis streams after retention evidence is no longer needed, deleted Firehose streams after buffers finish and backup is reconciled, stopped/deleted Flink applications after a required snapshot, and reviewed S3 objects, Lambda transforms, log groups, KCL lease tables, IAM roles, KMS grants, ENIs, endpoints, alarms, and dashboards. Access denied is not deletion proof. Review delayed billing.
14. Practical submission
Submit a p18-streaming/ package containing:
- event envelope and schema-evolution contract;
- event-time, ordering, duplication, loss, replay, and latency definitions;
- peak throughput and shard/key distribution workbook;
- KDS mode, retention, encryption, and scaling ADR;
- producer partial-failure and retry state machine;
- consumer mode and checkpoint/side-effect matrix;
- Firehose buffering, transformation, dynamic partition, backup, and reconciliation plan;
- Flink operator graph, watermarks, windows, late-data, and state;
- checkpoint/snapshot restore and sink-semantics proof;
- KPU/parallelism/backpressure/catch-up model;
- identity, KMS, private-network, and cross-account map;
- twelve failure-injection runbooks;
- monitoring and alert thresholds tied to business freshness;
- monthly cost and quota sensitivity;
- replay, Region recovery, and rollback runbook; and
- exact cleanup and delayed billing evidence.
Knowledge check
- What does a Kinesis partition key control? Hash-based shard placement and therefore an ordering/load domain.
- Does on-demand eliminate hot keys? No; a single key still follows shard-level constraints.
- Why can
PutRecordsretry duplicate data? A batch can partially succeed. - What does enhanced fan-out provide? Dedicated per-shard read throughput and low-latency consumer delivery.
- Is Firehose a replayable stream? No; it is managed buffered destination delivery.
- Why can dynamic partitioning raise latency? It buffers independently through multiple stages per active partition.
- What does a watermark mean? The processor's estimate of event-time progress, not a guarantee no older event exists.
- Is Flink exactly-once state recovery equal to exactly-once business effects? Only with compatible replayable sources and transactional/idempotent sinks.
- Can autoscaling fix a blocked sink? No; it can add cost while backpressure remains.
- What proves streaming recovery? Restored state/position, reconciled outputs, bounded lag, and no unexplained loss or duplication.
Lesson acceptance
Pass only when all are true:
- event, ingestion, and processing time plus ordering/replay/duplication are defined;
- shard and hot-key calculations use peak records and bytes, not averages;
- provisioned/on-demand and shared/enhanced fan-out decisions have evidence;
- producer partial failures and retries preserve event IDs and avoid whole-batch replay;
- consumer checkpoints are not misrepresented as atomic sink commits;
- Firehose buffering, dynamic partitioning, transformation, backup, and destination retries are explicit;
- Flink watermarks, windows, state, checkpoints, snapshots, parallelism, and sink semantics are testable;
- backpressure includes backlog and catch-up calculations;
- IAM, KMS, network, endpoint, secret, destination, and cross-account paths are mapped;
- twelve failures include hot key, poison event, destination outage, checkpoint failure, and retention overrun;
- monitoring measures business freshness, lag, failure retention, and reconciliation;
- service choice rejects both unnecessary Flink and misuse of Firehose;
- cost includes every connected service, retry, replay, retained state, and environment; and
- cleanup handles streams, applications, buffers, snapshots, data, logs, leases, roles, and network artifacts.
Official sources
- What is Kinesis Data Streams?
- Kinesis Data Streams capacity modes
- Kinesis Data Streams quotas
- Enhanced fan-out consumers
- Kinesis Client Library
- What is Amazon Data Firehose?
- Firehose dynamic partitioning
- Firehose dynamic-partition buffering
- What is Managed Service for Apache Flink?
- Managed Flink application resources
- Managed Flink automatic scaling
- Managed Flink best practices