Kinesis Streams, Firehose og analyse i realtid
Indlæs streamingdata med Kinesis Data Streams, lever dem til S3 eller Redshift med Firehose, og analysér dem i realtid med Managed Service for Apache Flink
Kinesis Streams, Firehose og analyse i realtid er en gratis Cloud & IT Cert Prep-lektion på CoddyKit. Dette er lektion 4 af 4. Du kan læse hele lektionen gratis nedenfor — og derefter øve dig praktisk i browseren med en indbygget kodeeditor og en AI-vejleder, der er tilgængelig døgnet rundt. Den er en del af læringsforløbet i Cloud & IT Cert Prep, og dine fremskridt synkroniseres på tværs af nettet og CoddyKit-appen. Cloud & IT Cert Prep-kurset indeholder 4 lektioner i alt.
Oversigt over Kinesis-familien
Amazon Kinesis er en familie af tjenester til indsamling, behandling og analyse af data i realtidsstreams. De tre centrale tjenester er: Kinesis Data Streams (lav latenstid og brugerdefineret behandling), Kinesis Data Firehose (fuldt administreret levering til S3/Redshift/OpenSearch) og Managed Service for Apache Flink (tidligere Kinesis Data Analytics) til realtids-SQL og Flink-behandling. Hver tjeneste er målrettet et forskelligt punkt i streaming-pipelinen.
Arkitektur for Kinesis Data Streams
En Kinesis Data Stream er en holdbar log i rækkefølge, som er opdelt i shards. Hver shard leverer en skrivekapacitet på 1 MB/s og en læsekapacitet på 2 MB/s. Datarecords opbevares som standard i 24 timer (kan udvides til 7 eller 365 dage). Producenter skriver records til en stream, mens forbrugere — Lambda, KCL-applikationer, Firehose eller Flink — læser fra en eller flere shards parallelt. Records kan ikke ændres, når de først er skrevet.
# 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, gennemløb og skalering
Antallet af shards bestemmer en streams samlede gennemløb. Du kan opdele en shard for at fordoble gennemløbet eller flette to shards for at reducere omkostningerne. Brug Enhanced Fan-Out til at give hver registreret forbruger sin egen læsekapacitet på 2 MB/s, uafhængigt af andre forbrugere, så læsebegrænsning undgås, når flere applikationer bruger den samme stream. Overvåg GetRecords.IteratorAgeMilliseconds for at registrere forsinkelser hos forbrugere.
# 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-appKinesis Data Firehose: Administreret levering
Kinesis Data Firehose er en fuldt administreret tjeneste, der indsamler, transformerer og leverer streamingdata til destinationer som S3, Amazon Redshift, Amazon OpenSearch Service, Splunk og HTTP-slutpunkter. Der er ingen shards, du skal administrere — Firehose skalerer automatisk. Du konfigurerer en bufferstørrelse (1-128 MB) og et bufferinterval (60-900 sekunder); Firehose leverer, så snart en af grænserne nås først.
# 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"
}'Datatransformation i Firehose med Lambda
Firehose kan kalde en Lambda-funktion for hver batch af records før levering for at transformere, berige eller filtrere data undervejs. Almindelige anvendelser omfatter konvertering af JSON til Parquet (via Glue-skema), maskering af PII-felter eller fjernelse af hændelser med lav værdi. Records, der ikke kan transformeres, sendes valgfrit til et separat S3-fejlpræfiks til genbehandling, så ingen data går tabt.
# 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 (tidligere Kinesis Data Analytics) kører Apache Flink-applikationer på fuldt administreret infrastruktur. Brug den til tilstandsholdende realtidsanalyse: aggregering i glidende vinduer, registrering af afvigelser, mønstergenkendelse i hændelsesforløb og sammenføjning af streamingdata med referencetabeller. Du skriver Flink-kode i Java, Python eller Scala, og Flink administrerer checkpoints og tilstand med præcis én behandling.
# 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);Valg mellem Streams og Firehose
Til SAA-C03-eksamen skal du kende beslutningskriterierne. Brug Kinesis Data Streams, når du har brug for latenstid under ét sekund, flere forbrugere, der læser samtidigt, eller brugerdefineret behandlingslogik med fuld kontrol over opbevaring og afspilning. Brug Firehose, når du blot har brug for pålidelig levering af streamingdata til S3, Redshift eller OpenSearch med minimal kode, valgfri transformation undervejs og automatisk skalering ved højere latenstid (60+ sekunder).
Kinesis eller SQS: Det klassiske eksamensvalg
Et almindeligt SAA-C03-spørgsmål beder dig vælge mellem Kinesis og SQS. De vigtigste forskelle er: Kinesis bevarer meddelelsesrækkefølgen inden for en shard, understøtter flere forbrugere, der læser de samme data samtidigt, og opbevarer records til afspilning. SQS fjerner meddelelser, efter de er blevet forbrugt (ingen afspilning), FIFO-tilstand garanterer streng rækkefølge, og SQS er bedre til afkobling af mikrotjenester. Hvis scenariet nævner realtidsanalyse eller afspilning, skal du vælge Kinesis.
Kinesis-producenter: SDK og KPL
Kinesis Producer Library (KPL) er en klient med højt gennemløb til skrivning til Kinesis Data Streams fra applikationer. KPL samler automatisk flere små records i ét API-kald (op til 1 MB) og håndterer gentagelser med eksponentiel tilbageventen. Det reducerer omkostningen pr. PUT betydeligt og øger gennemløbet pr. shard. Brug KPL til producenter med store datamængder, f.eks. webklikstreams, IoT-sensorer eller log-pipelines.
# 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() { ... });Mønsteret Firehose til Redshift
En almindelig arkitektur er at bruge Firehose som en administreret pipeline fra Kinesis-streams eller direkte producenter ind i Amazon Redshift til datalagring. Firehose skriver først data til en mellemliggende S3-stagingbucket og udfører derefter en COPY-kommando for at indlæse dataene i Redshift. Det er den mest effektive måde at masseindlæse streamingdata i Redshift på — direkte indsættelse række for række i Redshift ville være ekstremt langsom på grund af overhead for hver række.
# 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"
}
}'Overvågning af Kinesis med CloudWatch
Vigtige CloudWatch-metrikker for Kinesis er: IncomingBytes og IncomingRecords til måling af producenternes gennemløb, GetRecords.IteratorAgeMilliseconds til måling af forbrugerforsinkelse (en høj værdi betyder, at forbrugerne ikke kan følge med) og WriteProvisionedThroughputExceeded til registrering af, når producenter rammer shard-grænserne. Opret CloudWatch-alarmer for iteratoralder og overskredet tilgængeligt gennemløb for at udløse automatisk skalering af shards via Application Auto Scaling.
# 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-alertsHurtigt tjek
Test din forståelse af AWS Solutions Architect (SAA-C03)-koncepterne fra denne lektion.
Opsummering af lektionen
I denne lektion lærte du, at Kinesis Data Streams leverer holdbar streaming i rækkefølge og med shards samt understøtter flere forbrugere og afspilning, at Kinesis Data Firehose tilbyder fuldt administreret levering uden kode til S3, Redshift og OpenSearch, og at Managed Service for Apache Flink muliggør tilstandsholdende realtidsanalyse af streams. Dernæst ser vi nærmere på de 7 R'er i strategien for cloudmigrering.
Lær Cloud & IT Cert Prep med en AI-underviser — gratis
Skriv og kør rigtig kode i din browser, få øjeblikkelig hjælp fra en AI-underviser døgnet rundt, og fortsæt, hvor du slap, på web eller i appen.
- Kurser
- 150
- Lektioner
- 600
Ofte stillede spørgsmål
Er lektionen “Kinesis Streams, Firehose og analyse i realtid” gratis?
Ja — hele teksten til “Kinesis Streams, Firehose og analyse i realtid” kan læses gratis her på nettet. Hvis du vil øve dig interaktivt med en indbygget kodeeditor og en AI-vejleder døgnet rundt og få adgang til resten af Cloud & IT Cert Prep-kurset, skal du opgradere til CoddyKit PRO. Cloud & IT Cert Prep-kurset indeholder 4 lektioner i alt.
Hvad lærer jeg i “Kinesis Streams, Firehose og analyse i realtid”?
Indlæs streamingdata med Kinesis Data Streams, lever dem til S3 eller Redshift med Firehose, og analysér dem i realtid med Managed Service for Apache Flink Du øver dig i Cloud & IT Cert Prep med praktisk kode, som du kører direkte i browseren, og en AI-vejleder døgnet rundt besvarer dine spørgsmål, mens du arbejder dig gennem lektionen.
Skal jeg have erfaring for at begynde på Cloud & IT Cert Prep?
Der kræves ingen tidligere erfaring. Cloud & IT Cert Prep på CoddyKit er tilrettelagt for både begyndere og øvede, så du kan starte her eller fra begyndelsen og lære i dit eget tempo. Dette er lektion 4 af 4.
Hvor lang tid tager lektionen “Kinesis Streams, Firehose og analyse i realtid”?
De fleste CoddyKit-lektioner tager cirka 5–10 minutter. Hver lektion er kort og interaktiv, så du gør løbende fremskridt og kan fortsætte, hvor du slap – på både web og app.
Kan jeg skrive og køre kode i denne Cloud & IT Cert Prep-lektion?
Ja. Alle Cloud & IT Cert Prep-lektioner har en indbygget kodeeditor, så du kan skrive og køre rigtig kode direkte i din browser og få øjeblikkelig feedback fra AI – uden lokal opsætning.
Alle lektioner i dette kursus
- Opbygning af en datasø på S3
- AWS Glue: ETL og datakatalog
- Amazon Athena: Serverløs SQL på S3
- Kinesis Streams, Firehose og analyse i realtid