Kinesis Data Streams untuk Pemrosesan Peristiwa Waktu Nyata
Hasilkan dan konsumsi aliran peristiwa ber-throughput tinggi dengan Kinesis Data Streams, kelola shard untuk throughput, lalu gunakan Lambda sebagai konsumen.
Kinesis Data Streams untuk Pemrosesan Peristiwa Waktu Nyata adalah pelajaran AWS Solutions Architect gratis di CoddyKit. Ini adalah pelajaran 3 dari 4. Kamu bisa membaca pelajaran lengkapnya di bawah secara gratis — lalu praktikkan langsung di browser dengan editor kode bawaan dan tutor AI 24/7. Ini adalah bagian dari jalur belajar AWS Solutions Architect, dan progresmu tersinkronisasi di web dan aplikasi CoddyKit. Kursus AWS Solutions Architect mencakup 4 pelajaran total.
Konsep Inti Kinesis Data Streams
Kinesis Data Streams (KDS) adalah layanan streaming data waktu nyata yang tahan lama dan berurutan. Data diatur ke dalam sebuah aliran yang terdiri dari satu atau beberapa pecahan. Setiap pecahan merupakan urutan rekaman data yang berurutan. Produsen menempatkan rekaman ke pecahan menggunakan kunci partisi yang menentukan pecahan penerima rekaman tersebut. Konsumen membaca rekaman dari pecahan dan memprosesnya sesuai urutan kedatangannya dalam setiap pecahan.
Kapasitas Pecahan dan Batas Throughput
Setiap pecahan mendukung throughput penulisan 1 MB/detik atau 1.000 rekaman/detik dan throughput pembacaan 2 MB/detik (digunakan bersama oleh semua konsumen standar pada pecahan tersebut). Kapasitas total aliran meningkat secara linear sesuai jumlah pecahan. Gunakan rumus berikut: shards_needed = max(write_MB_per_s / 1, read_MB_per_s / 2). Jika produsen mencapai batas penulisan, akan muncul galat ProvisionedThroughputExceededException — atasi dengan membagi pecahan atau mendistribusikan kunci partisi secara lebih merata.
# 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 5Kunci Partisi dan Distribusi Data
Kunci partisi adalah string yang di-hash oleh Kinesis menggunakan MD5 untuk menentukan pecahan penerima suatu rekaman. Kunci partisi yang dipilih dengan baik mendistribusikan rekaman secara merata ke seluruh pecahan (pencegahan pecahan padat). Untuk IoT, gunakan ID perangkat. Untuk aliran klik, gunakan ID sesi atau ID pengguna. Hindari kunci dengan kardinalitas rendah (misalnya, nama negara yang hanya memiliki 5 nilai), karena kunci tersebut menyebabkan pecahan padat: satu pecahan menerima lalu lintas penulisan yang tidak proporsional, sementara pecahan lainnya tidak digunakan.
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
)Konsumen Standar dibandingkan dengan Enhanced Fan-Out
Konsumen standar berbagi throughput pembacaan 2 MB/detik per pecahan menggunakan GetRecords dengan polling. Jika terdapat 3 konsumen pada sebuah pecahan dan masing-masing memerlukan 2 MB/detik, semuanya akan dibatasi karena harus berbagi total 2 MB/detik. Enhanced Fan-Out (EFO) memberikan setiap konsumen terdaftar saluran pembacaan khusus sebesar 2 MB/detik melalui koneksi push HTTP/2 persisten (SubscribeToShard). EFO menambahkan biaya per konsumen-per pecahan-per jam, tetapi sepenuhnya menghilangkan persaingan pembacaan.
# 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-telemetryLambda sebagai Konsumen Kinesis
Lambda terintegrasi secara native dengan Kinesis Data Streams melalui Pemetaan Sumber Peristiwa. Lambda melakukan polling terhadap aliran, membaca kumpulan rekaman, lalu memanggil fungsi Anda. Konfigurasikan BatchSize (1–10.000 rekaman), StartingPosition (TRIM_HORIZON untuk rekaman terlama, LATEST untuk rekaman terbaru), dan BisectBatchOnFunctionError untuk membagi kumpulan rekaman yang gagal. Faktor Paralelisasi (1–10) memungkinkan Lambda menjalankan beberapa pemanggilan secara bersamaan per pecahan agar dapat mengikuti aliran yang bergerak cepat.
# 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) adalah kerangka kerja aplikasi untuk membangun konsumen Kinesis berbasis Java yang tangguh (atau multibahasa melalui MultiLangDaemon). KCL menangani enumerasi pecahan, pengelolaan sewa (mendistribusikan pecahan di antara instans Worker), pencatatan titik pemeriksaan kemajuan ke DynamoDB, serta pemisahan dan penggabungan pecahan secara mulus. Setiap Worker KCL memproses satu atau beberapa pecahan, dan KCL secara otomatis menyeimbangkan kembali pecahan saat Worker bertambah atau gagal. KCL lebih disarankan daripada polling SDK mentah untuk aplikasi konsumen produksi.
# 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();Penyimpanan Data dan Pemutaran Ulang
Kinesis Data Streams menyimpan rekaman selama 24 jam secara default (dapat diperpanjang hingga 7 hari atau hingga 365 hari dengan Penyimpanan Jangka Panjang dengan biaya tambahan). Berbeda dari SQS, rekaman yang telah dikonsumsi tidak dihapus setelah dikonsumsi — rekaman tetap tersedia hingga masa penyimpanan berakhir. Hal ini memungkinkan beberapa konsumen membaca rekaman yang sama secara independen, serta memungkinkan pemutaran ulang dengan mengatur ulang titik pemeriksaan konsumen ke nomor urutan yang lebih awal — sangat berguna untuk memperbaiki bug atau mengisi data untuk layanan baru.
# 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 100Mode Kapasitas Sesuai Permintaan
Kinesis Data Streams mendukung dua mode kapasitas. Mode Provisioned: Anda mengelola jumlah pecahan secara manual dan membayar per jam pecahan. Mode On-Demand: Kinesis secara otomatis menyesuaikan kapasitas pecahan berdasarkan throughput yang masuk (hingga 200 MB/detik untuk penulisan dan 400 MB/detik untuk pembacaan secara default), dan Anda membayar per GB data yang ditulis serta diambil. Mode On-Demand ideal untuk pola lalu lintas yang berubah-ubah atau tidak dapat diprediksi ketika Anda tidak ingin mengelola penskalaan pecahan.
# 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_DEMANDJaminan Pengurutan dalam Pecahan
Kinesis menjamin pengurutan dalam sebuah pecahan — rekaman dengan kunci partisi yang sama selalu masuk ke pecahan yang sama dan dibaca sesuai urutan penulisannya. Namun, tidak ada jaminan pengurutan antarpecahan. Jika aplikasi Anda memerlukan pengurutan global untuk semua rekaman, gunakan satu pecahan (yang membatasi throughput hingga 1 MB/detik) atau rancang ulang aplikasi agar pengurutan hanya diperlukan dalam kelompok kunci partisi (misalnya, pengurutan per perangkat). Ini adalah perbedaan umum dalam ujian dibandingkan dengan SQS FIFO, yang menyediakan deduplikasi dan pengurutan ketat.
Keamanan: Enkripsi dan VPC
Kinesis Data Streams mengenkripsi rekaman di sisi server saat tersimpan menggunakan AWS KMS (CMK atau kunci yang dikelola AWS) ketika enkripsi sisi server diaktifkan. Semua data saat transit dienkripsi dengan TLS. Untuk aplikasi yang berjalan di VPC dan tidak boleh mengirim data aliran melalui internet publik, gunakan VPC Interface Endpoint (PrivateLink) untuk Kinesis agar lalu lintas sepenuhnya tetap berada dalam jaringan inti AWS — hal yang penting untuk beban kerja yang diatur oleh persyaratan kepatuhan.
# 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-abc123Kinesis Data Streams dibandingkan dengan SQS dan Kafka
Untuk ujian SAA-C03, bandingkan Kinesis Data Streams dengan alternatifnya. Kinesis dibandingkan dengan SQS: Kinesis mempertahankan urutan dalam pecahan dan mendukung beberapa konsumen yang membaca data yang sama; SQS menghapus pesan setelah dikonsumsi. Kinesis dibandingkan dengan MSK (Kafka): MSK adalah Apache Kafka terkelola — gunakan MSK jika Anda memerlukan kompatibilitas protokol Kafka, konfigurasi topik tingkat lanjut, atau sedang bermigrasi dari Kafka lokal. Gunakan Kinesis untuk streaming native AWS dengan integrasi yang lebih erat ke Lambda, Firehose, dan Flink. Pilih Kinesis kecuali pertanyaan secara khusus menyebutkan Kafka atau persyaratan kompatibilitas Kafka.
Pemeriksaan Singkat
Uji pemahaman Anda tentang konsep AWS Solutions Architect (SAA-C03) dari pelajaran ini.
Ringkasan Pelajaran
Dalam pelajaran ini, Anda mempelajari bahwa: Kinesis Data Streams menyediakan aliran terpecah, berurutan, dan tahan lama dengan penyimpanan default selama 24 jam serta kemampuan pemutaran ulang, Enhanced Fan-Out memberikan setiap konsumen 2 MB/detik khusus per pecahan sehingga menghilangkan persaingan pembacaan, dan mode On-Demand secara otomatis menyesuaikan skala pecahan untuk lalu lintas yang tidak dapat diprediksi. Berikutnya, kita akan membahas pola choreography dan orchestration dalam arsitektur berbasis peristiwa.
Pertanyaan yang Sering Diajukan
Apakah pelajaran “Kinesis Data Streams untuk Pemrosesan Peristiwa Waktu Nyata” gratis?
Ya — teks lengkap “Kinesis Data Streams untuk Pemrosesan Peristiwa Waktu Nyata” gratis dibaca di sini di web. Untuk praktiknya secara interaktif (editor kode bawaan dan tutor AI 24/7) dan buka sisa kursus AWS Solutions Architect, upgrade ke CoddyKit PRO. Kursus AWS Solutions Architect mencakup 4 pelajaran total.
Apa yang akan aku pelajari di “Kinesis Data Streams untuk Pemrosesan Peristiwa Waktu Nyata”?
Hasilkan dan konsumsi aliran peristiwa ber-throughput tinggi dengan Kinesis Data Streams, kelola shard untuk throughput, lalu gunakan Lambda sebagai konsumen. Kamu berlatih AWS Solutions Architect dengan kode praktik yang langsung kamu jalankan di browser, dan tutor AI 24/7 menjawab pertanyaanmu saat kamu mengerjakan pelajaran ini.
Apakah aku perlu pengalaman untuk memulai AWS Solutions Architect?
Tidak diperlukan pengalaman sebelumnya. AWS Solutions Architect di CoddyKit dirancang untuk pemula hingga pelajar tingkat lanjut, jadi kamu bisa memulai di sini atau dari awal dan belajar sesuai kecepatan kamu sendiri. Ini adalah pelajaran 3 dari 4.
Berapa lama pelajaran “Kinesis Data Streams untuk Pemrosesan Peristiwa Waktu Nyata” memakan waktu?
Sebagian besar pelajaran CoddyKit memakan waktu sekitar 5–10 menit. Setiap pelajaran ringkas dan interaktif, jadi kamu membuat kemajuan stabil dan melanjutkan dari tempat kamu tinggalkan di web dan aplikasi.
Bisakah aku menulis dan menjalankan kode dalam pelajaran AWS Solutions Architect ini?
Ya. Setiap pelajaran AWS Solutions Architect menyertakan editor kode bawaan, jadi kamu menulis dan menjalankan kode nyata langsung di browser dan mendapatkan umpan balik AI instan — tidak diperlukan penyiapan lokal.
Semua pelajaran dalam kursus ini
- EventBridge: Bus Peristiwa dan Aturan
- Step Functions: Mengatur Alur Kerja Tanpa Server
- Kinesis Data Streams untuk Pemrosesan Peristiwa Waktu Nyata
- Pola Koreografi vs Orkestrasi