Cloud & IT Cert Prep · Oppitunti

Kinesis Streams, Firehose ja reaaliaikainen analytiikka

Vastaanota suoratoistodataa Kinesis Data Streamsilla, toimita sitä S3:een tai Redshiftiin Firehosella ja analysoi sitä reaaliajassa Managed Service for Apache Flinkillä.

Oppitunti 4/413 vaihetta

Kinesis Streams, Firehose ja reaaliaikainen analytiikka on ilmainen Cloud & IT Cert Prep-oppitunti CoddyKitissä. Tämä on oppitunti 4/4. Voit lukea koko oppitunnin alta ilmaiseksi ja harjoitella sen jälkeen käytännössä selaimessa sisäänrakennetulla koodieditorilla ja ympäri vuorokauden käytettävissä olevan tekoälytuutorin avulla. Oppitunti kuuluu Cloud & IT Cert Prep-oppimispolkuun, ja edistymisesi synkronoituu verkon ja CoddyKit-sovelluksen välillä. Cloud & IT Cert Prep-kurssilla on yhteensä 4 oppituntia.

Kinesis-palveluperheen yleiskatsaus

Amazon Kinesis on palveluperhe reaaliaikaisten suoratoistotietojen keräämiseen, käsittelyyn ja analysointiin. Kolme keskeistä palvelua ovat: Kinesis Data Streams (pieni viive ja mukautettu käsittely), Kinesis Data Firehose (täysin hallittu toimitus S3:een, Redshiftiin ja OpenSearchiin) sekä Managed Service for Apache Flink (aiemmin Kinesis Data Analytics) reaaliaikaista SQL- ja Flink-käsittelyä varten. Kukin palvelu palvelee suoratoistoputken eri vaihetta.

Kinesis Data Streams -arkkitehtuuri

Kinesis Data Stream on kestävä ja järjestetty loki, joka on jaettu shard-osioihin. Kukin shard tarjoaa 1 Mt/s kirjoituskapasiteetin ja 2 Mt/s lukukapasiteetin. Tietueita säilytetään oletusarvoisesti 24 tuntia, ja säilytysaikaa voidaan pidentää 7 päivään tai 365 päivään. Tuottajat kirjoittavat tietueita streamiin, ja kuluttajat – Lambda, KCL-sovellukset, Firehose tai Flink – lukevat tietoja yhdestä tai useammasta shard-osiosta rinnakkain. Tietueita ei voi muuttaa kirjoittamisen jälkeen.

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

Shardit, läpijuoksu ja skaalaus

Shardien määrä määrittää streamin kokonaisläpijuoksun. Voit jakaa shardin kaksinkertaistaaksesi läpijuoksun tai yhdistää kaksi shardia kustannusten pienentämiseksi. Käytä Enhanced Fan-Out -ominaisuutta, jos haluat antaa kullekin rekisteröidylle kuluttajalle oman, muista kuluttajista riippumattoman 2 Mt/s lukukapasiteetin. Näin vältät lukuoperaatioiden rajoittamisen, kun useat sovellukset kuluttavat samaa streamia. Seuraa mittaria GetRecords.IteratorAgeMilliseconds kuluttajien viiveen havaitsemiseksi.

# 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: hallittu toimitus

Kinesis Data Firehose on täysin hallittu palvelu, joka kerää, muuntaa ja toimittaa suoratoistotietoja muun muassa S3:een, Amazon Redshiftiin, Amazon OpenSearch Serviceen, Splunkiin ja HTTP-päätepisteisiin. Hallittavia shardeja ei ole, sillä Firehose skaalautuu automaattisesti. Määrität puskurikoon (1–128 Mt) ja puskurivälin (60–900 sekuntia); Firehose toimittaa tiedot, kun jompikumpi raja saavutetaan ensin.

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

Firehosen tietojen muuntaminen Lambdalla

Firehose voi kutsua Lambda-funktiota jokaiselle tietue-erälle ennen toimitusta ja muuntaa, rikastaa tai suodattaa tietoja niiden kulkiessa. Yleisiä käyttötapauksia ovat JSON:n muuntaminen Parquetiksi (Glue-skeeman avulla), henkilötietokenttien peittäminen ja vähäarvoisten tapahtumien poistaminen. Muunnoksessa epäonnistuvat tietueet voidaan valinnaisesti lähettää erilliseen S3-virheprefiksiin uudelleenkäsittelyä varten, joten tietoja ei menetetä.

# 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 (aiemmin Kinesis Data Analytics) suorittaa Apache Flink -sovelluksia täysin hallitussa infrastruktuurissa. Käytä sitä tilallisessa reaaliaikaisessa analytiikassa, kuten liukuvien ikkunoiden koostamisessa, poikkeamien havaitsemisessa, tapahtumasarjojen mallien tunnistamisessa ja suoratoistotietojen yhdistämisessä viitetauluihin. Kirjoitat Flink-koodin Javalla, Pythonilla tai Scalalla, ja Flink hallitsee tarkistuspisteet sekä exactly-once-tilan.

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

Streams- ja Firehose-palveluiden välillä valitseminen

SAA-C03-kokeessa on tärkeää tuntea valintaperusteet. Käytä Kinesis Data Streams -palvelua, kun tarvitset alle sekunnin viiveen, useita samanaikaisesti lukevia kuluttajia tai mukautettua käsittelylogiikkaa, jossa säilytyksen ja uudelleenluvun hallinta on täysin omissa käsissäsi. Käytä Firehosea, kun tarvitset vain suoratoistotietojen luotettavan toimituksen S3:een, Redshiftiin tai OpenSearchiin vähäisellä koodimäärällä, valinnaisella käsittelynaikaisella muunnoksella ja automaattisella skaalauksella, mutta hyväksyt suuremman viiveen (vähintään 60 sekuntia).

Kinesis vai SQS: klassinen koetehtävä

Yleisessä SAA-C03-kysymyksessä on valittava Kinesiksen ja SQS:n välillä. Keskeiset erot ovat seuraavat: Kinesis säilyttää viestien järjestyksen shardin sisällä, tukee useita samaa dataa samanaikaisesti lukevia kuluttajia ja säilyttää tietueet uudelleenlukua varten. SQS poistaa viestit kulutuksen jälkeen (uudelleenlukua ei ole), FIFO-tila takaa tiukan järjestyksen ja palvelu soveltuu paremmin mikropalveluiden irrottamiseen toisistaan. Jos tilanteessa mainitaan reaaliaikainen analytiikka tai uudelleenluku, valitse Kinesis.

Kinesiksen tuottajat: SDK ja KPL

Kinesis Producer Library (KPL) on suuren läpijuoksun asiakaskirjasto Kinesis Data Streams -palveluun kirjoittamiseen sovelluksista. KPL kokoaa automaattisesti useita pieniä tietueita yhteen API-kutsuun (enintään 1 Mt) ja käsittelee uudelleenyritykset back-off-mekanismilla. Tämä pienentää huomattavasti tietuekohtaista PUT-kustannusta ja kasvattaa shard-kohtaista läpijuoksua. Käytä KPL:ää suuren volyymin tuottajissa, kuten verkkosivujen klikkausvirroissa, IoT-antureissa tai lokiputkissa.

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

Firehose–Redshift-malli

Yleinen arkkitehtuuri on käyttää Firehosea hallittuna putkena Kinesis-streameista tai suoraan tuottajilta Amazon Redshiftiin tietovarastointia varten. Firehose kirjoittaa tiedot ensin väliaikaiseen S3-välivarastoon ja antaa sen jälkeen COPY-komennon tietojen lataamiseksi Redshiftiin. Tämä on tehokkain tapa ladata suoratoistotietoja erissä Redshiftiin – suorat rivi kerrallaan tehtävät lisäykset Redshiftiin olisivat erittäin hitaita rivikohtaisen käsittelykustannuksen vuoksi.

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

Kinesiksen valvonta CloudWatchilla

Kinesiksen keskeisiä CloudWatch-mittareita ovat IncomingBytes ja IncomingRecords tuottajien läpijuoksun mittaamiseen, GetRecords.IteratorAgeMilliseconds kuluttajien viiveen mittaamiseen (suuri arvo tarkoittaa, etteivät kuluttajat pysy mukana) sekä WriteProvisionedThroughputExceeded sen havaitsemiseen, milloin tuottajat saavuttavat shardien rajat. Määritä CloudWatch-hälytykset iteraattorin iälle ja varatun läpijuoksun ylittymiselle, jotta shardien Auto Scaling voidaan käynnistää Application Auto Scalingin avulla.

# 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

Pikatarkistus

Testaa, kuinka hyvin ymmärrät tämän oppitunnin AWS Solutions Architect (SAA-C03) -käsitteet.

Oppitunnin yhteenveto

Tässä oppitunnissa opit, että Kinesis Data Streams tarjoaa kestävän, järjestetyn ja shardeihin jaetun suoratoiston, jossa on useita kuluttajia ja uudelleenlukumahdollisuus, Kinesis Data Firehose tarjoaa täysin hallitun ja koodittoman toimituksen S3:een, Redshiftiin ja OpenSearchiin ja Managed Service for Apache Flink mahdollistaa tilallisen reaaliaikaisen analytiikan streameista. Seuraavaksi tutustumme pilvimigraatiostrategian seitsemään R:ään.

Aloita maksutta

Opi Cloud & IT Cert Prep 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
150
Oppitunnit
600

Usein kysytyt kysymykset

Onko oppitunti ”Kinesis Streams, Firehose ja reaaliaikainen analytiikka” ilmainen?

Kyllä – oppitunnin ”Kinesis Streams, Firehose ja reaaliaikainen analytiikka” koko tekstin voi lukea täällä verkossa ilmaiseksi. Jos haluat harjoitella interaktiivisesti sisäänrakennetulla koodieditorilla ja ympäri vuorokauden käytettävissä olevan tekoälytuutorin avulla sekä avata koko Cloud & IT Cert Prep-kurssin, päivitä CoddyKit PROhon. Cloud & IT Cert Prep-kurssilla on yhteensä 4 oppituntia.

Mitä opin oppitunnilla ”Kinesis Streams, Firehose ja reaaliaikainen analytiikka”?

Vastaanota suoratoistodataa Kinesis Data Streamsilla, toimita sitä S3:een tai Redshiftiin Firehosella ja analysoi sitä reaaliajassa Managed Service for Apache Flinkillä. Harjoittelet Cloud & IT Cert Prep-aihetta koodilla, jonka suoritat suoraan selaimessa. Ympäri vuorokauden käytettävissä oleva tekoälytuutori vastaa kysymyksiisi oppitunnin aikana.

Tarvitsenko kokemusta aloittaakseni Cloud & IT Cert Prep-opiskelun?

Aiempi kokemus ei ole tarpeen. CoddyKitin Cloud & IT Cert Prep-oppimispolku sopii vasta-alkajista edistyneisiin, joten voit aloittaa tästä tai alusta ja edetä omaan tahtiisi. Tämä on oppitunti 4/4.

Kuinka kauan ”Kinesis Streams, Firehose ja reaaliaikainen analytiikka”-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ä Cloud & IT Cert Prep-oppitunnilla?

Kyllä. Jokainen Cloud & IT Cert Prep-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. Data laken rakentaminen S3:een
  2. AWS Glue: ETL ja datakatalogi
  3. Amazon Athena: palvelimeton SQL S3:ssa
  4. Kinesis Streams, Firehose ja reaaliaikainen analytiikka
← Takaisin: Cloud & IT Cert Prep