Cloud & IT Cert Prep · Pelajaran

Kinesis Streams, Firehose dan Analitis Masa Nyata

Telan data strim dengan Kinesis Data Streams, hantar ke S3 atau Redshift menggunakan Firehose, dan analisis secara masa nyata dengan Managed Service for Apache Flink.

Pelajaran 4 daripada 413 langkah

Kinesis Streams, Firehose dan Analitis Masa Nyata ialah pelajaran Cloud & IT Cert Prep percuma di CoddyKit. Ini ialah pelajaran 4 daripada 4. Anda boleh membaca keseluruhan pelajaran di bawah secara percuma — kemudian berlatih secara praktikal dalam pelayar menggunakan penyunting kod terbina dalam dan tutor kecerdasan buatan 24/7. Pelajaran ini merupakan sebahagian daripada laluan pembelajaran Cloud & IT Cert Prep, dan kemajuan anda disegerakkan merentas web serta aplikasi CoddyKit. Kursus Cloud & IT Cert Prep merangkumi sejumlah 4 pelajaran.

Gambaran Keseluruhan Keluarga Kinesis

Amazon Kinesis ialah keluarga perkhidmatan untuk mengumpul, memproses dan menganalisis data penstriman masa nyata. Tiga perkhidmatan terasnya ialah: Kinesis Data Streams (pemprosesan tersuai dengan kependaman rendah), Kinesis Data Firehose (penghantaran terurus sepenuhnya ke S3/Redshift/OpenSearch), dan Managed Service for Apache Flink (dahulunya Kinesis Data Analytics) untuk SQL masa nyata dan pemprosesan Flink. Setiap perkhidmatan menyasarkan titik yang berbeza dalam saluran paip penstriman.

Seni Bina Kinesis Data Streams

Kinesis Data Stream ialah log tahan lama dan tersusun yang dipartitionkan kepada serpihan. Setiap serpihan menyediakan daya pemprosesan tulis 1 MB/s dan daya pemprosesan baca 2 MB/s. Rekod data disimpan selama 24 jam secara lalai (boleh dilanjutkan kepada 7 hari atau 365 hari). Pengeluar menulis rekod ke dalam stream; pengguna — Lambda, aplikasi KCL, Firehose atau Flink — membaca daripada satu atau lebih serpihan secara selari. Rekod tidak boleh diubah selepas ditulis.

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

Serpihan, Daya Pemprosesan dan Penskalaan

Bilangan serpihan menentukan jumlah daya pemprosesan stream. Anda boleh memisahkan serpihan untuk menggandakan daya pemprosesan atau menggabungkan dua serpihan untuk mengurangkan kos. Gunakan Enhanced Fan-Out untuk memberikan setiap pengguna berdaftar daya pemprosesan baca sendiri sebanyak 2 MB/s tanpa bergantung pada pengguna lain, sekali gus menghapuskan pendikitan bacaan apabila beberapa aplikasi menggunakan stream yang sama. Pantau GetRecords.IteratorAgeMilliseconds untuk mengesan kelewatan pengguna.

# 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: Penghantaran Terurus

Kinesis Data Firehose ialah perkhidmatan terurus sepenuhnya yang menangkap, mengubah dan menghantar data penstriman ke destinasi termasuk S3, Amazon Redshift, Amazon OpenSearch Service, Splunk dan titik akhir HTTP. Tiada serpihan untuk diuruskan — Firehose diskalakan secara automatik. Anda mengkonfigurasikan saiz penimbal (1–128 MB) dan selang penimbal (60–900 saat); Firehose menghantar data apabila salah satu had dicapai terlebih dahulu.

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

Transformasi Data Firehose dengan Lambda

Firehose boleh memanggil fungsi Lambda bagi setiap kelompok rekod sebelum penghantaran untuk mengubah, memperkaya atau menapis data semasa penghantaran. Kes penggunaan biasa termasuk menukar JSON kepada Parquet (melalui skema Glue), menyamarkan medan PII atau menggugurkan peristiwa bernilai rendah. Rekod yang gagal ditransformasikan secara pilihan dihantar ke awalan ralat S3 yang berasingan untuk diproses semula, supaya tiada data hilang.

# 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 (dahulunya Kinesis Data Analytics) menjalankan aplikasi Apache Flink pada infrastruktur terurus sepenuhnya. Gunakannya untuk analitik masa nyata berkeadaan: agregasi tetingkap gelongsor, pengesanan anomali, pemadanan corak merentasi jujukan peristiwa dan penggabungan data penstriman dengan jadual rujukan. Anda menulis kod Flink dalam Java, Python atau Scala, manakala Flink mengurus titik semak dan keadaan tepat sekali.

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

Memilih antara Streams dan Firehose

Untuk peperiksaan SAA-C03, ketahui kriteria pemilihannya. Gunakan Kinesis Data Streams apabila anda memerlukan kependaman kurang daripada satu saat, beberapa pengguna membaca secara serentak atau logik pemprosesan tersuai dengan kawalan penuh terhadap pengekalan dan main semula. Gunakan Firehose apabila anda hanya perlu menghantar data penstriman dengan pasti ke S3, Redshift atau OpenSearch menggunakan kod minimum, transformasi semasa penghantaran secara pilihan dan penskalaan automatik dengan kependaman lebih tinggi (60+ saat).

Kinesis berbanding SQS: Pilihan Klasik Peperiksaan

Soalan SAA-C03 yang biasa meminta anda memilih antara Kinesis dan SQS. Perbezaan utama: Kinesis mengekalkan tertib mesej dalam satu serpihan, menyokong beberapa pengguna membaca data yang sama secara serentak dan mengekalkan rekod untuk main semula. SQS mengalih keluar mesej selepas digunakan (tiada main semula), mod FIFO menjamin tertib yang ketat dan ia lebih sesuai untuk menyahgandingkan perkhidmatan mikro. Jika senario menyebut analitik masa nyata atau main semula, pilih Kinesis.

Pengeluar Kinesis: SDK dan KPL

Kinesis Producer Library (KPL) ialah klien berdaya pemprosesan tinggi untuk menulis ke Kinesis Data Streams daripada aplikasi. KPL secara automatik mengagregatkan beberapa rekod kecil ke dalam satu panggilan API (sehingga 1 MB) dan mengendalikan percubaan semula dengan pengunduran. Ini mengurangkan kos PUT bagi setiap rekod dengan ketara dan meningkatkan daya pemprosesan bagi setiap serpihan. Gunakan KPL untuk pengeluar volum tinggi seperti penstriman klik web, penderia IoT atau saluran paip log.

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

Corak Firehose ke Redshift

Seni bina yang biasa ialah menggunakan Firehose sebagai saluran paip terurus daripada stream Kinesis atau pengeluar terus ke Amazon Redshift untuk pergudangan data. Firehose terlebih dahulu menulis data ke bucket pementasan S3 perantara, kemudian mengeluarkan perintah COPY untuk memuatkan data ke dalam Redshift. Ini ialah cara paling cekap untuk memuatkan data penstriman secara pukal ke dalam Redshift — sisipan terus baris demi baris ke dalam Redshift akan menjadi sangat perlahan kerana overhed bagi setiap baris.

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

Memantau Kinesis dengan CloudWatch

Metrik CloudWatch utama untuk Kinesis: IncomingBytes dan IncomingRecords untuk mengukur daya pemprosesan pengeluar, GetRecords.IteratorAgeMilliseconds untuk mengukur kelewatan pengguna (nilai tinggi bermaksud pengguna tidak dapat mengimbangi kadar masuk), dan WriteProvisionedThroughputExceeded untuk mengesan apabila pengeluar mencapai had serpihan. Tetapkan penggera CloudWatch pada umur iterator dan lebihan daya pemprosesan diperuntukkan untuk mencetuskan Penskalaan Automatik bagi serpihan melalui 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

Semakan Pantas

Uji pemahaman anda tentang konsep AWS Solutions Architect (SAA-C03) daripada pelajaran ini.

Ringkasan Pelajaran

Dalam pelajaran ini, anda telah mempelajari bahawa: Kinesis Data Streams menyediakan penstriman tahan lama, tersusun dan berasaskan serpihan dengan sokongan untuk beberapa pengguna serta keupayaan main semula, Kinesis Data Firehose menawarkan penghantaran terurus sepenuhnya tanpa kod ke S3, Redshift dan OpenSearch, dan Managed Service for Apache Flink membolehkan analitik masa nyata berkeadaan pada stream. Seterusnya, kita akan meneroka 7 R strategi migrasi awan.

Percuma untuk bermula

Pelajari Cloud & IT Cert Prep dengan tutor kecerdasan buatan — percuma

Tulis dan jalankan kod sebenar dalam pelayar anda, dapatkan bantuan segera daripada tutor kecerdasan buatan yang tersedia 24/7, dan sambung semula dari tempat anda berhenti di web atau dalam aplikasi.

Kursus
150
Pelajaran
600

Soalan Lazim

Adakah pelajaran “Kinesis Streams, Firehose dan Analitis Masa Nyata” percuma?

Ya — teks penuh “Kinesis Streams, Firehose dan Analitis Masa Nyata” boleh dibaca secara percuma di web ini. Untuk berlatih secara interaktif menggunakan penyunting kod terbina dalam dan tutor kecerdasan buatan 24/7, serta membuka kunci baki kursus Cloud & IT Cert Prep, tingkat taraf kepada CoddyKit PRO. Kursus Cloud & IT Cert Prep merangkumi sejumlah 4 pelajaran.

Apakah yang akan saya pelajari dalam “Kinesis Streams, Firehose dan Analitis Masa Nyata”?

Telan data strim dengan Kinesis Data Streams, hantar ke S3 atau Redshift menggunakan Firehose, dan analisis secara masa nyata dengan Managed Service for Apache Flink. Anda berlatih Cloud & IT Cert Prep menggunakan kod praktikal yang dijalankan terus dalam pelayar, manakala tutor kecerdasan buatan 24/7 menjawab soalan anda semasa anda mengikuti pelajaran.

Adakah saya memerlukan pengalaman untuk memulakan Cloud & IT Cert Prep?

Tiada pengalaman terdahulu diperlukan. Pembelajaran Cloud & IT Cert Prep di CoddyKit disusun untuk pelajar daripada peringkat pemula hingga lanjutan, jadi anda boleh bermula di sini atau dari awal dan belajar mengikut kadar anda sendiri. Ini ialah pelajaran 4 daripada 4.

Berapa lamakah pelajaran “Kinesis Streams, Firehose dan Analitis Masa Nyata” diambil?

Kebanyakan pelajaran CoddyKit mengambil masa kira-kira 5–10 minit. Setiap pelajaran ringkas dan interaktif, jadi anda boleh membuat kemajuan secara berterusan dan menyambung tepat dari tempat anda berhenti di web atau aplikasi.

Bolehkah saya menulis dan menjalankan kod dalam pelajaran Cloud & IT Cert Prep ini?

Ya. Setiap pelajaran Cloud & IT Cert Prep menyertakan penyunting kod terbina dalam, jadi anda boleh menulis dan menjalankan kod sebenar terus dalam pelayar serta menerima maklum balas kecerdasan buatan serta-merta — tanpa memerlukan persediaan setempat.

Semua pelajaran dalam kursus ini

  1. Membina Tasik Data pada S3
  2. AWS Glue: ETL dan Katalog Data
  3. Amazon Athena: SQL Tanpa Pelayan pada S3
  4. Kinesis Streams, Firehose dan Analitis Masa Nyata
← Kembali ke Cloud & IT Cert Prep