View a markdown version of this page

Migration de KCL 1.x vers KCL 3.5.x+ - Amazon DynamoDB

Les traductions sont fournies par des outils de traduction automatique. En cas de conflit entre le contenu d'une traduction et celui de la version originale en anglais, la version anglaise prévaudra.

Migration de KCL 1.x vers KCL 3.5.x+

Présentation de

Ce guide fournit des instructions pour migrer votre application client de KCL 1.x vers KCL 3.5.x+. En raison des différences architecturales entre KCL 1.x et KCL 3.5.x+, la migration nécessite la mise à jour de plusieurs composants pour garantir la compatibilité.

KCL 1.x utilise différentes classes et interfaces par rapport à KCL 3.5.x+. Vous devez d'abord migrer le processeur d'enregistrement, la fabrique des processeurs d'enregistrement et les classes de travail vers le format compatible KCL 3.5.x+, puis suivre les étapes de migration pour la migration de KCL 1.x vers KCL 3.5.x+.

Étapes de la migration

Étape 1 : migrer le processeur d’enregistrements

L’exemple suivant illustre un processeur d’enregistrements implémenté pour l’adaptateur DynamoDB Streams Kinesis version 1.x :

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(); } } }
Pour migrer la RecordProcessor classe
  1. Remplacez les interfaces com.amazonaws.services.kinesis.clientlibrary.interfaces.v2.IRecordProcessor et com.amazonaws.services.kinesis.clientlibrary.interfaces.v2.IShutdownNotificationAware par com.amazonaws.services.dynamodbv2.streamsadapter.processor.DynamoDBStreamsShardRecordProcessor comme suit :

    // 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. Mettez à jour les instructions d’importation des méthodes initialize et 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. Remplacez la méthode shutdownRequested par les nouvelles méthodes suivantes : leaseLost, shardEnded et 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(); } }

Voici la version mise à jour de la classe du processeur d’enregistrements :

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(); } } }
Note

L’adaptateur DynamoDB Streams Kinesis utilise désormais le modèle d’enregistrement SDKv2. Dans SDKv2, les objets AttributeValue complexes (BS, NS, M, L et SS) ne renvoient jamais null. Vérifiez si ces valeurs existent à l’aide des méthodes hasBs(), hasNs(), hasM(), hasL() et hasSs().

Étape 2 : migrer la fabrique de processeurs d’enregistrements

La fabrique de processeurs d’enregistrements est responsable de la création des processeurs d’enregistrements lorsqu’un bail est acquis. Voici un exemple de fabrique KCL 1.x :

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); } }
Pour migrer la RecordProcessorFactory
  • Remplacez l’interface implémentée com.amazonaws.services.kinesis.clientlibrary.interfaces.v2.IRecordProcessorFactory par software.amazon.kinesis.processor.ShardRecordProcessorFactory comme suit :

    // 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() {

Voici un exemple de l'usine de processeurs d'enregistrement dans la version 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(); } }

Étape 3 : migrer le worker

Dans la version 3.5.x+ de la KCL, une nouvelle classe, appelée Scheduler, remplace la classe Worker. Voici un exemple de worker KCL 1.x :

final KinesisClientLibConfiguration config = new KinesisClientLibConfiguration(...) final IRecordProcessorFactory recordProcessorFactory = new RecordProcessorFactory(); final Worker worker = StreamsWorkerFactory.createDynamoDbStreamsWorker( recordProcessorFactory, workerConfig, adapterClient, amazonDynamoDB, amazonCloudWatchClient);
Pour migrer le worker
  1. Modifiez l’instruction import de la classe Worker pour les instructions d’importation pour les classes Scheduler et ConfigsBuilder.

    // import com.amazonaws.services.kinesis.clientlibrary.lib.worker.Worker; import software.amazon.kinesis.coordinator.Scheduler; import software.amazon.kinesis.common.ConfigsBuilder;
  2. Importez StreamTracker et remplacez l’importation StreamsWorkerFactory par StreamsSchedulerFactory.

    import software.amazon.kinesis.processor.StreamTracker; // import software.amazon.dynamodb.streamsadapter.StreamsWorkerFactory; import software.amazon.dynamodb.streamsadapter.StreamsSchedulerFactory;
  3. Choisissez la position à partir de laquelle vous souhaitez démarrer l’application. Vous avez le choix entre TRIM_HORIZON et LATEST.

    import software.amazon.kinesis.common.InitialPositionInStream; import software.amazon.kinesis.common.InitialPositionInStreamExtended;
  4. Créez une instance StreamTracker.

    StreamTracker streamTracker = StreamsSchedulerFactory.createSingleStreamTracker( streamArn, InitialPositionInStreamExtended.newInitialPosition(InitialPositionInStream.TRIM_HORIZON) );
  5. Créez l’objet 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. Créez l’objet ConfigsBuilder.

    import software.amazon.kinesis.common.ConfigsBuilder; ... ConfigsBuilder configsBuilder = new ConfigsBuilder( streamTracker, applicationName, adapterClient, dynamoDbAsyncClient, cloudWatchAsyncClient, UUID.randomUUID().toString(), new StreamsRecordProcessorFactory());
  7. Créez le Scheduler à l’aide de ConfigsBuilder, comme illustré dans l’exemple suivant :

    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 );
Note

La migration vers KCL 3.5.x+ comporte trois phases :

  • Phase 1 (CLIENT_VERSION_CONFIG_COMPATIBLE_WITH_2X_PHASE1) : mode compatible avec Pure KCL 1.x. Aucune métadonnée spécifique à la migration n'est écrite dans la table des baux. Revenez en toute sécurité à KCL v1 en redéployant le code précédent. Utilisez cette phase pour valider la stabilité.

  • Phase 2 (CLIENT_VERSION_CONFIG_COMPATIBLE_WITH_2X) : lance la migration. Écrit WORKER_METRIC_STATS et Migration3.0 saisit dans le tableau des baux. KCL passe automatiquement à l'équilibrage de charge 3.x complet lorsque tous les collaborateurs sont prêts. Le retour à la phase 1 est pris en charge (via l'outil de migration KCL). Le retour à KCL v1 n'est plus possible.

  • Phase 3 (CLIENT_VERSION_CONFIG_3X) : fonctionnalité complète de KCL 3.x. Défini explicitement par le client ou utilisé par défaut lorsque la configuration est supprimée. État du terminal, pas de retour en arrière.

Ces paramètres garantissent la compatibilité entre l'adaptateur DynamoDB Streams Kinesis pour KCL v3 et KCL v1, et non entre KCL v2 et v3.

Important

Vous devez commencer la migration par CLIENT_VERSION_CONFIG_COMPATIBLE_WITH_2X_PHASE1 (Phase 1). La phase 1 est rétrocompatible avec KCL v1 et n'écrit aucune entrée spécifique à la migration dans la table des baux, ce qui permet de revenir en toute sécurité à votre version précédente de KCL en redéployant simplement votre code précédent. Après des tests de cuisson approfondis au cours de la phase 1, vous pouvez passer à la phase 2 (CLIENT_VERSION_CONFIG_COMPATIBLE_WITH_2X) pour démarrer la migration complète. Si vous sautez la phase 1 et que vous commencez directement par la phase 2, les entrées non liées au bail sont immédiatement écrites dans la table des baux, ce qui empêche définitivement le retour à KCL v1 sans nettoyage manuel de DynamoDB.

Étape 4 : présentation de la configuration de KCL 3.5.x+ et recommandations

Pour une description détaillée des configurations introduites après KCL 1.x qui sont pertinentes dans KCL 3.5.x+, voir Configurations KCL et configuration du client de migration KCL. https://docs.aws.amazon.com//streams/latest/dev/kcl-migration.html#client-configuration

Important

Au lieu de créer directement des objets decheckpointConfig,,, processorConfig et coordinatorConfig leaseManagementConfig metricsConfigretrievalConfig, nous vous recommandons de les utiliser ConfigsBuilder pour définir les configurations dans KCL 3.5.x+ et les versions ultérieures afin d'éviter les problèmes d'initialisation du planificateur. ConfigsBuilderfournit un moyen plus flexible et plus facile à gérer de configurer votre application KCL.

Configurations avec mise à jour de la valeur par défaut dans KCL 3.5.x+

billingMode

Dans la KCL version 1.x, la valeur par défaut de billingMode est définie sur PROVISIONED. Cependant, avec la version 3.5.x+ de KCL, la valeur par défaut billingMode est PAY_PER_REQUEST (mode à la demande). Nous vous recommandons d’utiliser le mode de capacité à la demande pour votre table de baux, afin d’ajuster automatiquement la capacité en fonction de votre utilisation. Pour obtenir des conseils sur l’utilisation de la capacité allouée pour vos tables de baux, consultez Best practices for the lease table with provisioned capacity mode.

idleTimeBetweenReadsInMillis

Dans la KCL version 1.x, la valeur par défaut de idleTimeBetweenReadsInMillis est définie sur 1 000 (soit 1 seconde). La version 3.5.x+ de KCL définit la valeur par défaut sur 1 500 (soit 1,5 seconde), mais Amazon DynamoDB Streams Kinesis Adapter remplace la valeur par défaut sur 1 000 (ou 1 seconde). idleTimeBetweenReadsInMillis

Nouvelles configurations dans KCL 3.5.x+

leaseAssignmentIntervalMillis

Cette configuration définit l’intervalle de temps avant que les partitions récemment découvertes ne commencent à être traitées. Elle est calculée comme suit : 1,5 × leaseAssignmentIntervalMillis. Si ce paramètre n’est pas explicitement configuré, l’intervalle de temps est défini par défaut sur 1,5 × failoverTimeMillis. Le traitement des nouvelles partitions consiste à analyser la table de baux et à interroger un index secondaire global (GSI) de la table de baux. La réduction de leaseAssignmentIntervalMillis augmente la fréquence de ces opérations d’analyse et d’interrogation, ce qui entraîne une augmentation des coûts de DynamoDB. Nous vous recommandons de définir cette valeur sur 2 000 (soit 2 secondes) afin de réduire le délai de traitement des nouvelles partitions.

shardConsumerDispatchPollIntervalMillis

Cette configuration définit l’intervalle entre les interrogations successives effectuées par le consommateur de partitions pour déclencher des transitions d’état. Dans la KCL version 1.x, ce comportement était contrôlé par le paramètre idleTimeInMillis, qui n’était pas exposé en tant que paramètre configurable. Avec la version 3.5.x+ de KCL, nous vous recommandons de définir cette configuration pour qu'elle corresponde à la valeur utilisée idleTimeInMillis dans votre configuration de KCL version 1.x.

Étape 5 : Migrer de KCL 2.x vers KCL 3.5.x+

Pour garantir une transition fluide et une compatibilité avec la dernière version de Kinesis Client Library (KCL), suivez les étapes 5 à 8 des instructions du guide de migration pour passer de KCL 2.x à KCL 3.5.x+.

Pour les problèmes de dépannage courants liés à KCL 3.5.x+, voir Résolution des problèmes liés aux applications grand public KCL.