Terjemahan disediakan oleh mesin penerjemah. Jika konten terjemahan yang diberikan bertentangan dengan versi bahasa Inggris aslinya, utamakan versi bahasa Inggris.
Informasi KCL 1.x dan 2.x
penting
Amazon Kinesis Client Library (KCL) versi 1.x dan 2.x sudah usang. KCL 1.x akan mencapai akhir dukungan pada 30 Januari 2026. Kami sangat menyarankan Anda memigrasikan aplikasi KCL menggunakan versi 1.x ke versi KCL terbaru sebelum 30 Januari 2026. Untuk menemukan versi KCL terbaru, lihat halaman Perpustakaan Klien Amazon Kinesis di. GitHub
Salah satu metode pengembangan aplikasi konsumen khusus yang dapat memproses data dari aliran data KDS adalah dengan menggunakan Kinesis Client Library (KCL).
Topik
catatan
Untuk KCL 1.x dan KCL 2.x, Anda disarankan untuk meningkatkan ke versi KCL 1.x terbaru atau versi KCL 2.x, tergantung pada skenario penggunaan Anda. Baik KCL 1.x dan KCL 2.x diperbarui secara berkala dengan rilis baru yang mencakup patch ketergantungan dan keamanan terbaru, perbaikan bug, dan fitur baru yang kompatibel ke belakang. Untuk informasi selengkapnya, lihat https://github.com/awslabs/amazon-kinesis-client/releases
Tentang KCL (versi sebelumnya)
KCL membantu Anda mengkonsumsi dan memproses data dari aliran data Kinesis dengan menangani banyak tugas kompleks yang terkait dengan komputasi terdistribusi. Ini termasuk penyeimbangan beban di beberapa instans aplikasi konsumen, menanggapi kegagalan instans aplikasi konsumen, memeriksa catatan yang diproses, dan bereaksi terhadap reharding. KCL menangani semua subtugas ini sehingga Anda dapat memfokuskan upaya Anda pada penulisan logika pemrosesan catatan khusus Anda.
KCL berbeda dari API Kinesis Data Streams yang tersedia di AWS SDK. API Kinesis Data Streams membantu Anda mengelola banyak aspek Aliran Data Kinesis, termasuk membuat aliran, berbagi ulang, dan menempatkan dan mendapatkan catatan. KCL menyediakan lapisan abstraksi di sekitar semua subtugas ini, khususnya sehingga Anda dapat fokus pada logika pemrosesan data khusus aplikasi konsumen Anda. Untuk informasi tentang Kinesis Data Streams API, lihat Referensi API Amazon Kinesis.
penting
KCL adalah perpustakaan Java. Dukungan untuk bahasa selain Java disediakan menggunakan antarmuka multi-bahasa yang disebut. MultiLangDaemon Daemon ini Java-based dan berjalan di latar belakang ketika Anda menggunakan bahasa KCL selain Java. Misalnya, jika Anda menginstal KCL untuk Python dan menulis aplikasi konsumen Anda sepenuhnya dalam Python, Anda masih memerlukan Java diinstal pada sistem Anda karena. MultiLangDaemon Selanjutnya, MultiLangDaemon memiliki beberapa pengaturan default yang mungkin perlu Anda sesuaikan untuk kasus penggunaan Anda, misalnya, AWS wilayah yang terhubung. Untuk informasi lebih lanjut MultiLangDaemon tentang aktif GitHub, lihat MultiLangDaemon proyek https://github.com/awslabs/amazon-kinesis-client/tree/v1.x/src/main/java/com/amazonaws/services/kinesis/multilang
KCL bertindak sebagai perantara antara logika pemrosesan rekaman Anda dan Kinesis Data Streams.
KCL versi sebelumnya
Saat ini, Anda dapat menggunakan salah satu versi KCL yang didukung berikut untuk membangun aplikasi konsumen kustom Anda:
-
KCL 1.x
Untuk informasi selengkapnya, lihat Kembangkan konsumen KCL 1.x
-
KCL 2.x
Untuk informasi selengkapnya, lihat Kembangkan Konsumen KCL 2.x
Anda dapat menggunakan KCL 1.x atau KCL 2.x untuk membangun aplikasi konsumen yang menggunakan throughput bersama. Untuk informasi selengkapnya, lihat Kembangkan konsumen khusus dengan throughput bersama menggunakan KCL.
Untuk membangun aplikasi konsumen yang menggunakan throughput khusus (konsumen fan-out yang ditingkatkan), Anda hanya dapat menggunakan KCL 2.x. Untuk informasi selengkapnya, lihat Kembangkan konsumen fan-out yang ditingkatkan dengan throughput khusus.
Untuk informasi tentang perbedaan antara KCL 1.x dan KCL 2.x, dan petunjuk tentang cara bermigrasi dari KCL 1.x ke KCL 2.x, lihat. Migrasikan konsumen dari KCL 1.x ke KCL 2.x
Konsep KCL (versi sebelumnya)
-
Aplikasi konsumen KCL — aplikasi yang dibuat khusus menggunakan KCL dan dirancang untuk membaca dan memproses catatan dari aliran data.
-
Instans aplikasi konsumen - Aplikasi konsumen KCL biasanya didistribusikan, dengan satu atau lebih instance aplikasi berjalan secara bersamaan untuk mengoordinasikan kegagalan dan memproses data keseimbangan beban secara dinamis.
-
Wor ker — kelas tingkat tinggi yang digunakan instans aplikasi konsumen KCL untuk mulai memproses data.
penting
Setiap instance aplikasi konsumen KCL memiliki satu pekerja.
Pekerja menginisialisasi dan mengawasi berbagai tugas, termasuk menyinkronkan informasi pecahan dan sewa, melacak penugasan pecahan, dan memproses data dari pecahan. Seorang pekerja menyediakan KCL dengan informasi konfigurasi untuk aplikasi konsumen, seperti nama aliran data yang rekaman datanya aplikasi konsumen KCL ini akan diproses dan AWS kredenSIAL yang diperlukan untuk mengakses aliran data ini. Pekerja juga memulai instance aplikasi konsumen KCL tertentu untuk mengirimkan catatan data dari aliran data ke prosesor catatan.
penting
Di KCL 1.x kelas ini disebut Worker. Untuk informasi lebih lanjut, (ini adalah repositori Java KCL), lihat. https://github.com/awslabs/amazon-kinesis-client/blob/v1.x/src/main/java/com/amazonaws/services/kinesis/clientlibrary/lib/worker/Worker.java
Di KCL 2.x, kelas ini disebut Scheduler. Tujuan Scheduler di KCL 2.x identik dengan tujuan Pekerja di KCL 1.x. Untuk informasi selengkapnya tentang kelas Scheduler di KCL 2.x, lihat. https://github.com/awslabs/amazon-kinesis-client/blob/master/amazon-kinesis-client/src/main/java/software/amazon/kinesis/coordinator/Scheduler.java -
Sewa — data yang mendefinisikan ikatan antara pekerja dan pecahan. Aplikasi konsumen KCL terdistribusi menggunakan sewa untuk mempartisi pemrosesan catatan data di seluruh armada pekerja. Pada waktu tertentu, setiap pecahan catatan data terikat pada pekerja tertentu dengan sewa yang diidentifikasi oleh variabel leaseKey.
Secara default, seorang pekerja dapat memegang satu atau lebih sewa (tergantung pada nilai LeasesForWorker variabel maks) pada saat yang sama.
penting
Setiap pekerja akan bersaing untuk memegang semua sewa yang tersedia untuk semua pecahan yang tersedia dalam aliran data. Tetapi hanya satu pekerja yang berhasil memegang setiap sewa pada satu waktu.
Misalnya, jika Anda memiliki instance aplikasi konsumen A dengan pekerja A yang memproses aliran data dengan 4 pecahan, pekerja A dapat menahan sewa ke pecahan 1, 2, 3, dan 4 secara bersamaan. Tetapi jika Anda memiliki dua instance aplikasi konsumen: A dan B dengan pekerja A dan pekerja B, dan instans ini memproses aliran data dengan 4 pecahan, pekerja A dan pekerja B tidak dapat menahan sewa ke shard 1 secara bersamaan. Seorang pekerja memegang sewa untuk pecahan tertentu sampai siap untuk berhenti memproses catatan data pecahan ini atau sampai gagal. Ketika satu pekerja berhenti memegang sewa, pekerja lain mengambil dan menahan sewa.
Untuk informasi lebih lanjut, (ini adalah repositori Java KCL), lihat https://github.com/awslabs/amazon-kinesis-client/blob/v1.x/src/main/java/com/amazonaws/services/kinesis/leases/impl/Lease.java
untuk KCL 1.x dan https://github.com/awslabs/amazon-kinesis-client/blob/master/amazon-kinesis-client/src/main/java/software/amazon/kinesis/leases/Lease.java untuk KCL 2.x. -
Tabel sewa - tabel Amazon DynamoDB unik yang digunakan untuk melacak pecahan dalam aliran data KDS yang sedang disewa dan diproses oleh pekerja aplikasi konsumen KCL. Tabel sewa harus tetap sinkron (dalam pekerja dan di semua pekerja) dengan informasi shard terbaru dari aliran data saat aplikasi konsumen KCL sedang berjalan. Untuk informasi selengkapnya, lihat Gunakan tabel sewa untuk melacak pecahan yang diproses oleh aplikasi konsumen KCL.
-
Prosesor rekaman — logika yang mendefinisikan bagaimana aplikasi konsumen KCL Anda memproses data yang didapat dari aliran data. Saat runtime, instance aplikasi konsumen KCL membuat instance pekerja, dan pekerja ini membuat instance satu prosesor rekaman untuk setiap pecahan yang disewakan.
Gunakan tabel sewa untuk melacak pecahan yang diproses oleh aplikasi konsumen KCL
Topik
Apa itu tabel sewa
Untuk setiap aplikasi Amazon Kinesis Data Streams, KCL menggunakan tabel sewa unik (disimpan dalam tabel Amazon DynamoDB) untuk melacak pecahan dalam aliran data KDS yang sedang disewa dan diproses oleh pekerja aplikasi konsumen KCL.
penting
KCL menggunakan nama aplikasi konsumen untuk membuat nama tabel sewa yang digunakan aplikasi konsumen ini, oleh karena itu setiap nama aplikasi konsumen harus unik.
Anda dapat melihat tabel sewa menggunakan konsol Amazon DynamoDB saat aplikasi konsumen sedang berjalan.
Jika tabel sewa untuk aplikasi konsumen KCL Anda tidak ada saat aplikasi dimulai, salah satu pekerja membuat tabel sewa untuk aplikasi ini.
penting
Akun Anda dikenakan biaya untuk biaya yang terkait dengan tabel DynamoDB, selain biaya yang terkait dengan Kinesis Data Streams itu sendiri.
Setiap baris dalam tabel sewa mewakili pecahan yang sedang diproses oleh pekerja aplikasi konsumen Anda. Jika aplikasi konsumen KCL Anda hanya memproses satu aliran data, maka leaseKey yang merupakan kunci hash untuk tabel sewa adalah ID pecahan. Jika yaMemproses beberapa aliran data dengan KCL 2.x yang sama untuk aplikasi konsumen Java, maka struktur LeaseKey terlihat seperti ini:. account-id:StreamName:streamCreationTimestamp:ShardId Misalnya, 111111111:multiStreamTest-1:12345:shardId-000000000336.
Selain ID pecahan, setiap baris juga menyertakan data berikut:
-
pos pemeriksaan: Nomor urutan pos pemeriksaan terbaru untuk pecahan. Nilai ini unik di semua pecahan dalam aliran data.
-
checkpointSubSequenceNumber: Saat menggunakan fitur agregasi Kinesis Producer Library, ini adalah ekstensi ke pos pemeriksaan yang melacak catatan pengguna individu dalam catatan Kinesis.
-
LeaseCounter: Dig unakan untuk pembuatan versi sewa sehingga pekerja dapat mendeteksi bahwa sewa mereka telah diambil oleh pekerja lain.
-
LeaseKey: Pengi dentifikasi unik untuk sewa. Setiap sewa khusus untuk pecahan dalam aliran data dan dipegang oleh satu pekerja pada satu waktu.
-
LeaseOwner: Pekerja yang memegang sewa ini.
-
pemilikSwitchesSinceCheckpoint: Berapa kali sewa ini telah mengubah pekerja sejak terakhir kali pos pemeriksaan ditulis.
-
parentShardId: Digunakan untuk memastikan bahwa pecahan induk diproses sepenuhnya sebelum pemrosesan dimulai pada pecahan anak. Ini memastikan bahwa catatan diproses dalam urutan yang sama dengan yang dimasukkan ke dalam aliran.
-
hashrange: Dig unakan oleh
PeriodicShardSyncManageruntuk menjalankan sinkronisasi berkala untuk menemukan pecahan yang hilang dalam tabel sewa dan membuat sewa untuk mereka jika diperlukan.catatan
Data ini hadir dalam tabel sewa untuk setiap pecahan dimulai dengan KCL 1.14 dan KCL 2.3. Untuk informasi selengkapnya tentang
PeriodicShardSyncManagerdan sinkronisasi berkala antara sewa dan pecahan, lihat. Bagaimana tabel sewa disinkronkan dengan pecahan dalam aliran data Kinesis -
childshards: Digunakan oleh
LeaseCleanupManageruntuk meninjau status pemrosesan pecahan anak dan memutuskan apakah pecahan induk dapat dihapus dari tabel sewa.catatan
Data ini hadir dalam tabel sewa untuk setiap pecahan dimulai dengan KCL 1.14 dan KCL 2.3.
-
ShardID: ID pecahan.
catatan
Data ini hanya ada di tabel sewa jika Anda adaMemproses beberapa aliran data dengan KCL 2.x yang sama untuk aplikasi konsumen Java. Ini hanya didukung di KCL 2.x untuk Java, dimulai dengan KCL 2.3 untuk Java dan yang lebih baru.
-
nama aliran Pengidentifikasi aliran data dalam format berikut:
account-id:StreamName:streamCreationTimestamp.catatan
Data ini hanya ada di tabel sewa jika Anda adaMemproses beberapa aliran data dengan KCL 2.x yang sama untuk aplikasi konsumen Java. Ini hanya didukung di KCL 2.x untuk Java, dimulai dengan KCL 2.3 untuk Java dan yang lebih baru.
Throughput
Jika aplikasi Amazon Kinesis Data Streams menerima pengecualian throughput yang disediakan, Anda harus meningkatkan throughput yang disediakan untuk tabel DynamoDB. KCL membuat tabel dengan throughput yang disediakan 10 pembacaan per detik dan 10 penulisan per detik, tetapi ini mungkin tidak cukup untuk aplikasi Anda. Misalnya, jika aplikasi Amazon Kinesis Data Streams sering melakukan pemeriksaan atau beroperasi pada aliran yang terdiri dari banyak pecahan, Anda mungkin memerlukan lebih banyak throughput.
Untuk informasi tentang throughput yang disediakan di DynamoDB, lihat Mode Read/Write Kapasitas dan Bek erja dengan Tabel dan Data di Panduan Pengembang Amazon DynamoDB.
Bagaimana tabel sewa disinkronkan dengan pecahan dalam aliran data Kinesis
Pekerja di aplikasi konsumen KCL menggunakan sewa untuk memproses pecahan dari aliran data tertentu. Informasi tentang pekerja mana yang menyewa pecahan apa pada waktu tertentu disimpan dalam tabel sewa. Tabel sewa harus tetap sinkron dengan informasi shard terbaru dari aliran data saat aplikasi konsumen KCL sedang berjalan. KCL menyinkronkan tabel sewa dengan informasi pecahan yang diperoleh dari layanan Kinesis Data Streams selama bootstraping aplikasi konsumen (baik ketika aplikasi konsumen diinisialisasi atau dimulai ulang) dan juga setiap kali pecahan yang sedang diproses mencapai akhir (resharding). Dengan kata lain, pekerja atau aplikasi konsumen KCL disinkronkan dengan aliran data yang mereka proses selama bootstrap aplikasi konsumen awal dan setiap kali aplikasi konsumen menemukan peristiwa hard ulang aliran data.
Topik
Sinkronisasi di KCL 1.0 - 1.13 dan KCL 2.0 - 2.2
Di KCL 1.0 - 1.13 dan KCL 2.0 - 2.2, selama bootstraping aplikasi konsumen dan juga selama setiap peristiwa hard ulang aliran data, KCL menyinkronkan tabel sewa dengan informasi pecahan yang diperoleh dari layanan Kinesis Data Streams dengan memanggil atau API penemuan. ListShards DescribeStream Di semua versi KCL yang tercantum di atas, setiap pekerja aplikasi konsumen KCL menyelesaikan langkah-langkah berikut untuk melakukan proses lease/shard sinkronisasi selama bootstrap aplikasi konsumen dan pada setiap peristiwa hard ulang aliran:
-
Mengambil semua pecahan untuk data aliran yang sedang diproses
-
Mengambil semua sewa pecahan dari tabel sewa
-
Menyaring setiap pecahan terbuka yang tidak memiliki sewa di tabel sewa
-
Mengulangi semua pecahan terbuka yang ditemukan dan untuk setiap pecahan terbuka tanpa induk terbuka:
-
Melintasi pohon hierarki melalui jalur leluhurnya untuk menentukan apakah pecahan itu adalah keturunan. Pecahan dianggap sebagai keturunan, jika pecahan leluhur sedang diproses (entri sewa untuk pecahan leluhur ada di tabel sewa) atau jika pecahan leluhur harus diproses (misalnya, jika posisi awal adalah atau)
TRIM_HORIZONAT_TIMESTAMP -
Jika pecahan terbuka dalam konteks adalah keturunan, KCL memeriksa pecahan berdasarkan posisi awal dan membuat sewa untuk induknya, jika diperlukan
-
Sinkronisasi di KCL 2.x, dimulai dengan KCL 2.3 dan yang lebih baru
Dimulai dengan versi terbaru yang didukung KCL 2.x (KCL 2.3) dan yang lebih baru, pustaka sekarang mendukung perubahan berikut pada proses sinkronisasi. Perubahan lease/shard sinkronisasi ini secara signifikan mengurangi jumlah panggilan API yang dilakukan oleh aplikasi konsumen KCL ke layanan Kinesis Data Streams dan mengoptimalkan manajemen sewa di aplikasi konsumen KCL Anda.
-
Selama bootstraping aplikasi, jika tabel sewa kosong, KCL menggunakan opsi pemfilteran
ListShardAPI (parameter permintaanShardFilteropsional) untuk mengambil dan membuat sewa hanya untuk snapshot pecahan yang terbuka pada waktu yang ditentukan oleh parameter.ShardFilterShardFilterParameter ini memungkinkan Anda untuk menyaring responsListShardsAPI. Satu-satunya properti yang diperlukan dariShardFilterparameter adalahType. KCL menggunakan propertiTypefilter dan nilai-nilai validnya berikut ini untuk mengidentifikasi dan mengembalikan snapshot pecahan terbuka yang mungkin memerlukan sewa baru:-
AT_TRIM_HORIZON- respon mencakup semua pecahan yang terbuka diTRIM_HORIZON. -
AT_LATEST- respon hanya mencakup pecahan aliran data yang saat ini terbuka. -
AT_TIMESTAMP- respons mencakup semua pecahan yang cap waktu awalnya kurang dari atau sama dengan stempel waktu yang diberikan dan stempel waktu akhir lebih besar dari atau sama dengan stempel waktu yang diberikan atau masih terbuka.
ShardFilterdigunakan saat membuat sewa untuk tabel sewa kosong untuk menginisialisasi sewa untuk snapshot pecahan yang ditentukan di.RetrievalConfig#initialPositionInStreamExtendedUntuk informasi selengkapnya tentang
ShardFilter, lihat https://docs.aws.amazon.com/kinesis/latest/APIReference/API_ShardFilter.html. -
-
Alih-alih semua pekerja melakukan lease/shard sinkronisasi untuk menjaga tabel sewa tetap up to date dengan pecahan terbaru dalam aliran data, satu pemimpin pekerja terpilih melakukan lease/shard sinkronisasi.
-
KCL 2.3 menggunakan parameter
ChildShardspengembalianGetRecordsdanSubscribeToShardAPI untuk melakukan lease/shard sinkronisasi yang terjadi padaSHARD_ENDpecahan tertutup, memungkinkan pekerja KCL hanya membuat sewa untuk pecahan anak dari pecahan yang selesai diprosesnya. Untuk dibagikan di seluruh aplikasi konsumen, pengoptimalan lease/shard sinkronisasi ini menggunakanChildShardsparameterGetRecordsAPI. Untuk aplikasi konsumen throughput khusus (fan-out yang disempurnakan), pengoptimalan lease/shard sinkronisasi ini menggunakanChildShardsparameter API.SubscribeToShardLihat informasi selengkapnya di GetRecords, SubscribeToShards, dan ChildShard. -
Dengan perubahan di atas, perilaku KCL bergerak dari model semua pekerja yang belajar tentang semua pecahan yang ada ke model pekerja yang hanya belajar tentang pecahan anak-anak dari pecahan yang dimiliki setiap pekerja. Oleh karena itu, selain sinkronisasi yang terjadi selama peristiwa bootstraping dan reshard aplikasi konsumen, KCL sekarang juga melakukan shard/lease pemindaian berkala tambahan untuk mengidentifikasi lubang potensial dalam tabel sewa (dengan kata lain, untuk mempelajari semua pecahan baru) untuk memastikan rentang hash lengkap dari aliran data sedang diproses dan membuat sewa untuk mereka jika diperlukan.
PeriodicShardSyncManageradalah komponen yang bertanggung jawab untuk menjalankan lease/shard pemindaian berkala.Untuk informasi lebih lanjut tentang
PeriodicShardSyncManagerdi KCL 2.3, lihat https://github.com/awslabs/amazon-kinesis-client/blob/master/amazon-kinesis-client/src/main/java/software/amazon/kinesis/leases/LeaseManagementConfig.java # L201-L213. Di KCL 2.3, opsi konfigurasi baru tersedia untuk dikonfigurasi
PeriodicShardSyncManagerdi:LeaseManagementConfigNama Nilai default Deskripsi menyewa RecoveryAuditorExecutionFrequencyMillis 120000 (2 menit)
Frekuensi (dalam millis) pekerjaan auditor untuk memindai sewa parSIAL dalam tabel sewa. Jika auditor mendeteksi adanya lubang dalam sewa untuk aliran, maka itu akan memicu sinkronisasi pecahan berdasarkan.
leasesRecoveryAuditorInconsistencyConfidenceThresholdmenyewa RecoveryAuditorInconsistencyConfidenceThreshold 3
Ambang kepercayaan untuk pekerjaan auditor periodik untuk menentukan apakah sewa untuk aliran data dalam tabel sewa tidak konsisten. Jika auditor menemukan kumpulan inkonsistensi yang sama secara berurutan untuk aliran data untuk ini berkali-kali, maka itu akan memicu sinkronisasi pecahan.
CloudWatch Metrik baru juga sekarang dipancarkan untuk memantau kesehatan
PeriodicShardSyncManager. Untuk informasi selengkapnya, lihat PeriodicShardSyncManager. -
Termasuk pengoptimalan
HierarchicalShardSynceruntuk hanya membuat sewa untuk satu lapisan pecahan.
Sinkronisasi di KCL 1.x, dimulai dengan KCL 1.14 dan yang lebih baru
Dimulai dengan versi terbaru yang didukung KCL 1.x (KCL 1.14) dan yang lebih baru, pustaka sekarang mendukung perubahan berikut pada proses sinkronisasi. Perubahan lease/shard sinkronisasi ini secara signifikan mengurangi jumlah panggilan API yang dilakukan oleh aplikasi konsumen KCL ke layanan Kinesis Data Streams dan mengoptimalkan manajemen sewa di aplikasi konsumen KCL Anda.
-
Selama bootstraping aplikasi, jika tabel sewa kosong, KCL menggunakan opsi pemfilteran
ListShardAPI (parameter permintaanShardFilteropsional) untuk mengambil dan membuat sewa hanya untuk snapshot pecahan yang terbuka pada waktu yang ditentukan oleh parameter.ShardFilterShardFilterParameter ini memungkinkan Anda untuk menyaring responsListShardsAPI. Satu-satunya properti yang diperlukan dariShardFilterparameter adalahType. KCL menggunakan propertiTypefilter dan nilai-nilai validnya berikut ini untuk mengidentifikasi dan mengembalikan snapshot pecahan terbuka yang mungkin memerlukan sewa baru:-
AT_TRIM_HORIZON- respon mencakup semua pecahan yang terbuka diTRIM_HORIZON. -
AT_LATEST- respon hanya mencakup pecahan aliran data yang saat ini terbuka. -
AT_TIMESTAMP- respons mencakup semua pecahan yang cap waktu awalnya kurang dari atau sama dengan stempel waktu yang diberikan dan stempel waktu akhir lebih besar dari atau sama dengan stempel waktu yang diberikan atau masih terbuka.
ShardFilterdigunakan saat membuat sewa untuk tabel sewa kosong untuk menginisialisasi sewa untuk snapshot pecahan yang ditentukan di.KinesisClientLibConfiguration#initialPositionInStreamExtendedUntuk informasi selengkapnya tentang
ShardFilter, lihat https://docs.aws.amazon.com/kinesis/latest/APIReference/API_ShardFilter.html. -
-
Alih-alih semua pekerja melakukan lease/shard sinkronisasi untuk menjaga tabel sewa tetap up to date dengan pecahan terbaru dalam aliran data, satu pemimpin pekerja terpilih melakukan lease/shard sinkronisasi.
-
KCL 1.14 menggunakan parameter
ChildShardspengembalianGetRecordsdanSubscribeToShardAPI untuk melakukan lease/shard sinkronisasi yang terjadi padaSHARD_ENDpecahan tertutup, memungkinkan pekerja KCL hanya membuat sewa untuk pecahan anak dari pecahan yang selesai diprosesnya. Untuk informasi selengkapnya, lihat GetRecords dan ChildShard. -
Dengan perubahan di atas, perilaku KCL bergerak dari model semua pekerja yang belajar tentang semua pecahan yang ada ke model pekerja yang hanya belajar tentang pecahan anak-anak dari pecahan yang dimiliki setiap pekerja. Oleh karena itu, selain sinkronisasi yang terjadi selama peristiwa bootstraping dan reshard aplikasi konsumen, KCL sekarang juga melakukan shard/lease pemindaian berkala tambahan untuk mengidentifikasi lubang potensial dalam tabel sewa (dengan kata lain, untuk mempelajari semua pecahan baru) untuk memastikan rentang hash lengkap dari aliran data sedang diproses dan membuat sewa untuk mereka jika diperlukan.
PeriodicShardSyncManageradalah komponen yang bertanggung jawab untuk menjalankan lease/shard pemindaian berkala.Ketika
KinesisClientLibConfiguration#shardSyncStrategyTypedisetel keShardSyncStrategyType.SHARD_END,PeriodicShardSync leasesRecoveryAuditorInconsistencyConfidenceThresholddigunakan untuk menentukan ambang batas untuk jumlah pemindaian berturut-turut yang berisi lubang di tabel sewa setelah itu untuk menerapkan sinkronisasi pecahan. KetikaKinesisClientLibConfiguration#shardSyncStrategyTypedisetel keShardSyncStrategyType.PERIODIC,leasesRecoveryAuditorInconsistencyConfidenceThresholddiabaikan.Untuk informasi lebih lanjut tentang
PeriodicShardSyncManagerdi KCL 1.14, lihat https://github.com/awslabs/amazon-kinesis-client/blob/v1.x/src/main/java/com/amazonaws/services/kinesis/clientlibrary/lib/worker/KinesisClientLibConfiguration.java #. L987-L999Di KCL 1.14, opsi konfigurasi baru tersedia untuk dikonfigurasi
PeriodicShardSyncManagerdi:LeaseManagementConfigNama Nilai default Deskripsi menyewa RecoveryAuditorInconsistencyConfidenceThreshold 3
Ambang kepercayaan untuk pekerjaan auditor periodik untuk menentukan apakah sewa untuk aliran data dalam tabel sewa tidak konsisten. Jika auditor menemukan kumpulan inkonsistensi yang sama secara berurutan untuk aliran data untuk ini berkali-kali, maka itu akan memicu sinkronisasi pecahan.
CloudWatch Metrik baru juga sekarang dipancarkan untuk memantau kesehatan
PeriodicShardSyncManager. Untuk informasi selengkapnya, lihat PeriodicShardSyncManager. -
KCL 1.14 sekarang juga mendukung pembersihan sewa yang ditangguhkan. Sewa dihapus secara asinkron
LeaseCleanupManagersetelah mencapaiSHARD_END, ketika pecahan telah kedaluwarsa melewati periode retensi aliran data atau ditutup sebagai hasil dari operasi resharding.Opsi konfigurasi baru tersedia untuk dikonfigurasi
LeaseCleanupManager.Nama Nilai default Deskripsi menyewa CleanupIntervalMillis 1 menit
Interval untuk menjalankan thread pembersihan sewa.
selesai LeaseCleanupIntervalMillis 5 menit Interval untuk memeriksa apakah sewa selesai atau tidak.
sampah LeaseCleanupIntervalMillis 30 menit Interval untuk memeriksa apakah sewa adalah sampah (yaitu dipangkas melewati periode retensi aliran data) atau tidak.
-
Termasuk pengoptimalan
KinesisShardSynceruntuk hanya membuat sewa untuk satu lapisan pecahan.
Memproses beberapa aliran data dengan KCL 2.x yang sama untuk aplikasi konsumen Java
Bagian ini menjelaskan perubahan berikut di KCL 2.x untuk Java yang memungkinkan Anda membuat aplikasi konsumen KCL yang dapat memproses lebih dari satu aliran data pada saat yang sama.
penting
Pemrosesan multistream hanya didukung di KCL 2.x untuk Java, dimulai dengan KCL 2.3 untuk Java dan yang lebih baru.
Pemrosesan multistream TIDAK didukung untuk bahasa lain di mana KCL 2.x dapat diimplementasikan.
Pemrosesan multistream TIDAK didukung dalam versi KCL 1.x apa pun.
-
MultistreamTracker antarmuka
Untuk membangun aplikasi konsumen yang dapat memproses beberapa aliran pada saat yang sama, Anda harus menerapkan antarmuka baru yang disebut MultistreamTracker
. Antarmuka ini mencakup streamConfigListmetode yang mengembalikan daftar aliran data dan konfigurasinya untuk diproses oleh aplikasi konsumen KCL. Perhatikan bahwa aliran data yang sedang diproses dapat diubah selama runtime aplikasi konsumen.streamConfigListdipanggil secara berkala oleh KCL untuk mempelajari tentang perubahan aliran data untuk diproses.streamConfigListMetode mengisi StreamConfigdaftar. package software.amazon.kinesis.common; import lombok.Data; import lombok.experimental.Accessors; @Data @Accessors(fluent = true) public class StreamConfig { private final StreamIdentifier streamIdentifier; private final InitialPositionInStreamExtended initialPositionInStreamExtended; private String consumerArn; }Perhatikan bahwa kolom
StreamIdentifierdanInitialPositionInStreamExtendedwajib diisi, sementaraconsumerArnbersifat opsional. Anda harus menyediakanconsumerArnsatu-satunya jika Anda menggunakan KCL 2.x untuk mengimplementasikan aplikasi konsumen fan-out yang disempurnakan.Untuk informasi selengkapnya
StreamIdentifier, lihat https://github.com/awslabs/amazon-kinesis-client/blob/v2.5.8/amazon-kinesis-client/src/main/java/software/amazon/kinesis/common/StreamIdentifier.java #L129. Untuk membuat StreamIdentifier, kami sarankan Anda membuat instance multistream daristreamArndan yang tersedia di v2.5.0 danstreamCreationEpochyang lebih baru. Di KCL v2.3 dan v2.4, yang tidak mendukungstreamArm, buat instance multistream dengan menggunakan format.account-id:StreamName:streamCreationTimestampFormat ini tidak akan digunakan lagi dan tidak lagi didukung mulai dengan rilis besar berikutnya.MultistreamTrackerjuga mencakup strategi untuk menghapus sewa aliran lama di tabel sewa (formerStreamsLeasesDeletionStrategy). Perhatikan bahwa strategi TIDAK DAPAT diubah selama runtime aplikasi konsumen. Untuk informasi selengkapnya, lihat https://github.com/awslabs/amazon-kinesis-client/blob/0c5042dadf794fe988438436252a5a8fe70b6b0b/amazon-kinesis-client/src/main/java/software/amazon/kinesis/processor/FormerStreamsLeasesDeletionStrategy.java -
ConfigsBuilder
adalah kelas seluruh aplikasi yang dapat Anda gunakan untuk menentukan semua pengaturan konfigurasi KCL 2.x yang akan digunakan saat membangun aplikasi konsumen KCL Anda. ConfigsBuilderkelas sekarang memiliki dukungan untukMultistreamTrackerantarmuka. Anda dapat mengin ConfigsBuilder isialisasi salah satu dengan nama satu aliran data untuk menggunakan catatan dari:/** * Constructor to initialize ConfigsBuilder with StreamName * @param streamName * @param applicationName * @param kinesisClient * @param dynamoDBClient * @param cloudWatchClient * @param workerIdentifier * @param shardRecordProcessorFactory */ public ConfigsBuilder(@NonNull String streamName, @NonNull String applicationName, @NonNull KinesisAsyncClient kinesisClient, @NonNull DynamoDbAsyncClient dynamoDBClient, @NonNull CloudWatchAsyncClient cloudWatchClient, @NonNull String workerIdentifier, @NonNull ShardRecordProcessorFactory shardRecordProcessorFactory) { this.appStreamTracker = Either.right(streamName); this.applicationName = applicationName; this.kinesisClient = kinesisClient; this.dynamoDBClient = dynamoDBClient; this.cloudWatchClient = cloudWatchClient; this.workerIdentifier = workerIdentifier; this.shardRecordProcessorFactory = shardRecordProcessorFactory; }Atau Anda dapat menginisialisasi ConfigsBuilder dengan
MultiStreamTrackerjika Anda ingin mengimplementasikan aplikasi konsumen KCL yang memproses beberapa aliran secara bersamaan.* Constructor to initialize ConfigsBuilder with MultiStreamTracker * @param multiStreamTracker * @param applicationName * @param kinesisClient * @param dynamoDBClient * @param cloudWatchClient * @param workerIdentifier * @param shardRecordProcessorFactory */ public ConfigsBuilder(@NonNull MultiStreamTracker multiStreamTracker, @NonNull String applicationName, @NonNull KinesisAsyncClient kinesisClient, @NonNull DynamoDbAsyncClient dynamoDBClient, @NonNull CloudWatchAsyncClient cloudWatchClient, @NonNull String workerIdentifier, @NonNull ShardRecordProcessorFactory shardRecordProcessorFactory) { this.appStreamTracker = Either.left(multiStreamTracker); this.applicationName = applicationName; this.kinesisClient = kinesisClient; this.dynamoDBClient = dynamoDBClient; this.cloudWatchClient = cloudWatchClient; this.workerIdentifier = workerIdentifier; this.shardRecordProcessorFactory = shardRecordProcessorFactory; } -
Dengan dukungan multistream yang diterapkan untuk aplikasi konsumen KCL Anda, setiap baris tabel sewa aplikasi sekarang berisi ID pecahan dan nama aliran dari beberapa aliran data yang diproses aplikasi ini.
-
Ketika dukungan multistream untuk aplikasi konsumen KCL Anda diimplementasikan, LeaseKey mengambil struktur berikut:.
account-id:StreamName:streamCreationTimestamp:ShardIdMisalnya,111111111:multiStreamTest-1:12345:shardId-000000000336.penting
Ketika aplikasi konsumen KCL Anda yang ada dikonfigurasi untuk memproses hanya satu aliran data, leaseKey (yang merupakan kunci hash untuk tabel sewa) adalah ID pecahan. Jika Anda mengonfigurasi ulang aplikasi konsumen KCL yang ada ini untuk memproses beberapa aliran data, itu merusak tabel sewa Anda, karena dengan dukungan multistream, struktur LeaseKey harus sebagai berikut:.
account-id:StreamName:StreamCreationTimestamp:ShardId
Gunakan KCL dengan AWS Glue Registri Skema
Anda dapat mengintegrasikan aliran data Kinesis Anda dengan AWS Glue Schema Registry. Registry S AWS Glue chema memungkinkan Anda menemukan, mengontrol, dan mengembangkan skema secara terpusat, sambil memastikan data yang dihasilkan terus divalidasi oleh skema terdaftar. Sebuah skema mendefinisikan struktur dan format catatan data. Sebuah skema adalah sebuah spesifikasi berversi untuk publikasi data yang handal, konsumsi, atau penyimpanan. Registry S AWS Glue chema memungkinkan Anda meningkatkan kualitas data end-to-end dan tata kelola data dalam aplikasi streaming Anda. Untuk informasi selengkapnya, lihat AWS Glue Schema Registry. Salah satu cara untuk mengatur integrasi ini adalah melalui KCL di Java.
penting
Saat ini, integrasi Kinesis Data Streams dan AWS Glue Schema Registry hanya didukung untuk aliran data Kinesis yang menggunakan konsumen KCL 2.3 yang diimplementasikan di Java. Multi-language dukungan tidak diberikan. Konsumen KCL 1.0 tidak didukung. Konsumen KCL 2.x sebelum KCL 2.3 tidak didukung.
Untuk petunjuk terperinci tentang cara mengatur integrasi Aliran Data Kinesis dengan Schema Registry menggunakan KCL, lihat bagian “Berinteraksi dengan Data Menggunakan Per KPL/KCL pustakaan” dalam Kasus Penggunaan: Mengintegrasikan Amazon Kinesis Data Streams dengan AWS Glue Schema Registry.