Lesson 303 · AWS Learning Path

AWS 303: Kinesis Data Streams, Firehose and Managed Flink

· Published · 16 min read

Labelled process diagram for AWS 303: Producers and partition key to Durable stream to Delivery or stateful processing to Destination, checkpoint, replay, and failure evidence, with decision, proof and rejection...

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
RequirementQuestion that must be measurable
ThroughputRecords and bytes per second at average, peak, and burst
LatencyWhich percentile from event creation to which outcome?
OrderingGlobal, per entity, per partition key, or none?
DurabilityHow long must an accepted event remain replayable?
DuplicationCan producers, consumers, delivery, or retries duplicate it?
LossWhich failure boundary can lose data, and how is loss detected?
ReplayWho may reset position, how far, and into what safe sink?
Late dataHow late can events arrive and still change a window?
BackpressureWhat happens when processing is slower than ingestion?
RecoveryWhat 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.

ModeStrengthRiskEvidence
On-demandLess capacity administration; adapts to trafficCost and scaling assumptions can be vagueobserved peaks, hot-key test, quota and price
ProvisionedExplicit capacity and predictable shard modelReshard operations and forecastingshard 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.

ConsumerUse whenMain concern
KCL applicationCustom long-running logic and lease controlcheckpoint versus side-effect atomicity
LambdaShort event-driven handlingretries, poison records, concurrency, downstream limits
FirehoseManaged destination deliverybuffering, retry duration, failed-data handling
Managed FlinkStateful event-time computationstate, checkpoint, parallelism, sink semantics
Custom API consumerSpecialized protocol/controlmust 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.

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 pointRecovered fromDuplicate riskData-loss riskTest
task restartlatest completed checkpointsink dependentafter checkpoint if source not replayablekill task
application updatecompatible snapshotsink dependentincompatible state blocks restorerollback version
source lag beyond retentionno stream data remainspossible during manual repairhighcapacity/retention game day
sink outagebackpressure and retry/stateconnector dependentif retry/retention exhausteddeny sink
bad eventvalidation/side outputlow if quarantinedhigh if silently droppedpoison 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.

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

SymptomEvidenceLikely causeSafe response
Write throttling with low total trafficper-key metrics/profilehot partition keyredesign key or salt with order plan
Partial PutRecords errorsentry-level responseshard/key pressure or transient errorretry failed entries only
Iterator age risesconsumer rate, errors, downstreamconsumer slower than inputfix bottleneck, then controlled scale
Duplicate sink recordsevent IDs, retries, checkpointsat-least-once replaydeduplicate/idempotent commit
Firehose freshness risesdestination and transform metricsbuffer, throttling, Lambda, destinationinspect exact failed stage
S3 has tiny filesbuffer/partition cardinalitylow volume or excessive dynamic keysconsolidate partitions/buffer design
Flink watermark stallsper-partition activityidle source partition or timestamp issueconfigure idleness/fix timestamps
Checkpoints time outduration, state, backpressure, storagelarge state or blocked pipelinerepair bottleneck and tune safely
Autoscaling but lag worsenssource/sink limits and backpressuresink bottleneck or source partition ceilingscale correct layer
Restore rejects snapshotrelease, operator IDs, max parallelismincompatible state evolutionfollow 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:

  1. separate keys for order and page-view event families;
  2. peak records/bytes and hot-key calculations;
  3. on-demand versus provisioned KDS ADR;
  4. shared versus enhanced fan-out per consumer;
  5. retention and long-term S3 archive;
  6. Firehose buffer, transform, partition, error, and reconciliation design;
  7. Flink event-time, watermark, state, checkpoint, snapshot, parallelism, and sink design;
  8. duplicate-safe event envelope and sink contract;
  9. backpressure and catch-up model;
  10. identity, KMS, network, and cross-account paths;
  11. twelve failure injections; and
  12. 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:

  1. event envelope and schema-evolution contract;
  2. event-time, ordering, duplication, loss, replay, and latency definitions;
  3. peak throughput and shard/key distribution workbook;
  4. KDS mode, retention, encryption, and scaling ADR;
  5. producer partial-failure and retry state machine;
  6. consumer mode and checkpoint/side-effect matrix;
  7. Firehose buffering, transformation, dynamic partition, backup, and reconciliation plan;
  8. Flink operator graph, watermarks, windows, late-data, and state;
  9. checkpoint/snapshot restore and sink-semantics proof;
  10. KPU/parallelism/backpressure/catch-up model;
  11. identity, KMS, private-network, and cross-account map;
  12. twelve failure-injection runbooks;
  13. monitoring and alert thresholds tied to business freshness;
  14. monthly cost and quota sensitivity;
  15. replay, Region recovery, and rollback runbook; and
  16. exact cleanup and delayed billing evidence.

Knowledge check

  1. What does a Kinesis partition key control? Hash-based shard placement and therefore an ordering/load domain.
  2. Does on-demand eliminate hot keys? No; a single key still follows shard-level constraints.
  3. Why can PutRecords retry duplicate data? A batch can partially succeed.
  4. What does enhanced fan-out provide? Dedicated per-shard read throughput and low-latency consumer delivery.
  5. Is Firehose a replayable stream? No; it is managed buffered destination delivery.
  6. Why can dynamic partitioning raise latency? It buffers independently through multiple stages per active partition.
  7. What does a watermark mean? The processor's estimate of event-time progress, not a guarantee no older event exists.
  8. Is Flink exactly-once state recovery equal to exactly-once business effects? Only with compatible replayable sources and transactional/idempotent sinks.
  9. Can autoscaling fix a blocked sink? No; it can add cost while backpressure remains.
  10. 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

Advertisement