View a markdown version of this page

Entwickeln Sie Verbraucher mit KCL in Java - Amazon Kinesis Data Streams

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.

Entwickeln Sie Verbraucher mit KCL in Java

Voraussetzungen

Bevor Sie mit der Verwendung von KCL 3.x beginnen, stellen Sie sicher, dass Sie über Folgendes verfügen:

  • Java Development Kit (JDK) 8 oder höher

  • AWS SDK für Java 2.x

  • Maven oder Gradle für das Abhängigkeitsmanagement

KCL erfasst Kennzahlen zur CPU-Auslastung, wie z. B. die CPU-Auslastung des Rechenhosts, auf dem die Workers gerade laufen, um die Auslastung auszugleichen und eine gleichmäßige Ressourcenauslastung aller Mitarbeiter zu erreichen. Damit KCL Kennzahlen zur CPU-Auslastung von Workern erfassen kann, müssen Sie die folgenden Voraussetzungen erfüllen:

Amazon Elastic Compute Cloud(Amazon EC2)

  • Ihr Betriebssystem muss Linux OS sein.

  • Sie müssen IMDSv2 in Ihrer EC2-Instance aktivieren.

Amazon Elastic Container Service (Amazon ECS) auf Amazon EC2

Amazon ECS aktiviert AWS Fargate

  • Sie müssen Version 4 des Fargate-Endpunkts für Aufgabenmetadaten aktivieren. Wenn Sie die Fargate-Plattformversion 1.4.0 oder höher verwenden, ist dies standardmäßig aktiviert.

  • Fargate-Plattformversion 1.4.0 oder höher.

Amazon Elastic Kubernetes Service (Amazon EKS) auf Amazon EC2

  • Ihr Betriebssystem muss Linux OS sein.

Amazon EKS an AWS Fargate

  • Fargate-Plattform 1.3.0 oder höher.

Wichtig

Wenn KCL die CPU-Auslastungskennzahlen von Mitarbeitern nicht erfassen kann, verwendet KCL den Durchsatz pro Mitarbeiter, um Leasingverträge zuzuweisen und die Auslastung der Mitarbeiter in der Flotte zu verteilen. Weitere Informationen finden Sie unter Wie KCL den Arbeitern Leasingverträge zuweist und die Last ausgleicht.

Abhängigkeiten installieren und hinzufügen

Wenn Sie Maven verwenden, fügen Sie Ihrer pom.xml Datei die folgende Abhängigkeit hinzu. Stellen Sie sicher, dass Sie 3.x.x durch die neueste KCL-Version ersetzt haben.

<dependency> <groupId>software.amazon.kinesis</groupId> <artifactId>amazon-kinesis-client</artifactId> <version>3.x.x</version> <!-- Use the latest version --> </dependency>

Wenn Sie Gradle verwenden, fügen Sie Ihrer Datei Folgendes hinzu. build.gradle Stellen Sie sicher, dass Sie 3.x.x durch die neueste KCL-Version ersetzt haben.

implementation 'software.amazon.kinesis:amazon-kinesis-client:3.x.x'

Sie können im Maven Central Repository nach der neuesten Version der KCL suchen. https://search.maven.org/artifact/software.amazon.kinesis/amazon-kinesis-client

Implementieren Sie den Consumer

Eine KCL-Verbraucheranwendung besteht aus den folgenden Schlüsselkomponenten:

RecordProcessor

RecordProcessor ist die Kernkomponente, in der sich Ihre Geschäftslogik für die Verarbeitung von Kinesis-Datenstromdatensätzen befindet. Sie definiert, wie Ihre Anwendung die Daten verarbeitet, die sie vom Kinesis-Stream empfängt.

Wichtigste Aufgaben:

  • Initialisieren Sie die Verarbeitung für einen Shard

  • Verarbeiten Sie Stapel von Datensätzen aus dem Kinesis-Stream

  • Beendet die Verarbeitung für einen Shard (z. B. wenn der Shard geteilt oder zusammengeführt wird oder das Leasing an einen anderen Host übergeben wird)

  • Übernehmen Sie Checkpoints, um den Fortschritt zu verfolgen

Das Folgende zeigt ein Implementierungsbeispiel:

import org.slf4j.Logger; import org.slf4j.LoggerFactory; import org.slf4j.MDC; import software.amazon.kinesis.exceptions.InvalidStateException; import software.amazon.kinesis.exceptions.ShutdownException; import software.amazon.kinesis.lifecycle.events.*; import software.amazon.kinesis.processor.ShardRecordProcessor; public class SampleRecordProcessor implements ShardRecordProcessor { private static final String SHARD_ID_MDC_KEY = "ShardId"; private static final Logger log = LoggerFactory.getLogger(SampleRecordProcessor.class); private String shardId; @Override public void initialize(InitializationInput initializationInput) { shardId = initializationInput.shardId(); MDC.put(SHARD_ID_MDC_KEY, shardId); try { log.info("Initializing @ Sequence: {}", initializationInput.extendedSequenceNumber()); } finally { MDC.remove(SHARD_ID_MDC_KEY); } } @Override public void processRecords(ProcessRecordsInput processRecordsInput) { MDC.put(SHARD_ID_MDC_KEY, shardId); try { log.info("Processing {} record(s)", processRecordsInput.records().size()); processRecordsInput.records().forEach(r -> log.info("Processing record pk: {} -- Seq: {}", r.partitionKey(), r.sequenceNumber()) ); // Checkpoint periodically processRecordsInput.checkpointer().checkpoint(); } catch (Throwable t) { log.error("Caught throwable while processing records. Aborting.", t); } finally { MDC.remove(SHARD_ID_MDC_KEY); } } @Override public void leaseLost(LeaseLostInput leaseLostInput) { MDC.put(SHARD_ID_MDC_KEY, shardId); try { log.info("Lost lease, so terminating."); } finally { MDC.remove(SHARD_ID_MDC_KEY); } } @Override public void shardEnded(ShardEndedInput shardEndedInput) { MDC.put(SHARD_ID_MDC_KEY, shardId); try { log.info("Reached shard end checkpointing."); shardEndedInput.checkpointer().checkpoint(); } catch (ShutdownException | InvalidStateException e) { log.error("Exception while checkpointing at shard end. Giving up.", e); } finally { MDC.remove(SHARD_ID_MDC_KEY); } } @Override public void shutdownRequested(ShutdownRequestedInput shutdownRequestedInput) { MDC.put(SHARD_ID_MDC_KEY, shardId); try { log.info("Scheduler is shutting down, checkpointing."); shutdownRequestedInput.checkpointer().checkpoint(); } catch (ShutdownException | InvalidStateException e) { log.error("Exception while checkpointing at requested shutdown. Giving up.", e); } finally { MDC.remove(SHARD_ID_MDC_KEY); } } }

Im Folgenden finden Sie eine detaillierte Erläuterung der einzelnen im Beispiel verwendeten Methoden:

initialisieren (InitializationInputInitializationInput)

  • Zweck: Richten Sie alle erforderlichen Ressourcen oder den Status für die Verarbeitung von Datensätzen ein.

  • Wann es aufgerufen wird: Einmal, wenn KCL diesem Datensatzprozessor einen Shard zuweist.

  • Die wichtigsten Punkte:

    • initializationInput.shardId(): Die ID des Shards, das dieser Prozessor verarbeiten wird.

    • initializationInput.extendedSequenceNumber(): Die Sequenznummer, von der aus die Verarbeitung gestartet werden soll.

processRecords (ProcessRecordsInputProzessRecordsInput)

  • Zweck: Verarbeitung der eingehenden Datensätze und optional Checkpoint-Fortschritt.

  • Wann es aufgerufen wird: Wiederholt, solange der Record Processor den Leasingvertrag für den Shard hält.

  • Die wichtigsten Punkte:

    • processRecordsInput.records(): Liste der zu verarbeitenden Datensätze.

    • processRecordsInput.checkpointer(): Wird verwendet, um den Fortschritt zu überprüfen.

    • Stellen Sie sicher, dass Sie während der Verarbeitung alle Ausnahmen behandelt haben, um zu verhindern, dass KCL fehlschlägt.

    • Diese Methode sollte idempotent sein, da derselbe Datensatz in einigen Szenarien mehrmals verarbeitet werden kann, z. B. Daten, die vor einem unerwarteten Absturz oder Neustart des Workers nicht überprüft wurden.

    • Leeren Sie vor dem Checkpoint stets alle gepufferten Daten, um die Datenkonsistenz sicherzustellen.

Leasing Verloren (Lease) LeaseLostInput LostInput

  • Zweck: Bereinigen Sie alle Ressourcen, die für die Verarbeitung dieses Shards spezifisch sind.

  • Wenn es aufgerufen wird: Wenn ein anderer Scheduler das Leasing für diesen Shard übernimmt.

  • Die wichtigsten Punkte:

    • Checkpointing ist bei dieser Methode nicht zulässig.

ShardEnded (geteilt) ShardEndedInput EndedInput

  • Zweck: Beenden Sie die Verarbeitung für diesen Shard und Checkpoint.

  • Wann es aufgerufen wird: Wenn der Shard geteilt oder zusammengeführt wird, was bedeutet, dass alle Daten für diesen Shard verarbeitet wurden.

  • Die wichtigsten Punkte:

    • shardEndedInput.checkpointer(): Wird verwendet, um das letzte Checkpointing durchzuführen.

    • Bei dieser Methode ist Checkpointing erforderlich, um die Verarbeitung abzuschließen.

    • Wenn die Daten und der Checkpoint hier nicht geleert werden, kann dies zu Datenverlust oder doppelter Verarbeitung führen, wenn der Shard erneut geöffnet wird.

ShutdownRequested (Herunterfahren) ShutdownRequestedInput RequestedInput

  • Zweck: Überprüfen und bereinigen Sie Ressourcen, wenn KCL heruntergefahren wird.

  • Wenn es aufgerufen wird: Wenn KCL heruntergefahren wird, z. B. wenn die Anwendung beendet wird).

  • Die wichtigsten Punkte:

    • shutdownRequestedInput.checkpointer(): Wird verwendet, um Checkpoints vor dem Herunterfahren durchzuführen.

    • Stellen Sie sicher, dass Sie Checkpointing in der Methode implementiert haben, damit der Fortschritt gespeichert wird, bevor die Anwendung beendet wird.

    • Wenn Daten und Checkpoint hier nicht geleert werden, kann es beim Neustart der Anwendung zu Datenverlust oder zur erneuten Verarbeitung von Datensätzen kommen.

Wichtig

KCL 3.x sorgt dafür, dass weniger Daten erneut verarbeitet werden müssen, wenn der Leasingvertrag von einem Mitarbeiter an einen anderen Mitarbeiter übergeben wird. Dies erfolgt durch Checkpoints, bevor der vorherige Mitarbeiter heruntergefahren wird. Wenn Sie die Checkpointing-Logik nicht in der shutdownRequested() Methode implementieren, werden Sie diesen Vorteil nicht sehen. Stellen Sie sicher, dass Sie eine Checkpointing-Logik in der shutdownRequested() Methode implementiert haben.

RecordProcessorFactory

RecordProcessorFactory ist verantwortlich für die Erstellung neuer RecordProcessor Instanzen. KCL verwendet diese Factory, um RecordProcessor für jeden Shard, den die Anwendung verarbeiten muss, einen neuen zu erstellen.

Wichtigste Aufgaben:

  • Erstellen Sie bei Bedarf neue RecordProcessor Instanzen

  • Stellen Sie sicher, dass jede RecordProcessor korrekt initialisiert ist

Das Folgende ist ein Implementierungsbeispiel:

import software.amazon.kinesis.processor.ShardRecordProcessor; import software.amazon.kinesis.processor.ShardRecordProcessorFactory; public class SampleRecordProcessorFactory implements ShardRecordProcessorFactory { @Override public ShardRecordProcessor shardRecordProcessor() { return new SampleRecordProcessor(); } }

In diesem Beispiel erstellt die Factory bei SampleRecordProcessor jedem Aufruf von shard RecordProcessor () ein neues. Sie können dies um die erforderliche Initialisierungslogik erweitern.

Scheduler

Der Scheduler ist eine übergeordnete Komponente, die alle Aktivitäten der KCL-Anwendung koordiniert. Es ist für die gesamte Orchestrierung der Datenverarbeitung verantwortlich.

Wichtigste Aufgaben:

  • Managen Sie den Lebenszyklus von RecordProcessors

  • Kümmern Sie sich um das Leasingmanagement für Shards

  • Koordinieren Sie den Checkpoint

  • Verteilen Sie die Auslastung der Shard-Verarbeitung auf mehrere Worker Ihrer Anwendung

  • Sorgen Sie für reibungsloses Herunterfahren und Beenden der Anwendung

Der Scheduler wird in der Regel in der Hauptanwendung erstellt und gestartet. Das Implementierungsbeispiel von Scheduler finden Sie im folgenden Abschnitt, Main Consumer Application.

Hauptanwendung für Verbraucher

Die Hauptanwendung für Verbraucher verbindet alle Komponenten miteinander. Sie ist verantwortlich für die Einrichtung des KCL-Consumers, die Erstellung der erforderlichen Clients, die Konfiguration des Schedulers und die Verwaltung des Lebenszyklus der Anwendung.

Wichtigste Aufgaben:

  • Richten Sie AWS Service-Clients ein (Kinesis, DynamoDB,) CloudWatch

  • Konfigurieren Sie die KCL-Anwendung

  • Erstellen und starten Sie den Scheduler

  • Das Herunterfahren der Anwendung durchführen

Das Folgende ist ein Implementierungsbeispiel:

import software.amazon.awssdk.regions.Region; import software.amazon.awssdk.services.cloudwatch.CloudWatchAsyncClient; import software.amazon.awssdk.services.dynamodb.DynamoDbAsyncClient; import software.amazon.awssdk.services.kinesis.KinesisAsyncClient; import software.amazon.kinesis.common.ConfigsBuilder; import software.amazon.kinesis.common.KinesisClientUtil; import software.amazon.kinesis.coordinator.Scheduler; import java.util.UUID; public class SampleConsumer { private final String streamName; private final Region region; private final KinesisAsyncClient kinesisClient; public SampleConsumer(String streamName, Region region) { this.streamName = streamName; this.region = region; this.kinesisClient = KinesisClientUtil.createKinesisAsyncClient(KinesisAsyncClient.builder().region(this.region)); } public void run() { DynamoDbAsyncClient dynamoDbAsyncClient = DynamoDbAsyncClient.builder().region(region).build(); CloudWatchAsyncClient cloudWatchClient = CloudWatchAsyncClient.builder().region(region).build(); ConfigsBuilder configsBuilder = new ConfigsBuilder( streamName, streamName, kinesisClient, dynamoDbAsyncClient, cloudWatchClient, UUID.randomUUID().toString(), new SampleRecordProcessorFactory() ); Scheduler scheduler = new Scheduler( configsBuilder.checkpointConfig(), configsBuilder.coordinatorConfig(), configsBuilder.leaseManagementConfig(), configsBuilder.lifecycleConfig(), configsBuilder.metricsConfig(), configsBuilder.processorConfig(), configsBuilder.retrievalConfig() ); Thread schedulerThread = new Thread(scheduler); schedulerThread.setDaemon(true); schedulerThread.start(); } public static void main(String[] args) { String streamName = "your-stream-name"; // replace with your stream name Region region = Region.US_EAST_1; // replace with your region new SampleConsumer(streamName, region).run(); } }

KCL erstellt standardmäßig einen Enhanced Fan-out (EFO) Consumer mit dediziertem Durchsatz. Weitere Informationen zu Enhanced finden Sie Fan-out unter. Entwickeln Sie verbesserte Fan-Out-Verbraucher mit dediziertem Durchsatz Wenn Sie weniger als 2 Verbraucher haben oder keine Verzögerungen bei der Leseverteilung unter 200 ms benötigen, müssen Sie die folgende Konfiguration im Scheduler-Objekt festlegen, um Verbraucher mit gemeinsamem Durchsatz zu verwenden:

configsBuilder.retrievalConfig().retrievalSpecificConfig(new PollingConfig(streamName, kinesisClient))

Der folgende Code ist ein Beispiel für die Erstellung eines Scheduler-Objekts, das Verbraucher mit gemeinsamem Durchsatz verwendet:

Importiert:

import software.amazon.kinesis.retrieval.polling.PollingConfig;

Kode:

Scheduler scheduler = new Scheduler( configsBuilder.checkpointConfig(), configsBuilder.coordinatorConfig(), configsBuilder.leaseManagementConfig(), configsBuilder.lifecycleConfig(), configsBuilder.metricsConfig(), configsBuilder.processorConfig(), configsBuilder.retrievalConfig().retrievalSpecificConfig(new PollingConfig(streamName, kinesisClient)) );/