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+.
Note
KCL 3.5.x+ est pris en charge avec Amazon DynamoDB Streams Kinesis Adapter version 2.4.x+ sur le site Web.
Étapes de la migration
Rubriques
É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
-
Remplacez les interfaces
com.amazonaws.services.kinesis.clientlibrary.interfaces.v2.IRecordProcessoretcom.amazonaws.services.kinesis.clientlibrary.interfaces.v2.IShutdownNotificationAwareparcom.amazonaws.services.dynamodbv2.streamsadapter.processor.DynamoDBStreamsShardRecordProcessorcomme 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; -
Mettez à jour les instructions d’importation des méthodes
initializeetprocessRecords:// 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; -
Remplacez la méthode
shutdownRequestedpar les nouvelles méthodes suivantes :leaseLost,shardEndedetshutdownRequested.// @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.IRecordProcessorFactoryparsoftware.amazon.kinesis.processor.ShardRecordProcessorFactorycomme 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
-
Modifiez l’instruction
importde la classeWorkerpour les instructions d’importation pour les classesScheduleretConfigsBuilder.// import com.amazonaws.services.kinesis.clientlibrary.lib.worker.Worker; import software.amazon.kinesis.coordinator.Scheduler; import software.amazon.kinesis.common.ConfigsBuilder; -
Importez
StreamTrackeret remplacez l’importationStreamsWorkerFactoryparStreamsSchedulerFactory.import software.amazon.kinesis.processor.StreamTracker; // import software.amazon.dynamodb.streamsadapter.StreamsWorkerFactory; import software.amazon.dynamodb.streamsadapter.StreamsSchedulerFactory; -
Choisissez la position à partir de laquelle vous souhaitez démarrer l’application. Vous avez le choix entre
TRIM_HORIZONetLATEST.import software.amazon.kinesis.common.InitialPositionInStream; import software.amazon.kinesis.common.InitialPositionInStreamExtended; -
Créez une instance
StreamTracker.StreamTracker streamTracker = StreamsSchedulerFactory.createSingleStreamTracker( streamArn, InitialPositionInStreamExtended.newInitialPosition(InitialPositionInStream.TRIM_HORIZON) ); -
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); -
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()); -
Créez le
Schedulerà l’aide deConfigsBuilder, 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. ÉcritWORKER_METRIC_STATSetMigration3.0saisit 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
billingModeest définie surPROVISIONED. Cependant, avec la version 3.5.x+ de KCL, la valeur par défautbillingModeestPAY_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
idleTimeBetweenReadsInMillisest 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 deleaseAssignmentIntervalMillisaugmente 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éeidleTimeInMillisdans 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.