AWS Solutions Architect · Oppitunti

Kinesis Data Streams reaaliaikaiseen tapahtumankäsittelyyn

Tuota ja kuluta suuren läpimenon tapahtumavirtoja Kinesis Data Streamsilla, hallitse shardit läpimenon säätämiseksi ja käytä Lambdaa kuluttajana.

Oppitunti 3/413 vaihetta

Kinesis Data Streams reaaliaikaiseen tapahtumankäsittelyyn on ilmainen AWS Solutions Architect-oppitunti CoddyKitissä. Tämä on oppitunti 3/4. Voit lukea tästä oppimispolusta kokonaan mitkä tahansa 3 oppituntia ilmaiseksi — sen jälkeen CoddyKit PRO avaa kaikki oppitunnit sekä käytännön harjoittelun sisäänrakennetulla koodieditorilla ja ympäri vuorokauden toimivalla tekoälytuutorilla. Oppitunti kuuluu AWS Solutions Architect-oppimispolkuun, ja edistymisesi synkronoituu verkon ja CoddyKit-sovelluksen välillä. AWS Solutions Architect-kurssilla on yhteensä 4 oppituntia.

Kinesis Data Streamsin keskeiset käsitteet

Kinesis Data Streams (KDS) on kestävä, järjestystä säilyttävä reaaliaikainen tietovirtojen suoratoistopalvelu. Data järjestetään yhdeksi tai useammaksi shardiksi koostuvaksi streamiksi. Jokainen shard on järjestetty tietueiden sarja. Tuottajat sijoittavat tietueet shardeihin käyttämällä partition-avainta, joka määrittää, mikä shard vastaanottaa tietueen. Kuluttajat lukevat tietueita shardeista ja käsittelevät ne kussakin shardissa saapumisjärjestyksessä.

Shardin kapasiteetti ja suorituskykyrajat

Kukin shard tukee 1 Mt/s:n tai 1 000 tietueen/s:n kirjoitussuorituskykyä ja 2 Mt/s:n lukusuorituskykyä (joka jaetaan shardin kaikkien tavallisten kuluttajien kesken). Streamin kokonaiskapasiteetti kasvaa lineaarisesti shardien määrän mukaan. Käyttäkää kaavaa: shards_needed = max(write_MB_per_s / 1, read_MB_per_s / 2). Jos tuottajat saavuttavat kirjoitusrajan, näkyviin tulee ProvisionedThroughputExceededException-virheitä. Ratkaiskaa ongelma jakamalla shardeja tai hajauttamalla partition-avaimia tasaisemmin.

# 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

Partition-avaimet ja datan jakautuminen

Partition-avain on merkkijono, josta Kinesis laskee tiivisteen (MD5) määrittääkseen, mikä shard vastaanottaa tietueen. Hyvin valittu partition-avain jakaa tietueet tasaisesti shardien kesken (kuumien shardien ehkäiseminen). IoT-järjestelmissä käyttäkää laitetunnistetta. Klikkivirroissa käyttäkää istunto- tai käyttäjätunnistetta. Välttäkää vähän erilaisia arvoja tuottavia avaimia (esimerkiksi maan nimeä, jolla on vain viisi mahdollista arvoa), sillä ne aiheuttavat kuumia shardeja: yksi shard vastaanottaa suhteettoman paljon kirjoitusliikennettä muiden jäädessä käyttämättömiksi.

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
)

Tavalliset kuluttajat ja Enhanced Fan-Out

Tavalliset kuluttajat jakavat shardikohtaisen 2 Mt/s:n lukusuorituskyvyn käyttämällä GetRecords-kutsua ja kyselyitä. Jos shardilla on kolme kuluttajaa, joista kukin tarvitsee 2 Mt/s, niiden suorituskykyä rajoitetaan, koska ne jakavat yhteensä 2 Mt/s. Enhanced Fan-Out (EFO) antaa jokaiselle rekisteröidylle kuluttajalle oman erillisen 2 Mt/s:n lukukanavan jatkuvan HTTP/2-push-yhteyden kautta (SubscribeToShard). EFO aiheuttaa kuluttaja-shardituntiin perustuvia kustannuksia, mutta poistaa lukukilpailun kokonaan.

# 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 Kinesis-kuluttajana

Lambda integroituu Kinesis Data Streams -palveluun natiivisti Event Source Mapping -määrityksen kautta. Lambda kyselee streamia, lukee tietue-eriä ja kutsuu funktiotanne. Määrittäkää BatchSize (1–10 000 tietuetta), StartingPosition (TRIM_HORIZON vanhimmille ja LATEST uusimmille tietueille) sekä BisectBatchOnFunctionError epäonnistuneiden erien jakamista varten. Parallelisation Factor (1–10) mahdollistaa sen, että Lambda käynnistää shardia kohden useita samanaikaisia kutsuja pysyäkseen nopeasti kasvavien streamien tahdissa.

# 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) on sovelluskehys kestävien Java-pohjaisten (tai MultiLangDaemonin avulla monikielisten) Kinesis-kuluttajien rakentamiseen. KCL huolehtii shardien luetteloinnista, lease-hallinnasta (shardien jakamisesta työntekijäinstanssien kesken), edistymisen tarkistuspisteiden tallentamisesta DynamoDB:hen sekä shardien jakamisen ja yhdistämisen hallitusta käsittelystä. Kukin KCL-työntekijä käsittelee yhtä tai useampaa shardia, ja KCL tasapainottaa shardit automaattisesti uudelleen työntekijöiden määrän kasvaessa tai työntekijöiden vikaantuessa. KCL:ää suositellaan tuotannon kuluttajasovelluksiin suoran SDK-kyselyn sijaan.

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

Datan säilytys ja uudelleentoisto

Kinesis Data Streams säilyttää tietueita oletusarvoisesti 24 tuntia (säilytysaikaa voi pidentää seitsemään päivään tai enintään 365 päivään Long-Term Retention -ominaisuudella lisämaksusta). Toisin kuin SQS:ssä, kulutettuja tietueita ei poisteta kulutuksen jälkeen, vaan ne ovat käytettävissä säilytysajan päättymiseen asti. Näin useat kuluttajat voivat lukea samoja tietueita toisistaan riippumatta, ja tietueet voidaan toistaa uudelleen siirtämällä kuluttajan tarkistuspiste aiempaan järjestysnumeroon. Tämä on erittäin hyödyllistä virheenkorjauksissa ja uusien palvelujen tietojen täydentämisessä.

# 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-kapasiteettitila

Kinesis Data Streams tukee kahta kapasiteettitilaa. Provisioned-tila: hallitsette shardien määrää manuaalisesti ja maksatte sharditunneista. On-Demand-tila: Kinesis skaalaa shardien kapasiteettia automaattisesti saapuvan suorituskyvyn perusteella (oletusarvoisesti enintään 200 Mt/s kirjoitusta ja 400 Mt/s lukua), ja maksatte kirjoitetun ja noudetun datan gigatavuista. On-Demand-tila sopii vaihtelevaan tai ennakoimattomaan liikenteeseen, kun ette halua hallita shardien skaalausta.

# 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

Järjestyksen säilyminen shardien sisällä

Kinesis takaa järjestyksen shardin sisällä: saman partition-avaimen tietueet päätyvät aina samaan shardiin, ja ne luetaan siinä järjestyksessä kuin ne kirjoitettiin. Shardien välillä järjestystä ei kuitenkaan taata. Jos sovelluksenne edellyttää kaikkien tietueiden välistä globaalia järjestystä, käyttäkää yhtä shardia (jolloin suorituskyky rajoittuu 1 Mt/s:iin) tai suunnitelkaa järjestelmä uudelleen niin, että järjestystä tarvitaan vain partition-avaimen ryhmän sisällä, esimerkiksi laitekohtaisesti. Tämä on yleinen koetehtävissä testattava ero SQS FIFO:on nähden, sillä SQS FIFO tarjoaa tiukan deduplikoinnin ja järjestyksen.

Tietoturva: salaus ja VPC

Kinesis Data Streams salaa tietueet palvelinpuolella levossa AWS KMS:n avulla (CMK- tai AWS-hallinnoitu avain), kun palvelinpuolen salaus on käytössä. Kaikki siirrettävä data salataan TLS:llä. Jos VPC:ssä toimivien sovellusten ei pitäisi lähettää stream-dataa julkisen internetin kautta, käyttäkää Kinesis-palvelulle VPC Interface Endpoint -päätepistettä (PrivateLink), jolloin liikenne pysyy kokonaan AWS-verkon rungossa. Tämä on tärkeää säännellyissä työkuormissa.

# 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, SQS ja Kafka

Vertailkaa SAA-C03-kokeessa Kinesis Data Streamsia vaihtoehtoihin. Kinesis ja SQS: Kinesis säilyttää järjestyksen shardin sisällä ja tukee useita kuluttajia, jotka lukevat samaa dataa; SQS poistaa viestit kulutuksen jälkeen. Kinesis ja MSK (Kafka): MSK on hallittu Apache Kafka. Käyttäkää sitä, kun tarvitsette Kafka-protokollan yhteensopivuutta, edistyneitä topic-määrityksiä tai olette siirtymässä paikallisesta Kafka-ympäristöstä. Käyttäkää Kinesisiä AWS-natiiviseen suoratoistoon, jossa integrointi Lambdaan, Firehoseen ja Flinkiin on tiiviimpi. Valitkaa Kinesis, ellei tehtävässä nimenomaisesti mainita Kafkaa tai Kafka-yhteensopivuuden vaatimuksia.

Pikatarkistus

Testatkaa tämän oppitunnin AWS Solutions Architect (SAA-C03) -käsitteiden ymmärtämistänne.

Oppitunnin yhteenveto

Tässä oppitunnissa opitte, että Kinesis Data Streams tarjoaa shardeihin jaettuja, järjestettyjä ja kestäviä streameja, joissa oletussäilytysaika on 24 tuntia ja jotka tukevat uudelleentoistoa, Enhanced Fan-Out antaa kullekin kuluttajalle shardikohtaisen erillisen 2 Mt/s:n yhteyden ja poistaa lukukilpailun ja että On-Demand-tila skaalaa shardit automaattisesti ennakoimattoman liikenteen mukaan. Seuraavaksi tutustumme tapahtumapohjaisen arkkitehtuurin choreography- ja orchestration-malleihin.

Aloita maksutta

Opi AWS Solutions Architect tekoälytuutorin avulla — ilmaiseksi

Kirjoita ja suorita oikeaa koodia selaimessa, saa välitöntä apua tekoälytuutorilta ympäri vuorokauden ja jatka siitä, mihin jäit, verkossa tai sovelluksessa.

Kurssit
30
Oppitunnit
120

Usein kysytyt kysymykset

Onko oppitunti ”Kinesis Data Streams reaaliaikaiseen tapahtumankäsittelyyn” ilmainen?

Kyllä — voit lukea täällä verkossa kokonaan ilmaiseksi mitkä tahansa AWS Solutions Architect-oppimispolun 3 oppituntia, myös oppitunnin “Kinesis Data Streams reaaliaikaiseen tapahtumankäsittelyyn”. Sen jälkeen CoddyKit PRO avaa kaikki oppitunnit sekä interaktiiviset harjoitukset sisäänrakennetulla koodieditorilla ja ympäri vuorokauden toimivalla tekoälytuutorilla. AWS Solutions Architect-kurssilla on yhteensä 4 oppituntia.

Mitä opin oppitunnilla ”Kinesis Data Streams reaaliaikaiseen tapahtumankäsittelyyn”?

Tuota ja kuluta suuren läpimenon tapahtumavirtoja Kinesis Data Streamsilla, hallitse shardit läpimenon säätämiseksi ja käytä Lambdaa kuluttajana. Harjoittelet AWS Solutions Architect-aihetta koodilla, jonka suoritat suoraan selaimessa. Ympäri vuorokauden käytettävissä oleva tekoälytuutori vastaa kysymyksiisi oppitunnin aikana.

Tarvitsenko kokemusta aloittaakseni AWS Solutions Architect-opiskelun?

Aiempi kokemus ei ole tarpeen. CoddyKitin AWS Solutions Architect-oppimispolku sopii vasta-alkajista edistyneisiin, joten voit aloittaa tästä tai alusta ja edetä omaan tahtiisi. Tämä on oppitunti 3/4.

Kuinka kauan ”Kinesis Data Streams reaaliaikaiseen tapahtumankäsittelyyn”-oppitunnin suorittaminen kestää?

Useimmat CoddyKitin oppitunnit kestävät noin 5–10 minuuttia. Jokainen oppitunti on lyhyt ja interaktiivinen, joten edistyt tasaisesti ja voit jatkaa siitä, mihin jäit – sekä verkossa että sovelluksessa.

Voinko kirjoittaa ja suorittaa koodia tällä AWS Solutions Architect-oppitunnilla?

Kyllä. Jokainen AWS Solutions Architect-oppitunti sisältää sisäänrakennetun koodieditorin, joten voit kirjoittaa ja suorittaa oikeaa koodia suoraan selaimessa ja saada välitöntä palautetta tekoälyltä – paikallista asennusta ei tarvita.

Kaikki tämän kurssin oppitunnit

  1. EventBridge: tapahtumaväylä ja säännöt
  2. Step Functions: palvelimettomien työnkulkujen orkestrointi
  3. Kinesis Data Streams reaaliaikaiseen tapahtumankäsittelyyn
  4. Choreography- ja orchestration-mallit
← Takaisin: AWS Solutions Architect