Migrar de KCL 1.x para KCL 3.5.x+
Visão geral
Este guia oferece instruções para migrar um aplicativo de consumidor de KCL 1.x para KCL 3.5.x+. Devido a diferenças de arquitetura entre KCL 1.x e KCL 3.5.x+, a migração requer a atualização de vários componentes para garantir compatibilidade.
KCL 1.x usa classes e interfaces diferentes em comparação com KCL 3.5.x+. É necessário primeiro migrar o processador de registros, a fábrica do processador de registros e as classes de operador para o formato compatível com KCL 3.5.x+ e depois seguir as etapas de migração de KCL 1.x para KCL 3.5.x+.
nota
KCL 3.5.x+ é compatível com o Amazon DynamoDB Streams Kinesis Adapter
Etapas da migração
Tópicos
Etapa 1: migrar o processador de registros
Este exemplo mostra um processador de registros implementado para o DynamoDB Streams Kinesis Adapter da KCL 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(); } } }
Como migrar a classe RecordProcessor
-
Altere as interfaces de
com.amazonaws.services.kinesis.clientlibrary.interfaces.v2.IRecordProcessorecom.amazonaws.services.kinesis.clientlibrary.interfaces.v2.IShutdownNotificationAwareparacom.amazonaws.services.dynamodbv2.streamsadapter.processor.DynamoDBStreamsShardRecordProcessorda seguinte forma:// 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; -
Atualize as instruções de importação para os métodos
initializeeprocessRecords.// 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; -
Substitua o método
shutdownRequestedpelos novos métodos a seguir:leaseLost,shardEnded, eshutdownRequested.// @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(); } }
Esta é a versão atualizada da classe de processador de registros:
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(); } } }
nota
O DynamoDB Streams Kinesis Adapter agora usa o modelo de registro SDKv2. No SDKv2, objetos AttributeValue complexos (BS, NS, M, L e SS) nunca retornam null. Use os métodos hasBs(), hasNs(), hasM(), hasL() e hasSs() para verificar se esses valores existem.
Etapa 2: migrar a fábrica do processador de registros
A fábrica do processador de registros é responsável por criar processadores de registro quando uma concessão é realizada. Veja o seguinte exemplo de fábrica da 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); } }
Migrar para RecordProcessorFactory
-
Altere a interface implementada de
com.amazonaws.services.kinesis.clientlibrary.interfaces.v2.IRecordProcessorFactoryparasoftware.amazon.kinesis.processor.ShardRecordProcessorFactoryda seguinte forma:// 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() {
Veja o seguinte exemplo de fábrica do processador de registros em 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(); } }
Etapa 3: migrar o operador
Na versão 3.5.x+ da KCL, uma nova classe, chamada Scheduler, substitui a classe Worker. Veja o seguinte exemplo de operador da KCL 1.x:
final KinesisClientLibConfiguration config = new KinesisClientLibConfiguration(...) final IRecordProcessorFactory recordProcessorFactory = new RecordProcessorFactory(); final Worker worker = StreamsWorkerFactory.createDynamoDbStreamsWorker( recordProcessorFactory, workerConfig, adapterClient, amazonDynamoDB, amazonCloudWatchClient);
Para migrar o operador
-
Altere a instrução
importpara a classeWorkerpara as instruções de importação para as classesSchedulereConfigsBuilder.// import com.amazonaws.services.kinesis.clientlibrary.lib.worker.Worker; import software.amazon.kinesis.coordinator.Scheduler; import software.amazon.kinesis.common.ConfigsBuilder; -
Importe
StreamTrackere altere a importação deStreamsWorkerFactoryparaStreamsSchedulerFactory.import software.amazon.kinesis.processor.StreamTracker; // import software.amazon.dynamodb.streamsadapter.StreamsWorkerFactory; import software.amazon.dynamodb.streamsadapter.StreamsSchedulerFactory; -
Escolha a posição por meio da qual iniciar a aplicação. Ela pode ser
TRIM_HORIZONouLATEST.import software.amazon.kinesis.common.InitialPositionInStream; import software.amazon.kinesis.common.InitialPositionInStreamExtended; -
Crie uma instância de
StreamTracker.StreamTracker streamTracker = StreamsSchedulerFactory.createSingleStreamTracker( streamArn, InitialPositionInStreamExtended.newInitialPosition(InitialPositionInStream.TRIM_HORIZON) ); -
Crie o objeto
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); -
Crie o objeto
ConfigsBuilder.import software.amazon.kinesis.common.ConfigsBuilder; ... ConfigsBuilder configsBuilder = new ConfigsBuilder( streamTracker, applicationName, adapterClient, dynamoDbAsyncClient, cloudWatchAsyncClient, UUID.randomUUID().toString(), new StreamsRecordProcessorFactory()); -
Crie o
SchedulerusandoConfigsBuildercomo mostrado no seguinte exemplo: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 );
nota
A migração da KCL 3.5.x+ tem três fases:
-
Fase 1 (
CLIENT_VERSION_CONFIG_COMPATIBLE_WITH_2X_PHASE1): modo compatível com KCL 1.x puro. Nenhum metadado específico da migração é gravado na tabela de concessão. Reversão segura para KCL v1 reimplantando o código anterior. Use essa fase para validar a estabilidade. -
Fase 2 (
CLIENT_VERSION_CONFIG_COMPATIBLE_WITH_2X): inicia a migração. Grava entradasWORKER_METRIC_STATSeMigration3.0na tabela de concessão. A KCL faz a transição automática para o balanceamento de carga 3.x completo quando todos os operadores estão prontos. A reversão para a Fase 1 é permitida (por meio da Ferramenta de Migração da KCL). A reversão para KCL v1 não é mais possível. -
Fase 3 (
CLIENT_VERSION_CONFIG_3X): funcionalidade completa da KCL 3.x. Definido explicitamente pelo cliente ou usado como padrão quando a configuração é removida. Estado de terminal, sem reversão.
As configurações mantêm a compatibilidade entre o DynamoDB Streams Kinesis Adapter para KCL v3 e KCL v1, não entre KCL v2 e v3.
Importante
Você deve iniciar a migração com CLIENT_VERSION_CONFIG_COMPATIBLE_WITH_2X_PHASE1 (Fase 1). A Fase 1 é compatível com versões anteriores da KCL v1 e não grava nenhuma entrada específica de migração na tabela de concessão, permitindo a reversão segura para a versão anterior da KCL simplesmente reimplantando seu código anterior. Após um teste completo na Fase 1, você pode prosseguir para a Fase 2 (CLIENT_VERSION_CONFIG_COMPATIBLE_WITH_2X) para iniciar a migração completa. Se você pular a Fase 1 e começar diretamente com a Fase 2, as entradas sem concessão serão gravadas imediatamente na tabela de concessão, impedindo permanentemente a reversão para a KCL v1 sem a limpeza manual do DynamoDB.
Etapa 4: visão geral e recomendações de configuração de KCL 3.5.x+
Para obter uma descrição detalhada das configurações introduzidas após KCL 1.x que são relevantes em KCL 3.5.x+, consulte Configurações de KCL e Configuração do cliente de migração da KCL.
Importante
Em vez de criar objetos diretamente de checkpointConfig, coordinatorConfig, leaseManagementConfig, metricsConfig, processorConfig e retrievalConfig, recomendamos usar o ConfigsBuilder para definir configurações em KCL 3.5.x+ e versões posteriores para evitar problemas de inicialização do Agendador. O ConfigsBuilder oferece uma maneira mais flexível e sustentável de configurar o aplicativo da KCL.
Configurações com valor padrão de atualização em KCL 3.5.x+
billingMode-
Na KCL versão 1.x, o valor padrão para
billingModeé definido comoPROVISIONED. No entanto, na KCL versão 3.5.x+, obillingModepadrão éPAY_PER_REQUEST(modo sob demanda). Recomendamos que você use o modo de capacidade sob demanda em sua tabela de concessões para ajustar automaticamente a capacidade com base no uso. Para obter orientações sobre como usar a capacidade provisionada para suas tabelas de concessões, consulte Best practices for the lease table with provisioned capacity mode. idleTimeBetweenReadsInMillis-
Na KCL versão 1.x, o valor padrão para
idleTimeBetweenReadsInMillisé definido como 1.000 (ou 1 segundo). A KCL versão 3.5.x+ define o valor padrão paraidleTimeBetweenReadsInMilliscomo 1.500 (ou 1,5 segundo), mas o Amazon DynamoDB Streams Kinesis Adapter substitui esse valor padrão, definindo-o como 1.000 (ou 1 segundo).
Novas configurações em KCL 3.5.x+
leaseAssignmentIntervalMillis-
Essa configuração define o intervalo de tempo antes que os fragmentos recém-descobertos comecem a ser processados, e é calculada como 1,5 ×
leaseAssignmentIntervalMillis. Se essa configuração não for definida explicitamente, o intervalo de tempo será padronizado como 1,5 ×failoverTimeMillis. O processamento de novos fragmentos exige a verificação da tabela de concessões e a consulta a um índice secundário global (GSI) na tabela de concessões. A redução deleaseAssignmentIntervalMillisaumenta a frequência dessas operações de verificação e consulta, aumentando os custos do DynamoDB. Recomendamos definir esse valor como 2 mil (ou 2 segundos) para minimizar o atraso no processamento de novos fragmentos. shardConsumerDispatchPollIntervalMillis-
Essa configuração define o intervalo entre pesquisas sucessivas feitas pelo consumidor do fragmento para acionar transições de estado. Na KCL versão 1.x, esse comportamento era controlado pelo parâmetro
idleTimeInMillis, que não era exposto como uma definição configurável. Na KCL versão 3.5.x+, recomendamos definir essa configuração para corresponder ao valor usado emidleTimeInMillisna configuração da KCL versão 1.x.
Etapa 5: migrar de KCL 2.x para KCL 3.5.x+
Para garantir uma transição tranquila e a compatibilidade com a versão mais recente da Kinesis Client Library (KCL), siga as etapas de 5 a 8 nas instruções do guia de migração para atualizar de KCL 2.x para KCL 3.5.x+.
Para solucionar problemas comuns da KCL 3.5.x+, consulte Solução de problemas de aplicativos de consumidor da KCL.