View a markdown version of this page

Migration von KCL 1.x zu KCL 3.5.x+ - Amazon DynamoDB

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.

Migration von KCL 1.x zu KCL 3.5.x+

-Übersicht

Dieses Handbuch enthält Anweisungen zur Migration Ihrer Consumer-Anwendung von KCL 1.x auf KCL 3.5.x+. Aufgrund der architektonischen Unterschiede zwischen KCL 1.x und KCL 3.5.x+ müssen bei der Migration mehrere Komponenten aktualisiert werden, um die Kompatibilität sicherzustellen.

KCL 1.x verwendet im Vergleich zu KCL 3.5.x+ andere Klassen und Schnittstellen. Sie müssen zuerst die Klassen Record Processor, Record Processor Factory und Worker in das mit KCL 3.5.x+ kompatible Format migrieren und die Migrationsschritte für die Migration von KCL 1.x zu KCL 3.5.x+ befolgen.

Anmerkung

KCL 3.5.x+ wird mit der Amazon DynamoDB Streams Kinesis Adapter-Version 2.4.x+ auf der Website unterstützt. https://github.com/awslabs/dynamodb-streams-kinesis-adapter GitHub

Schritte zur Migration

Schritt 1: Migrieren des Datensatzprozessors

Das folgende Beispiel zeigt einen Datensatzprozessor, der für die Version KCL 1.x des DynamoDB-Streams-Kinesis-Adapters implementiert wurde:

package com.amazonaws.kcl; import com.amazonaws.services.kinesis.clientlibrary.exceptions.InvalidStateException; import com.amazonaws.services.kinesis.clientlibrary.exceptions.ShutdownException; import com.amazonaws.services.kinesis.clientlibrary.interfaces.IRecordProcessorCheckpointer; import com.amazonaws.services.kinesis.clientlibrary.interfaces.v2.IRecordProcessor; import com.amazonaws.services.kinesis.clientlibrary.interfaces.v2.IShutdownNotificationAware; import com.amazonaws.services.kinesis.clientlibrary.lib.worker.ShutdownReason; import com.amazonaws.services.kinesis.clientlibrary.types.InitializationInput; import com.amazonaws.services.kinesis.clientlibrary.types.ProcessRecordsInput; import com.amazonaws.services.kinesis.clientlibrary.types.ShutdownInput; public class StreamsRecordProcessor implements IRecordProcessor, IShutdownNotificationAware { @Override public void initialize(InitializationInput initializationInput) { // // Setup record processor // } @Override public void processRecords(ProcessRecordsInput processRecordsInput) { for (Record record : processRecordsInput.getRecords()) { String data = new String(record.getData().array(), Charset.forName("UTF-8")); System.out.println(data); if (record instanceof RecordAdapter) { // record processing and checkpointing logic } } } @Override public void shutdown(ShutdownInput shutdownInput) { if (shutdownInput.getShutdownReason() == ShutdownReason.TERMINATE) { try { shutdownInput.getCheckpointer().checkpoint(); } catch (ShutdownException | InvalidStateException e) { throw new RuntimeException(e); } } } @Override public void shutdownRequested(IRecordProcessorCheckpointer checkpointer) { try { checkpointer.checkpoint(); } catch (ShutdownException | InvalidStateException e) { // // Swallow exception // e.printStackTrace(); } } }
Um RecordProcessor die Klasse zu migrieren
  1. Ändern Sie die Schnittstellen von com.amazonaws.services.kinesis.clientlibrary.interfaces.v2.IRecordProcessor und com.amazonaws.services.kinesis.clientlibrary.interfaces.v2.IShutdownNotificationAware folgendermaßen zu com.amazonaws.services.dynamodbv2.streamsadapter.processor.DynamoDBStreamsShardRecordProcessor:

    // import com.amazonaws.services.kinesis.clientlibrary.interfaces.v2.IRecordProcessor; // import com.amazonaws.services.kinesis.clientlibrary.interfaces.v2.IShutdownNotificationAware; import com.amazonaws.services.dynamodbv2.streamsadapter.processor.DynamoDBStreamsShardRecordProcessor;
  2. Aktualisieren Sie die Importanweisungen für die Methoden initialize und processRecords:

    // import com.amazonaws.services.kinesis.clientlibrary.types.InitializationInput; import software.amazon.kinesis.lifecycle.events.InitializationInput; // import com.amazonaws.services.kinesis.clientlibrary.types.ProcessRecordsInput; import com.amazonaws.services.dynamodbv2.streamsadapter.model.DynamoDBStreamsProcessRecordsInput;
  3. Ersetzen Sie die Methode shutdownRequested durch die folgenden neuen Methoden: leaseLost, shardEnded und shutdownRequested.

    // @Override // public void shutdownRequested(IRecordProcessorCheckpointer checkpointer) { // // // // This is moved to shardEnded(...) and shutdownRequested(ShutdownReauestedInput) // // // try { // checkpointer.checkpoint(); // } catch (ShutdownException | InvalidStateException e) { // // // // Swallow exception // // // e.printStackTrace(); // } // } @Override public void leaseLost(LeaseLostInput leaseLostInput) { } @Override public void shardEnded(ShardEndedInput shardEndedInput) { try { shardEndedInput.checkpointer().checkpoint(); } catch (ShutdownException | InvalidStateException e) { // // Swallow the exception // e.printStackTrace(); } } @Override public void shutdownRequested(ShutdownRequestedInput shutdownRequestedInput) { try { shutdownRequestedInput.checkpointer().checkpoint(); } catch (ShutdownException | InvalidStateException e) { // // Swallow the exception // e.printStackTrace(); } }

Nachstehend finden Sie die aktualisierte Version der Datensatzprozessorklasse:

package com.amazonaws.codesamples; import software.amazon.kinesis.exceptions.InvalidStateException; import software.amazon.kinesis.exceptions.ShutdownException; import software.amazon.kinesis.lifecycle.events.InitializationInput; import software.amazon.kinesis.lifecycle.events.LeaseLostInput; import com.amazonaws.services.dynamodbv2.streamsadapter.model.DynamoDBStreamsProcessRecordsInput; import software.amazon.kinesis.lifecycle.events.ShardEndedInput; import software.amazon.kinesis.lifecycle.events.ShutdownRequestedInput; import software.amazon.dynamodb.streamsadapter.processor.DynamoDBStreamsShardRecordProcessor; import software.amazon.dynamodb.streamsadapter.adapter.DynamoDBStreamsKinesisClientRecord; import com.amazonaws.services.dynamodbv2.streamsadapter.processor.DynamoDBStreamsShardRecordProcessor; import com.amazonaws.services.dynamodbv2.streamsadapter.adapter.DynamoDBStreamsClientRecord; import software.amazon.awssdk.services.dynamodb.model.Record; public class StreamsRecordProcessor implements DynamoDBStreamsShardRecordProcessor { @Override public void initialize(InitializationInput initializationInput) { } @Override public void processRecords(DynamoDBStreamsProcessRecordsInput processRecordsInput) { for (DynamoDBStreamsKinesisClientRecord record: processRecordsInput.records()) Record ddbRecord = record.getRecord(); // processing and checkpointing logic for the ddbRecord } } @Override public void leaseLost(LeaseLostInput leaseLostInput) { } @Override public void shardEnded(ShardEndedInput shardEndedInput) { try { shardEndedInput.checkpointer().checkpoint(); } catch (ShutdownException | InvalidStateException e) { // // Swallow the exception // e.printStackTrace(); } } @Override public void shutdownRequested(ShutdownRequestedInput shutdownRequestedInput) { try { shutdownRequestedInput.checkpointer().checkpoint(); } catch (ShutdownException | InvalidStateException e) { // // Swallow the exception // e.printStackTrace(); } } }
Anmerkung

Der DynamoDB-Streams-Kinesis-Adapter verwendet jetzt das SDKv2-Datensatzmodell. In SDKv2 geben komplexe AttributeValue-Objekte (BS, NS, M, L, SS) niemals Null zurück. Verwenden Sie die Methoden hasBs(), hasNs(), hasM(), hasL(), hasSs(), um zu überprüfen, ob diese Werte existieren.

Schritt 2: Migrieren der Datensatzprozessor-Factory

Die Datensatzprozessor-Factory ist für das Erstellen von Prozessoren verantwortlich, wenn eine Lease erworben wird. Nachfolgend sehen Sie ein Beispiel für eine KCL-1.x-Factory:

package com.amazonaws.codesamples; import software.amazon.dynamodb.AmazonDynamoDB; import com.amazonaws.services.kinesis.clientlibrary.interfaces.v2.IRecordProcessor; import com.amazonaws.services.kinesis.clientlibrary.interfaces.v2.IRecordProcessorFactory; public class StreamsRecordProcessorFactory implements IRecordProcessorFactory { @Override public IRecordProcessor createProcessor() { return new StreamsRecordProcessor(dynamoDBClient, tableName); } }
So migrieren Sie die RecordProcessorFactory
  • Ändern Sie die implementierte Schnittstelle von com.amazonaws.services.kinesis.clientlibrary.interfaces.v2.IRecordProcessorFactory folgendermaßen zu software.amazon.kinesis.processor.ShardRecordProcessorFactory:

    // import com.amazonaws.services.kinesis.clientlibrary.interfaces.v2.IRecordProcessor; import software.amazon.kinesis.processor.ShardRecordProcessor; // import com.amazonaws.services.kinesis.clientlibrary.interfaces.v2.IRecordProcessorFactory; import software.amazon.kinesis.processor.ShardRecordProcessorFactory; // public class TestRecordProcessorFactory implements IRecordProcessorFactory { public class StreamsRecordProcessorFactory implements ShardRecordProcessorFactory { Change the return signature for createProcessor. // public IRecordProcessor createProcessor() { public ShardRecordProcessor shardRecordProcessor() {

Das Folgende ist ein Beispiel für die Record Processor Factory in 3.5.x+:

package com.amazonaws.codesamples; import software.amazon.kinesis.processor.ShardRecordProcessor; import software.amazon.kinesis.processor.ShardRecordProcessorFactory; public class StreamsRecordProcessorFactory implements ShardRecordProcessorFactory { @Override public ShardRecordProcessor shardRecordProcessor() { return new StreamsRecordProcessor(); } }

Schritt 3: Migrieren des Workers

In Version 3.5.x+ der KCL ersetzt eine neue Klasse namens Scheduler die Worker-Klasse. Nachfolgend sehen Sie ein Beispiel für einen KCL-1.x-Worker.

final KinesisClientLibConfiguration config = new KinesisClientLibConfiguration(...) final IRecordProcessorFactory recordProcessorFactory = new RecordProcessorFactory(); final Worker worker = StreamsWorkerFactory.createDynamoDbStreamsWorker( recordProcessorFactory, workerConfig, adapterClient, amazonDynamoDB, amazonCloudWatchClient);
So migrieren Sie den Worker
  1. Ändern Sie die import-Anweisung für die Worker-Klasse in die Import-Anweisungen für die Klassen Scheduler und ConfigsBuilder.

    // import com.amazonaws.services.kinesis.clientlibrary.lib.worker.Worker; import software.amazon.kinesis.coordinator.Scheduler; import software.amazon.kinesis.common.ConfigsBuilder;
  2. Importieren Sie StreamTracker und ändern Sie den Import von StreamsWorkerFactory zu StreamsSchedulerFactory.

    import software.amazon.kinesis.processor.StreamTracker; // import software.amazon.dynamodb.streamsadapter.StreamsWorkerFactory; import software.amazon.dynamodb.streamsadapter.StreamsSchedulerFactory;
  3. Wählen Sie die Position, von der aus die Anwendung gestartet werden soll. Möglich sind TRIM_HORIZON oder LATEST.

    import software.amazon.kinesis.common.InitialPositionInStream; import software.amazon.kinesis.common.InitialPositionInStreamExtended;
  4. Erstellen Sie eine StreamTracker-Instance.

    StreamTracker streamTracker = StreamsSchedulerFactory.createSingleStreamTracker( streamArn, InitialPositionInStreamExtended.newInitialPosition(InitialPositionInStream.TRIM_HORIZON) );
  5. Erstellen Sie das Objekt AmazonDynamoDBStreamsAdapterClient.

    import software.amazon.dynamodb.streamsadapter.AmazonDynamoDBStreamsAdapterClient; import software.amazon.awssdk.auth.credentials.AwsCredentialsProvider; import software.amazon.awssdk.auth.credentials.DefaultCredentialsProvider; ... AwsCredentialsProvider credentialsProvider = DefaultCredentialsProvider.create(); AmazonDynamoDBStreamsAdapterClient adapterClient = new AmazonDynamoDBStreamsAdapterClient( credentialsProvider, awsRegion);
  6. Erstellen Sie das Objekt ConfigsBuilder.

    import software.amazon.kinesis.common.ConfigsBuilder; ... ConfigsBuilder configsBuilder = new ConfigsBuilder( streamTracker, applicationName, adapterClient, dynamoDbAsyncClient, cloudWatchAsyncClient, UUID.randomUUID().toString(), new StreamsRecordProcessorFactory());
  7. Erstellen Sie den Scheduler mithilfe von ConfigsBuilder, wie im folgenden Beispiel gezeigt:

    import java.util.UUID; import software.amazon.awssdk.regions.Region; import software.amazon.awssdk.services.dynamodb.DynamoDbAsyncClient; import software.amazon.awssdk.services.cloudwatch.CloudWatchAsyncClient; import software.amazon.awssdk.services.kinesis.KinesisAsyncClient; import software.amazon.kinesis.common.KinesisClientUtil; import software.amazon.kinesis.coordinator.Scheduler; ... DynamoDbAsyncClient dynamoClient = DynamoDbAsyncClient.builder().region(region).build(); CloudWatchAsyncClient cloudWatchClient = CloudWatchAsyncClient.builder().region(region).build(); DynamoDBStreamsPollingConfig pollingConfig = new DynamoDBStreamsPollingConfig(adapterClient); pollingConfig.idleTimeBetweenReadsInMillis(idleTimeBetweenReadsInMillis); // Use ConfigsBuilder to configure settings RetrievalConfig retrievalConfig = configsBuilder.retrievalConfig(); retrievalConfig.retrievalSpecificConfig(pollingConfig); CoordinatorConfig coordinatorConfig = configsBuilder.coordinatorConfig(); coordinatorConfig.clientVersionConfig(CoordinatorConfig.ClientVersionConfig.CLIENT_VERSION_CONFIG_COMPATIBLE_WITH_2X_PHASE1); Scheduler scheduler = StreamsSchedulerFactory.createScheduler( configsBuilder.checkpointConfig(), coordinatorConfig, configsBuilder.leaseManagementConfig(), configsBuilder.lifecycleConfig(), configsBuilder.metricsConfig(), configsBuilder.processorConfig(), retrievalConfig, adapterClient );
Anmerkung

Die KCL 3.5.x+ Migration erfolgt in drei Phasen:

  • Phase 1 (CLIENT_VERSION_CONFIG_COMPATIBLE_WITH_2X_PHASE1): Reiner KCL 1.x-kompatibler Modus. Es werden keine migrationsspezifischen Metadaten in die Leasetabelle geschrieben. Sicheres Rollback zu KCL v1, indem Sie vorherigen Code erneut bereitstellen. Verwenden Sie diese Phase, um die Stabilität zu überprüfen.

  • Phase 2 (CLIENT_VERSION_CONFIG_COMPATIBLE_WITH_2X): Startet die Migration. Schreibt WORKER_METRIC_STATS und schreibt Migration3.0 Einträge in die Leasing-Tabelle. KCL wechselt automatisch zum vollständigen 3.x-Load Balancing, wenn alle Worker bereit sind. Ein Rollback zu Phase 1 wird unterstützt (über das KCL-Migrationstool). Ein Rollback auf KCL v1 ist nicht mehr möglich.

  • Phase 3 (CLIENT_VERSION_CONFIG_3X): Volle KCL 3.x-Funktionalität. Explizit vom Kunden festgelegt oder als Standard verwendet, wenn die Konfiguration entfernt wird. Terminalstatus, kein Rollback.

Diese Einstellungen gewährleisten die Kompatibilität zwischen dem DynamoDB Streams Kinesis Adapter für KCL v3 und KCL v1, nicht zwischen KCL v2 und v3.

Wichtig

Sie müssen die Migration mit (Phase 1) starten. CLIENT_VERSION_CONFIG_COMPATIBLE_WITH_2X_PHASE1 Phase 1 ist abwärtskompatibel mit KCL v1 und schreibt keine migrationsspezifischen Einträge in die Leasetabelle, sodass ein sicheres Rollback zu Ihrer vorherigen KCL-Version möglich ist, indem Sie einfach Ihren vorherigen Code erneut bereitstellen. Nach gründlichen Bake-Tests in Phase 1 können Sie mit Phase 2 () fortfahren, um die vollständige Migration zu starten. CLIENT_VERSION_CONFIG_COMPATIBLE_WITH_2X Wenn Sie Phase 1 überspringen und direkt mit Phase 2 beginnen, werden Nicht-Lease-Einträge sofort in die Leasing-Tabelle geschrieben, sodass ein Rollback zu KCL v1 ohne manuelle DynamoDB-Bereinigung dauerhaft verhindert wird.

Schritt 4: Überblick und Empfehlungen zur Konfiguration von KCL 3.5.x+

Eine ausführliche Beschreibung der nach KCL 1.x eingeführten Konfigurationen, die in KCL 3.5.x+ relevant sind, finden Sie unter KCL-Konfigurationen und Konfiguration des KCL-Migrationsclients. https://docs.aws.amazon.com//streams/latest/dev/kcl-configuration.html https://docs.aws.amazon.com//streams/latest/dev/kcl-migration.html#client-configuration

Wichtig

Anstatt Objekte aus,, und direkt zu erstellen, empfehlen wir checkpointConfig coordinatorConfig leaseManagementConfigmetricsConfig, die Konfiguration in KCL processorConfig retrievalConfig 3.5.x+ und späteren Versionen ConfigsBuilder zu verwenden, um Probleme mit der Scheduler-Initialisierung zu vermeiden. ConfigsBuilderbietet eine flexiblere und wartbarere Methode zur Konfiguration Ihrer KCL-Anwendung.

Konfigurationen mit aktualisiertem Standardwert in KCL 3.5.x+

billingMode

In der KCL-Version 1.x ist der Standardwert für billingMode auf PROVISIONED eingestellt. Bei KCL-Version 3.5.x+ ist die Standardeinstellung jedoch (On-Demand-Modus). billingMode PAY_PER_REQUEST Wir empfehlen Ihnen, den On-Demand-Kapazitätsmodus für Ihre Leasetabelle zu verwenden, um die Kapazität automatisch an die Nutzung anzupassen. Anleitungen zur Verwendung der bereitgestellten Kapazität für Ihre Leasetabellen finden Sie unter Best practices for the lease table with provisioned capacity mode.

idleTimeBetweenReadsInMillis

In der KCL-Version 1.x ist der Standardwert für idleTimeBetweenReadsInMillis auf 1 000 (oder 1 Sekunde) eingestellt. KCL-Version 3.5.x+ legt den Standardwert idleTimeBetweenReadsInMillis auf 1.500 (oder 1,5 Sekunden) fest, aber Amazon DynamoDB Streams Kinesis Adapter überschreibt den Standardwert auf 1.000 (oder 1 Sekunde).

Neue Konfigurationen in KCL 3.5.x+

leaseAssignmentIntervalMillis

Diese Konfiguration definiert das Zeitintervall, bis die Verarbeitung neu erkannter Shards beginnt. Es wird nach der Formel 1,5 × leaseAssignmentIntervalMillis berechnet. Wenn diese Einstellung nicht explizit konfiguriert ist, beträgt das Zeitintervall standardmäßig 1,5 × failoverTimeMillis. Die Verarbeitung neuer Shards beinhaltet das Scannen der Leasetabelle und das Abfragen eines globalen sekundären Index (GSI) in der Leasetabelle. Eine Absenkung des leaseAssignmentIntervalMillis-Werts erhöht die Häufigkeit dieser Scan- und Abfragevorgänge, was zu höheren DynamoDB-Kosten führt. Wir empfehlen, diesen Wert auf 2 000 (oder 2 Sekunden) einzustellen, um die Verzögerung bei der Verarbeitung neuer Shards zu minimieren.

shardConsumerDispatchPollIntervalMillis

Diese Konfiguration definiert das Intervall zwischen aufeinanderfolgenden Abfragen durch den Shard-Verbraucher, um Zustandsübergänge auszulösen. In KCL-Version 1.x wurde dieses Verhalten durch den Parameter idleTimeInMillis gesteuert, der nicht als konfigurierbare Einstellung verfügbar war. Bei KCL-Version 3.5.x+ empfehlen wir, diese Konfiguration so einzustellen, dass sie dem Wert entspricht, der in Ihrer KCL-Version 1.x-Setup verwendet wurde. idleTimeInMillis

Schritt 5: Migrieren Sie von KCL 2.x zu KCL 3.5.x+

Um einen reibungslosen Übergang und die Kompatibilität mit der neuesten Version der Kinesis Client Library (KCL) zu gewährleisten, folgen Sie den Schritten 5-8 in den Anweisungen für das Upgrade von KCL 2.x auf KCL 3.5.x+ im Migrationshandbuch. https://docs.aws.amazon.com//streams/latest/dev/kcl-migration-from-2-3.html#kcl-migration-from-2-3-worker-metrics

Informationen zu häufig auftretenden Problemen bei der Behebung von KCL 3.5.x+ finden Sie unter Problembehandlung bei KCL-Verbraucheranwendungen. https://docs.aws.amazon.com//streams/latest/dev/troubleshooting-consumers.html