Cloud & IT Cert Prep · leksjon

Kinesis Streams, Firehose og sanntidsanalyse

Innhent datastrømmer med Kinesis Data Streams, lever dem til S3 eller Redshift med Firehose, og analyser dem i sanntid med Managed Service for Apache Flink

Leksjon 4 av 413 trinn

Kinesis Streams, Firehose og sanntidsanalyse er en gratis leksjon i Cloud & IT Cert Prep på CoddyKit. Dette er leksjon 4 av 4. Du kan lese hele leksjonen gratis nedenfor – og deretter øve praktisk i nettleseren med en innebygd kodeeditor og en AI-veileder som er tilgjengelig døgnet rundt. Den er en del av læringsløpet i Cloud & IT Cert Prep, og fremdriften din synkroniseres mellom nettet og CoddyKit-appen. Kurset i Cloud & IT Cert Prep inneholder totalt 4 leksjoner.

Oversikt over Kinesis-familien

Amazon Kinesis er en familie av tjenester for innsamling, behandling og analyse av datastrømmer i sanntid. De tre viktigste tjenestene er: Kinesis Data Streams (lav forsinkelse og egendefinert behandling), Kinesis Data Firehose (fullt administrert levering til S3/Redshift/OpenSearch) og Managed Service for Apache Flink (tidligere Kinesis Data Analytics) for sanntids-SQL og Flink-behandling. Hver tjeneste er rettet mot et forskjellig punkt i strømmingsprosessen.

Arkitektur for Kinesis Data Streams

En Kinesis Data Stream er en holdbar, sortert logg som er partisjonert i shards. Hver shard gir 1 MB/s skrivegjennomstrømming og 2 MB/s lesegjennomstrømming. Dataoppføringer beholdes som standard i 24 timer (kan utvides til 7 eller 365 dager). Produsenter skriver oppføringer til en strøm, mens konsumenter – Lambda, KCL-applikasjoner, Firehose eller Flink – leser fra én eller flere shards parallelt. Oppføringer kan ikke endres etter at de 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, gjennomstrømming og skalering

Antallet shards bestemmer den totale gjennomstrømmingen i en strøm. De kan dele en shard for å doble gjennomstrømmingen eller slå sammen to shards for å redusere kostnadene. Bruk Enhanced Fan-Out for å gi hver registrerte konsument sin egen lesegjennomstrømming på 2 MB/s, uavhengig av andre konsumenter. Dette eliminerer lesebegrensning når flere applikasjoner bruker den samme strømmen. Overvåk GetRecords.IteratorAgeMilliseconds for å oppdage forsinkelser hos konsumentene.

# 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: administrert levering

Kinesis Data Firehose er en fullt administrert tjeneste som samler inn, transformerer og leverer datastrømmer til blant annet S3, Amazon Redshift, Amazon OpenSearch Service, Splunk og HTTP-endepunkter. Det finnes ingen shards som må administreres – Firehose skalerer automatisk. De konfigurerer en bufferstørrelse (1–128 MB) og et bufferintervall (60–900 sekunder). Firehose leverer dataene så snart én av grensene nås.

# 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"
  }'

Datatransformasjon i Firehose med Lambda

Firehose kan starte en Lambda function for hver gruppe med oppføringer før levering, slik at dataene kan transformeres, berikes eller filtreres underveis. Vanlige bruksområder er å konvertere JSON til Parquet (via Glue-skjema), maskere PII-felt eller forkaste hendelser med lav verdi. Oppføringer som ikke kan transformeres, kan valgfritt sendes til et separat S3-feilprefiks for ny behandling, slik at ingen data går tapt.

# 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) kjører Apache Flink-applikasjoner på fullt administrert infrastruktur. Bruk tjenesten til tilstandsbevarende sanntidsanalyse, for eksempel aggregasjoner over glidende vinduer, avviksdeteksjon, mønstergjenkjenning i hendelsessekvenser og sammenføyning av datastrømmer med referansetabeller. De skriver Flink-kode i Java, Python eller Scala, mens Flink håndterer kontrollpunkter og tilstand med exactly-once-semantikk.

# 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);

Velge mellom Streams og Firehose

Til SAA-C03-eksamen må De kjenne beslutningskriteriene. Bruk Kinesis Data Streams når De trenger forsinkelse på under ett sekund, flere konsumenter som leser samtidig, eller egendefinert behandlingslogikk med full kontroll over oppbevaring og avspilling. Bruk Firehose når De bare trenger å levere datastrømmer pålitelig til S3, Redshift eller OpenSearch, med minimalt med kode, valgfri transformasjon underveis og automatisk skalering, men med høyere forsinkelse (60 sekunder eller mer).

Kinesis kontra SQS: Det klassiske eksamensvalget

Et vanlig SAA-C03-spørsmål ber Dem velge mellom Kinesis og SQS. Viktige forskjeller er: Kinesis bevarer meldingsrekkefølgen innenfor en shard, støtter flere konsumenter som leser de samme dataene samtidig, og beholder oppføringer for avspilling. SQS fjerner meldinger etter at de er konsumert (ingen avspilling), FIFO-modus garanterer streng rekkefølge, og tjenesten egner seg bedre til frakobling av mikrotjenester. Hvis scenarioet nevner sanntidsanalyse eller avspilling, velger De Kinesis.

Kinesis-produsenter: SDK og KPL

Kinesis Producer Library (KPL) er en klient med høy gjennomstrømming for skriving til Kinesis Data Streams fra applikasjoner. KPL aggregerer automatisk flere små oppføringer i ett API-kall (opptil 1 MB) og håndterer nye forsøk med backoff. Dette reduserer kostnaden per PUT-oppføring betydelig og øker gjennomstrømmingen per shard. Bruk KPL for produsenter med stort volum, for eksempel nettbaserte klikkstrømmer, IoT-sensorer eller loggkanaler.

# 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 vanlig arkitektur er å bruke Firehose som en administrert datapipeline fra Kinesis-strømmer eller direkte fra produsenter til Amazon Redshift for datavarehus. Firehose skriver først dataene til en mellomliggende S3-oppsamlingsbøtte og utfører deretter en COPY-kommando for å laste dataene inn i Redshift. Dette er den mest effektive måten å masseinnlaste strømmende data i Redshift på – direkte innsetting rad for rad i Redshift ville vært svært tregt på grunn av kostnaden for hver enkelt rad.

# 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åke Kinesis med CloudWatch

Viktige CloudWatch-metrikker for Kinesis er: IncomingBytes og IncomingRecords for å måle produsentens gjennomstrømming, GetRecords.IteratorAgeMilliseconds for å måle forsinkelse hos konsumentene (en høy verdi betyr at konsumentene ikke klarer å holde tritt), og WriteProvisionedThroughputExceeded for å oppdage når produsenter når shard-grensene. Angi CloudWatch-alarmer for iteratoralder og overskredet tildelt gjennomstrømming for å utløse Auto Scaling av 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-alerts

Hurtigsjekk

Test forståelsen Deres av AWS Solutions Architect-konsepter (SAA-C03) fra denne leksjonen.

Oppsummering av leksjonen

I denne leksjonen lærte De at: Kinesis Data Streams tilbyr holdbar, sortert og shard-basert strømming med flere konsumenter og støtte for avspilling, Kinesis Data Firehose tilbyr fullt administrert levering uten kode til S3, Redshift og OpenSearch, og Managed Service for Apache Flink muliggjør tilstandsbevarende sanntidsanalyse av datastrømmer. Neste tema er de 7 R-ene i strategien for skymigrering.

Gratis å komme i gang

Lær deg Cloud & IT Cert Prep med en AI-veileder – gratis

Skriv og kjør ekte kode i nettleseren, få umiddelbar hjelp fra en AI-veileder som er tilgjengelig døgnet rundt, og fortsett der du slapp – på nettet eller i appen.

Kurs
150
Leksjoner
600

Ofte stilte spørsmål

Er leksjonen «Kinesis Streams, Firehose og sanntidsanalyse» gratis?

Ja – hele teksten i «Kinesis Streams, Firehose og sanntidsanalyse» er gratis å lese her på nettet. For å øve interaktivt med en innebygd kodeeditor og en AI-veileder som er tilgjengelig døgnet rundt, og for å låse opp resten av Cloud & IT Cert Prep-kurset, kan du oppgradere til CoddyKit PRO. Kurset i Cloud & IT Cert Prep inneholder totalt 4 leksjoner.

Hva lærer jeg i «Kinesis Streams, Firehose og sanntidsanalyse»?

Innhent datastrømmer med Kinesis Data Streams, lever dem til S3 eller Redshift med Firehose, og analyser dem i sanntid med Managed Service for Apache Flink Du øver på Cloud & IT Cert Prep med praktisk kode som du kjører direkte i nettleseren, mens en AI-veileder som er tilgjengelig døgnet rundt, svarer på spørsmålene dine mens du jobber deg gjennom leksjonen.

Trenger jeg erfaring for å begynne med Cloud & IT Cert Prep?

Ingen tidligere erfaring er nødvendig. Cloud & IT Cert Prep på CoddyKit er lagt opp for både nybegynnere og viderekomne, så De kan begynne her eller helt fra start og lære i Deres eget tempo. Dette er leksjon 4 av 4.

Hvor lang tid tar leksjonen «Kinesis Streams, Firehose og sanntidsanalyse»?

De fleste CoddyKit-leksjoner tar omtrent 5–10 minutter. Hver leksjon er kort og interaktiv, slik at De gjør jevne fremskritt og kan fortsette akkurat der De slapp – både på nettet og i appen.

Kan jeg skrive og kjøre kode i denne Cloud & IT Cert Prep-leksjonen?

Ja. Alle Cloud & IT Cert Prep-leksjoner har en innebygd kodeeditor, slik at De kan skrive og kjøre ekte kode direkte i nettleseren og få umiddelbar tilbakemelding fra AI – uten lokal konfigurering.

Alle leksjonene i dette kurset

  1. Bygging av en datasjø på S3
  2. AWS Glue: ETL og datakatalog
  3. Amazon Athena: Serverløs SQL på S3
  4. Kinesis Streams, Firehose og sanntidsanalyse
← Tilbake til Cloud & IT Cert Prep