Die vorliegende Übersetzung wurde maschinell erstellt. Im Falle eines Konflikts oder eines Widerspruchs zwischen dieser übersetzten Fassung und der englischen Fassung (einschließlich infolge von Verzögerungen bei der Übersetzung) ist die englische Fassung maßgeblich.
Multi-stream Verarbeitung mit KCL
In diesem Abschnitt werden die erforderlichen Änderungen in KCL beschrieben, die es Ihnen ermöglichen, KCL-Verbraucheranwendungen zu erstellen, die mehr als einen Datenstrom gleichzeitig verarbeiten können.
Wichtig
-
Multi-stream Die Verarbeitung wird nur in KCL 2.3 oder höher unterstützt.
-
Multi-stream Die Verarbeitung wird für KCL-Benutzer nicht unterstützt, die in anderen Sprachen geschrieben sind als Java, die mit laufen.
multilangdaemon -
Multi-stream Die Verarbeitung wird in keiner Version von KCL 1.x unterstützt.
-
MultistreamTracker Schnittstelle
-
Um eine Verbraucheranwendung zu erstellen, die mehrere Streams gleichzeitig verarbeiten kann, müssen Sie eine neue Schnittstelle mit dem Namen implementieren MultistreamTracker
. Diese Schnittstelle enthält die streamConfigList-Methode, die die Liste der Datenströme und ihrer Konfigurationen zurückgibt, die von der KCL-Konsumentenanwendung verarbeitet werden sollen. Beachten Sie, dass die Datenströme, die verarbeitet werden, während der Laufzeit der Consumer-Anwendung geändert werden können.streamConfigListwird regelmäßig von KCL aufgerufen, um sich über die Änderungen der zu verarbeitenden Datenströme zu informieren. -
Der
streamConfigListfüllt die Liste aus. StreamConfig
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; }-
Die Felder
StreamIdentifierundInitialPositionInStreamExtendedsind Pflichtfelder, obwohl sie optionalconsumerArnsind. Sie müssen das FeldconsumerArnnur angeben, wenn Sie KCL verwenden, um eine erweiterte Fan-Out-Verbraucheranwendung zu implementieren. -
Weitere Informationen zu finden Sie
StreamIdentifierunter https://github.com/awslabs/amazon-kinesis-client/blob/v2.5.8/amazon-kinesis-client/src/main/java/software/amazon/kinesis/common/StreamIdentifier.java #L129.Um eine zu erstellen StreamIdentifier, empfehlen wir, dass Sie eine Multistream-Instanz ausstreamArnund der erstellenstreamCreationEpoch, die in KCL 2.5.0 oder höher verfügbar ist. Erstellen Sie in KCL v2.3 und v2.4, die dies nicht unterstützenstreamArm, eine Multistream-Instanz mithilfe des Formats.account-id:StreamName:streamCreationTimestampDieses Format ist veraltet und wird ab der nächsten Hauptversion nicht mehr unterstützt. -
MultistreamTracker beinhaltet auch eine Strategie zum Löschen von Leasings alter Streams in der Leasing-Tabelle (früher). StreamsLeasesDeletionStrategy Beachten Sie, dass die Strategie während der Laufzeit der Konsumentenanwendung NICHT geändert werden kann. Weitere Informationen finden Sie unter https://github.com/awslabs/amazon-kinesis-client/blob/0c5042dadf794fe988438436252a5a8fe70b6b0b/amazon-kinesis-client/src/main/java/software/amazon/kinesis/processor/FormerStreamsLeasesDeletionStrategy.java
.
-
Oder Sie können ConfigsBuilder mit initialisieren, MultiStreamTracker wenn Sie eine KCL-Verbraucheranwendung implementieren möchten, die mehrere Streams gleichzeitig verarbeitet.
* 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; }
-
Da für Ihre KCL-Consumer-Anwendung die Unterstützung mehrerer Streams implementiert ist, enthält jede Zeile der Leasetabelle der Anwendung jetzt die Shard-ID und den Streamnamen der verschiedenen Datenströme, die diese Anwendung verarbeitet.
-
Wenn die Multistream-Unterstützung für Ihre KCL-Verbraucheranwendung implementiert ist, hat der LeaseKey die folgende Struktur:.
account-id:StreamName:streamCreationTimestamp:ShardIdBeispiel,111111111:multiStreamTest-1:12345:shardId-000000000336.
Wichtig
Wenn Ihre bestehende KCL-Verbraucheranwendung so konfiguriert ist, dass sie nur einen Datenstrom verarbeitet, ist leaseKey (das ist der Partitionsschlüssel für die Leasetabelle) die Shard-ID. Wenn Sie eine bestehende KCL-Verbraucheranwendung für die Verarbeitung mehrerer Datenströme neu konfigurieren, wird Ihre Leasetabelle beschädigt, da die leaseKey Struktur wie folgt aussehen muss: account-id:StreamName:StreamCreationTimestamp:ShardId um mehrere Datenströme zu unterstützen.