Terjemahan disediakan oleh mesin penerjemah. Jika konten terjemahan yang diberikan bertentangan dengan versi bahasa Inggris aslinya, utamakan versi bahasa Inggris.
Multi-stream Pengolahan dengan KCL
Bagian ini menjelaskan perubahan yang diperlukan dalam KCL yang memungkinkan Anda membuat aplikasi konsumen KCL yang dapat memproses lebih dari satu aliran data secara bersamaan.
penting
-
Multi-stream pemrosesan hanya didukung di KCL 2.3 atau yang lebih baru.
-
Multi-stream pemrosesan tidak didukung untuk konsumen KCL yang ditulis dalam bahasa non-Java yang dijalankan dengan.
multilangdaemon -
Multi-stream pemrosesan tidak didukung dalam versi KCL 1.x manapun.
-
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. -
I
streamConfigListsi 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; }-
Bid
StreamIdentifierang danInitialPositionInStreamExtendedwajib diisi, sementaraconsumerArnbersifat opsional. Anda harus menyediakanconsumerArnsatu-satunya jika Anda menggunakan KCL 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, sebaiknya buat instance multistream daristreamArndan yang tersedia di KCL 2.5.0 ataustreamCreationEpochyang 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. -
MultistreamTracker juga mencakup strategi untuk menghapus sewa aliran lama di tabel sewa (sebelumnyaStreamsLeasesDeletionStrategy). 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
.
-
Atau Anda dapat menginisialisasi ConfigsBuilder dengan MultiStreamTracker jika 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 multi-aliran 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 multi-stream 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 yang ada dikonfigurasi untuk memproses hanya satu aliran data, leaseKey (yang merupakan kunci partisi untuk tabel sewa) adalah ID pecahan. Jika Anda mengkonfigurasi ulang aplikasi konsumen KCL yang ada untuk memproses beberapa aliran data, itu merusak tabel sewa Anda, karena leaseKey strukturnya harus sebagai berikut: account-id:StreamName:StreamCreationTimestamp:ShardId untuk mendukung multi-stream.