Cloud & IT Cert Prep · leksjon

Kinesis Data Streams for sanntidsbehandling av hendelser

Produser og konsumer hendelsesstrømmer med høy gjennomstrømming med Kinesis Data Streams, administrer shards for gjennomstrømming, og bruk Lambda som konsument

Leksjon 3 av 413 trinn

Kinesis Data Streams for sanntidsbehandling av hendelser er en gratis leksjon i Cloud & IT Cert Prep på CoddyKit. Dette er leksjon 3 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.

Grunnleggende konsepter i Kinesis Data Streams

Kinesis Data Streams (KDS) er en robust datastrømmetjeneste for ordnede data i sanntid. Data organiseres i en strøm som består av én eller flere shards. Hver shard er en ordnet sekvens av dataposter. Produsenter legger poster i shards ved hjelp av en partisjonsnøkkel som bestemmer hvilken shard som mottar posten. Konsumenter leser poster fra shards og behandler dem i den rekkefølgen de ankom i hver shard.

Kapasitet og grenser for gjennomstrømming per shard

Hver shard støtter skrivegjennomstrømming på 1 MB/s eller 1 000 poster/s og lesegjennomstrømming på 2 MB/s (delt mellom alle standardkonsumenter på den sharden). Den totale kapasiteten til strømmen skalerer lineært med antallet shards. Bruk formelen: shards_needed = max(write_MB_per_s / 1, read_MB_per_s / 2). Hvis produsentene når skrivegrensen, oppstår feil av typen ProvisionedThroughputExceededException — løs dette ved å dele opp shards eller fordele partisjonsnøklene jevnere.

# Calculate shards needed for a stream:
# - Ingest rate: 5 MB/s writes
# - Read rate: 3 consumers x 2 MB/s = 6 MB/s reads
# shards = max(5/1, 6/2) = max(5, 3) = 5 shards needed

aws kinesis create-stream \
  --stream-name iot-telemetry \
  --shard-count 5

Partisjonsnøkler og datafordeling

Partisjonsnøkkelen er en streng som Kinesis hasher (MD5) for å bestemme hvilken shard som mottar en post. En godt valgt partisjonsnøkkel fordeler poster jevnt mellom shards (forebygging av overbelastede shards). For IoT kan De bruke enhets-ID. For klikkstrømmer kan De bruke økt-ID eller bruker-ID. Unngå partisjonsnøkler med få mulige verdier (for eksempel landenavn med bare fem verdier), siden de fører til overbelastede shards der én shard mottar en uforholdsmessig stor del av skrivebelastningen, mens andre står ubenyttet.

import boto3, json

client = boto3.client('kinesis', region_name='us-east-1')

# Good: use device_id as partition key for even distribution
event = {'deviceId': 'sensor-42', 'temp': 23.5, 'ts': '2024-01-15T10:00:00Z'}
client.put_record(
    StreamName='iot-telemetry',
    Data=json.dumps(event),
    PartitionKey='sensor-42'  # high-cardinality -> even distribution
)

Standardkonsumenter kontra Enhanced Fan-Out

Standardkonsumenter deler lesegjennomstrømmingen på 2 MB/s per shard ved å bruke GetRecords med polling. Hvis De har tre konsumenter på en shard og hver trenger 2 MB/s, blir de begrenset når de deler totalt 2 MB/s. Enhanced Fan-Out (EFO) gir hver registrerte konsument sin egen dedikerte lesestrøm på 2 MB/s via en vedvarende HTTP/2-tilkobling for push (SubscribeToShard). EFO medfører en kostnad per konsument-shard-time, men eliminerer all konkurranse om lesekapasiteten.

# Register an Enhanced Fan-Out consumer
aws kinesis register-stream-consumer \
  --stream-arn arn:aws:kinesis:us-east-1:123456789012:stream/iot-telemetry \
  --consumer-name real-time-analytics

# List registered consumers
aws kinesis list-stream-consumers \
  --stream-arn arn:aws:kinesis:us-east-1:123456789012:stream/iot-telemetry

Lambda som Kinesis-konsument

Lambda integreres direkte med Kinesis Data Streams via en Event Source Mapping. Lambda poller strømmen, leser grupper med poster og starter funksjonen Deres. Konfigurer BatchSize (1–10 000 poster), StartingPosition (TRIM_HORIZON for de eldste, LATEST for de nyeste) og BisectBatchOnFunctionError for å dele grupper som feiler. Parallelisation Factor (1–10) lar Lambda starte flere samtidige kjøringer per shard for å holde tritt med strømmer som beveger seg raskt.

# Create a Lambda event source mapping for Kinesis
aws lambda create-event-source-mapping \
  --function-name ProcessIoTEvents \
  --event-source-arn arn:aws:kinesis:us-east-1:123456789012:stream/iot-telemetry \
  --batch-size 100 \
  --starting-position LATEST \
  --parallelization-factor 5 \
  --bisect-batch-on-function-error true \
  --destination-config '{
    "OnFailure": {
      "Destination": "arn:aws:sqs:us-east-1:123456789012:kinesis-dlq"
    }
  }'

Kinesis Client Library (KCL)

Kinesis Client Library (KCL) er et applikasjonsrammeverk for å bygge robuste Kinesis-konsumenter i Java (eller på flere språk via MultiLangDaemon). KCL håndterer opptelling av shards, lease-administrasjon (fordeling av shards mellom arbeiderinstanser), kontrollpunktlagring av fremdrift i DynamoDB og kontrollert håndtering av oppsplitting og sammenslåing av shards. Hver KCL-arbeider behandler én eller flere shards, og KCL omfordeler automatisk shards når arbeidere skaleres ut eller feiler. KCL foretrekkes fremfor rå polling med SDK for konsumentapplikasjoner i produksjon.

# KCL stores checkpoints in a DynamoDB table automatically
# Each shard has one row tracking the last successfully processed sequence number
# KCL lease table structure:
# leaseKey (shardId) | checkpoint (sequenceNumber) | leaseOwner (workerId)
#
# To start a KCL application (pseudocode):
# KinesisClientLibConfiguration config = new KinesisClientLibConfiguration(
#   'iot-app', 'iot-telemetry', credentialsProvider, 'worker-1');
# Worker worker = new Worker.Builder().config(config).recordProcessorFactory(factory).build();
# worker.run();

Datalagring og avspilling

Kinesis Data Streams lagrer poster i 24 timer som standard (kan utvides til 7 dager eller opptil 365 dager med Long-Term Retention mot en tilleggskostnad). I motsetning til SQS blir konsumerte poster ikke slettet etter konsumering — de forblir tilgjengelige til lagringstiden utløper. Dette gjør det mulig for flere konsumenter å lese de samme postene uavhengig av hverandre, og muliggjør avspilling ved å tilbakestille en konsuments kontrollpunkt til et tidligere sekvensnummer — svært nyttig ved feilretting eller utfylling av historiske data for nye tjenester.

# Extend stream retention to 7 days
aws kinesis increase-stream-retention-period \
  --stream-name iot-telemetry \
  --retention-period-hours 168

# Get records from the oldest available record (replay)
SHARD_ITERATOR=$(aws kinesis get-shard-iterator \
  --stream-name iot-telemetry \
  --shard-id shardId-000000000000 \
  --shard-iterator-type TRIM_HORIZON \
  --query 'ShardIterator' --output text)

aws kinesis get-records --shard-iterator $SHARD_ITERATOR --limit 100

On-Demand-kapasitetsmodus

Kinesis Data Streams støtter to kapasitetsmoduser. Provisioned mode: De administrerer antallet shards manuelt og betaler per shard-time. On-Demand mode: Kinesis skalerer automatisk shard-kapasiteten basert på innkommende gjennomstrømming (som standard opptil 200 MB/s skriving og 400 MB/s lesing), og De betaler per GB med data som skrives og hentes. On-demand-modus egner seg godt for varierende eller uforutsigbare trafikkmønstre der De ikke ønsker å administrere skaleringen av shards.

# Switch an existing stream to On-Demand mode
aws kinesis update-stream-mode \
  --stream-arn arn:aws:kinesis:us-east-1:123456789012:stream/iot-telemetry \
  --stream-mode-details StreamMode=ON_DEMAND

Ordregarantier innenfor shards

Kinesis garanterer rekkefølge innenfor en shard — poster med samme partisjonsnøkkel går alltid til samme shard og leses i den rekkefølgen de ble skrevet. Det finnes imidlertid ingen garanti for rekkefølge på tvers av shards. Hvis applikasjonen Deres krever global rekkefølge for alle poster, kan De bruke én enkelt shard (som begrenser gjennomstrømmingen til 1 MB/s) eller utforme systemet på nytt slik at rekkefølge bare er nødvendig innenfor en gruppe med samme partisjonsnøkkel (for eksempel rekkefølge per enhet). Dette er en vanlig forskjell på eksamen sammenlignet med SQS FIFO, som tilbyr streng deduplisering og rekkefølge.

Sikkerhet: Kryptering og VPC

Kinesis Data Streams krypterer poster på tjenersiden når de er lagret ved hjelp av AWS KMS (CMK eller AWS-administrert nøkkel) når kryptering på tjenersiden er aktivert. Alle data under overføring krypteres med TLS. For applikasjoner som kjører i en VPC og ikke skal sende strømmens data over det offentlige internett, kan De bruke et VPC Interface Endpoint (PrivateLink) for Kinesis, slik at trafikken forblir i AWS-nettverkets interne ryggrad — noe som er viktig for regulerte arbeidsbelastninger.

# Enable server-side encryption on a Kinesis stream
aws kinesis start-stream-encryption \
  --stream-name iot-telemetry \
  --encryption-type KMS \
  --key-id arn:aws:kms:us-east-1:123456789012:key/mrk-abc123

Kinesis Data Streams kontra SQS kontra Kafka

På SAA-C03-eksamen bør De sammenligne Kinesis Data Streams med alternativene. Kinesis kontra SQS: Kinesis bevarer rekkefølgen innenfor en shard og støtter flere konsumenter som leser de samme dataene; SQS sletter meldinger etter konsumering. Kinesis kontra MSK (Kafka): MSK er administrert Apache Kafka — bruk det når De trenger kompatibilitet med Kafka-protokollen, avansert konfigurasjon av topics eller migrerer fra Kafka lokalt. Bruk Kinesis for AWS-basert strømming med tettere integrasjon med Lambda, Firehose og Flink. Velg Kinesis med mindre oppgaven spesifikt nevner Kafka eller krav om Kafka-kompatibilitet.

Hurtigsjekk

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

Oppsummering av leksjonen

I denne leksjonen lærte De: Kinesis Data Streams tilbyr oppdelte, ordnede og robuste datastrømmer med 24 timers standardlagring og mulighet for avspilling, Enhanced Fan-Out gir hver konsument en dedikert kapasitet på 2 MB/s per shard og eliminerer konkurranse om lesekapasiteten, og On-Demand-modus skalerer shards automatisk for uforutsigbar trafikk. Neste tema er mønstrene choreography og orchestration i hendelsesdrevet arkitektur.

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 Data Streams for sanntidsbehandling av hendelser» gratis?

Ja – hele teksten i «Kinesis Data Streams for sanntidsbehandling av hendelser» 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 Data Streams for sanntidsbehandling av hendelser»?

Produser og konsumer hendelsesstrømmer med høy gjennomstrømming med Kinesis Data Streams, administrer shards for gjennomstrømming, og bruk Lambda som konsument 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 3 av 4.

Hvor lang tid tar leksjonen «Kinesis Data Streams for sanntidsbehandling av hendelser»?

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. EventBridge: Hendelsesbuss og regler
  2. Step Functions: Orkestrering av serverløse arbeidsflyter
  3. Kinesis Data Streams for sanntidsbehandling av hendelser
  4. Mønstre for koreografi kontra orkestrering
← Tilbake til Cloud & IT Cert Prep