AWS Solutions Architect · Les

Kinesis Streams, Firehose en realtimeanalyse

Neem streaminggegevens op met Kinesis Data Streams, lever ze met Firehose aan S3 of Redshift en analyseer ze in realtime met Managed Service for Apache Flink.

Les 4 van 413 stappen

Kinesis Streams, Firehose en realtimeanalyse is een gratis AWS Solutions Architect-les op CoddyKit. Dit is les 4 van 4. Je kunt 3 lessen uit dit leerpad gratis volledig lezen — daarna ontgrendelt CoddyKit PRO alle lessen, plus praktische oefeningen met een ingebouwde code-editor en een AI-tutor die 24/7 beschikbaar is. Deze les maakt deel uit van het leertraject AWS Solutions Architect. Je voortgang wordt gesynchroniseerd op het web en in de CoddyKit-app. De cursus AWS Solutions Architect bevat in totaal 4 lessen.

Overzicht van de Kinesis-familie

Amazon Kinesis is een familie van services voor het verzamelen, verwerken en analyseren van realtime streaminggegevens. De drie belangrijkste services zijn: Kinesis Data Streams (lage latentie en aangepaste verwerking), Kinesis Data Firehose (volledig beheerde levering aan S3/Redshift/OpenSearch) en Managed Service for Apache Flink (voorheen Kinesis Data Analytics) voor realtime SQL- en Flink-verwerking. Elke service richt zich op een ander onderdeel van de streamingverwerkingsketen.

Architectuur van Kinesis Data Streams

Een Kinesis Data Stream is een duurzame, geordende log die is opgedeeld in shards. Elke shard biedt een schrijfcapaciteit van 1 MB/s en een leescapaciteit van 2 MB/s. Gegevensrecords worden standaard 24 uur bewaard (uit te breiden naar 7 of 365 dagen). Producenten schrijven records naar een stream; consumenten — Lambda, KCL-toepassingen, Firehose of Flink — lezen parallel uit een of meer shards. Records kunnen na het schrijven niet meer worden gewijzigd.

# Create a Kinesis Data Stream with 4 shards
aws kinesis create-stream \
  --stream-name clickstream \
  --shard-count 4

# Put a record into the stream
aws kinesis put-record \
  --stream-name clickstream \
  --partition-key 'user-123' \
  --data 'eyJldmVudCI6ICJjbGljayJ9'

Shards, verwerkingscapaciteit en schalen

Het aantal shards bepaalt de totale verwerkingscapaciteit van een stream. Je kunt een shard splitsen om de verwerkingscapaciteit te verdubbelen of twee shards samenvoegen om kosten te verlagen. Gebruik Enhanced Fan-Out om elke geregistreerde consument onafhankelijk van andere consumenten een eigen leescapaciteit van 2 MB/s te geven. Zo voorkom je leesbeperking wanneer meerdere toepassingen dezelfde stream gebruiken. Houd GetRecords.IteratorAgeMilliseconds in de gaten om achterstand bij consumenten te detecteren.

# Split shard to increase throughput
aws kinesis split-shard \
  --stream-name clickstream \
  --shard-to-split shardId-000000000001 \
  --new-starting-hash-key 170141183460469231731687303715884105728

# Register an enhanced fan-out consumer
aws kinesis register-stream-consumer \
  --stream-arn arn:aws:kinesis:us-east-1:123456789012:stream/clickstream \
  --consumer-name analytics-app

Kinesis Data Firehose: beheerde levering

Kinesis Data Firehose is een volledig beheerde service die streaminggegevens vastlegt, transformeert en levert aan bestemmingen zoals S3, Amazon Redshift, Amazon OpenSearch Service, Splunk en HTTP-eindpunten. Je hoeft geen shards te beheren: Firehose schaalt automatisch. Je configureert een buffergrootte (1–128 MB) en een bufferinterval (60–900 seconden); Firehose levert zodra een van beide limieten als eerste wordt bereikt.

# Create a Firehose delivery stream to S3
aws firehose create-delivery-stream \
  --delivery-stream-name clickstream-to-s3 \
  --s3-destination-configuration '{
    "RoleARN": "arn:aws:iam::123456789012:role/FirehoseRole",
    "BucketARN": "arn:aws:s3:::my-data-lake-123",
    "Prefix": "landing/clickstream/year=!{timestamp:yyyy}/month=!{timestamp:MM}/",
    "BufferingHints": {"SizeInMBs": 64, "IntervalInSeconds": 300},
    "CompressionFormat": "GZIP"
  }'

Gegevenstransformatie in Firehose met Lambda

Firehose kan vóór de levering voor elke batch records een Lambda-functie aanroepen om de gegevens tijdens de verwerking te transformeren, verrijken of filteren. Veelvoorkomende toepassingen zijn JSON converteren naar Parquet (via een Glue-schema), PII-velden maskeren of gebeurtenissen met weinig waarde verwijderen. Records waarvoor de transformatie mislukt, kunnen optioneel naar een afzonderlijk S3-foutprefix worden gestuurd voor herverwerking, zodat er geen gegevens verloren gaan.

# Lambda transform function signature for Firehose
def lambda_handler(event, context):
    output = []
    for record in event['records']:
        import base64, json
        payload = json.loads(base64.b64decode(record['data']))
        # Drop events with no user_id
        if not payload.get('user_id'):
            output.append({'recordId': record['recordId'], 'result': 'Dropped', 'data': record['data']})
        else:
            output.append({'recordId': record['recordId'], 'result': 'Ok', 'data': record['data']})
    return {'records': output}

Managed Service for Apache Flink

Amazon Managed Service for Apache Flink (voorheen Kinesis Data Analytics) voert Apache Flink-toepassingen uit op volledig beheerde infrastructuur. Gebruik deze service voor stateful realtime-analyses: aggregaties met schuivende vensters, anomaliedetectie, patroonherkenning in reeksen gebeurtenissen en het koppelen van streaminggegevens aan referentietabellen. Je schrijft Flink-code in Java, Python of Scala en Flink beheert controlepunten en de state voor exact-één-keer-verwerking.

# Flink SQL-style tumbling window (conceptual)
# Count page views per URL every 5 minutes
CREATE TABLE clickstream (
  url STRING,
  event_time TIMESTAMP(3),
  WATERMARK FOR event_time AS event_time - INTERVAL '5' SECOND
) WITH ('connector' = 'kinesis', 'stream' = 'clickstream', ...);

SELECT
  url,
  COUNT(*) AS views,
  TUMBLE_START(event_time, INTERVAL '5' MINUTE) AS window_start
FROM clickstream
GROUP BY url, TUMBLE(event_time, INTERVAL '5' MINUTE);

Kiezen tussen Streams en Firehose

Voor het SAA-C03-examen moet je de beslissingscriteria kennen. Gebruik Kinesis Data Streams wanneer je een latentie van minder dan een seconde, meerdere gelijktijdig lezende consumenten of aangepaste verwerkingslogica nodig hebt, met volledige controle over bewaartermijn en opnieuw afspelen. Gebruik Firehose wanneer je streaminggegevens alleen betrouwbaar met minimale code naar S3, Redshift of OpenSearch wilt leveren, met optionele transformatie tijdens de verwerking en automatisch schalen, maar met een hogere latentie (60 seconden of meer).

Kinesis versus SQS: de klassieke examenvraag

Een veelvoorkomende SAA-C03-vraag vraagt je te kiezen tussen Kinesis en SQS. De belangrijkste verschillen: Kinesis behoudt de volgorde binnen een shard, ondersteunt meerdere consumenten die dezelfde gegevens gelijktijdig lezen en bewaart records zodat je ze opnieuw kunt afspelen. SQS verwijdert berichten nadat ze zijn verwerkt (geen opnieuw afspelen), de FIFO-modus garandeert een strikte volgorde en SQS is beter geschikt voor het loskoppelen van microservices. Als het scenario realtime-analyse of opnieuw afspelen vermeldt, kies je Kinesis.

Kinesis-producenten: SDK en KPL

De Kinesis Producer Library (KPL) is een client met hoge verwerkingscapaciteit voor het schrijven naar Kinesis Data Streams vanuit toepassingen. KPL voegt automatisch meerdere kleine records samen in één API-aanroep (tot 1 MB) en verwerkt nieuwe pogingen met oplopende wachttijden. Hierdoor dalen de kosten per record voor PUT-aanroepen aanzienlijk en neemt de verwerkingscapaciteit per shard toe. Gebruik KPL voor producenten met grote volumes, zoals webkliks, IoT-sensoren of logverwerkingsketens.

# Basic KPL usage (Java pseudocode, no backticks)
KinesisProducer producer = new KinesisProducer();
byte[] data = 'hello world'.getBytes();
ByteBuffer buf = ByteBuffer.wrap(data);
// addUserRecord handles aggregation and retry internally
ListenableFuture future = producer.addUserRecord('clickstream', 'partitionKey', buf);
Futures.addCallback(future, new FutureCallback() { ... });

Patroon voor Firehose naar Redshift

Een veelvoorkomende architectuur gebruikt Firehose als beheerde verwerkingsketen van Kinesis-streams of rechtstreekse producenten naar Amazon Redshift voor gegevensopslag en -analyse. Firehose schrijft gegevens eerst naar een tussentijdse S3-stagingbucket en geeft daarna een COPY-opdracht om de gegevens in Redshift te laden. Dit is de efficiëntste manier om streaminggegevens in bulk in Redshift te laden: rechtstreekse invoegingen per rij in Redshift zouden door de overhead per rij extreem traag zijn.

# Firehose Redshift destination (CLI snippet)
--redshift-destination-configuration '{
  "RoleARN": "arn:aws:iam::123456789012:role/FirehoseRole",
  "ClusterJDBCURL": "jdbc:redshift://cluster.xyz.us-east-1.redshift.amazonaws.com:5439/sales",
  "CopyCommand": {
    "DataTableName": "clickevents",
    "CopyOptions": "JSON 'auto'"
  },
  "Username": "firehose_user",
  "Password": "{{resolve:secretsmanager:redshift-pw}}",
  "S3Configuration": {
    "RoleARN": "...",
    "BucketARN": "arn:aws:s3:::firehose-staging"
  }
}'

Kinesis bewaken met CloudWatch

Belangrijke CloudWatch-statistieken voor Kinesis zijn: IncomingBytes en IncomingRecords om de verwerkingscapaciteit van producenten te meten, GetRecords.IteratorAgeMilliseconds om de achterstand bij consumenten te meten (een hoge waarde betekent dat consumenten het niet kunnen bijhouden) en WriteProvisionedThroughputExceeded om te detecteren wanneer producenten de shardlimieten bereiken. Stel CloudWatch-alarmen in voor de iteratorleeftijd en overschreden ingerichte verwerkingscapaciteit, zodat Application Auto Scaling het automatisch schalen van shards kan starten.

# CloudWatch alarm on high consumer lag
aws cloudwatch put-metric-alarm \
  --alarm-name kinesis-high-lag \
  --metric-name GetRecords.IteratorAgeMilliseconds \
  --namespace AWS/Kinesis \
  --dimensions Name=StreamName,Value=clickstream \
  --statistic Maximum \
  --period 60 \
  --threshold 60000 \
  --comparison-operator GreaterThanThreshold \
  --evaluation-periods 3 \
  --alarm-actions arn:aws:sns:us-east-1:123456789012:ops-alerts

Korte controle

Test je begrip van de concepten uit deze les voor AWS Solutions Architect (SAA-C03).

Samenvatting van de les

In deze les heb je geleerd dat Kinesis Data Streams duurzame, geordende streaming met shards, meerdere consumenten en de mogelijkheid om gegevens opnieuw af te spelen biedt, dat Kinesis Data Firehose volledig beheerde levering zonder code aan S3, Redshift en OpenSearch biedt en dat Managed Service for Apache Flink stateful realtime-analyses op streams mogelijk maakt. Hierna bekijken we de 7 R's van de strategie voor cloudmigratie.

Gratis beginnen

Leer AWS Solutions Architect met een AI-tutor — gratis

Schrijf echte code en voer die uit in je browser, krijg direct hulp van een AI-tutor die 24/7 beschikbaar is en ga verder waar je gebleven bent op het web of in de app.

Cursussen
30
Lessen
120

Veelgestelde vragen

Is de les “Kinesis Streams, Firehose en realtimeanalyse” gratis?

Ja — je kunt hier op het web alle 3 lessen van het leerpad AWS Solutions Architect, waaronder “Kinesis Streams, Firehose en realtimeanalyse”, gratis volledig lezen. Daarna ontgrendelt CoddyKit PRO alle lessen, plus interactieve oefeningen met een ingebouwde code-editor en een AI-tutor die 24/7 beschikbaar is. De cursus AWS Solutions Architect bevat in totaal 4 lessen.

Wat leer ik in “Kinesis Streams, Firehose en realtimeanalyse”?

Neem streaminggegevens op met Kinesis Data Streams, lever ze met Firehose aan S3 of Redshift en analyseer ze in realtime met Managed Service for Apache Flink. Je oefent met AWS Solutions Architect door code rechtstreeks in de browser uit te voeren. Een AI-begeleider die 24/7 beschikbaar is beantwoordt je vragen terwijl je de les doorwerkt.

Heb ik ervaring nodig om met AWS Solutions Architect te beginnen?

Ervaring vooraf is niet nodig. AWS Solutions Architect op CoddyKit is opgebouwd voor beginners tot gevorderden, zodat je hier of bij het begin kunt starten en in je eigen tempo kunt leren. Dit is les 4 van 4.

Hoe lang duurt de les “Kinesis Streams, Firehose en realtimeanalyse”?

De meeste lessen van CoddyKit duren ongeveer 5–10 minuten. Elke les is kort en interactief, zodat je gestaag vooruitgaat en op het web en in de app precies verdergaat waar je was gebleven.

Kan ik code schrijven en uitvoeren in deze les over AWS Solutions Architect?

Ja. Elke les over AWS Solutions Architect bevat een ingebouwde code-editor, zodat je rechtstreeks in je browser echte code kunt schrijven en uitvoeren en direct feedback van AI krijgt — lokale installatie is niet nodig.

Alle lessen in deze cursus

  1. Een data lake op S3 bouwen
  2. AWS Glue: ETL en datapcatalogus
  3. Amazon Athena: serverloze SQL op S3
  4. Kinesis Streams, Firehose en realtimeanalyse
← Terug naar AWS Solutions Architect