View a markdown version of this page

从 KCL 1.x 迁移到 KCL 3.5.x+ - Amazon DynamoDB

从 KCL 1.x 迁移到 KCL 3.5.x+

概述

本指南提供有关将使用者应用程序从 KCL 1.x 迁移到 KCL 3.5.x+ 的说明。由于 KCL 1.x 和 KCL 3.5.x+ 之间的架构差异,迁移需要更新多个组件以确保兼容性。

与 KCL 3.5.x+ 相比,KCL 1.x 使用不同的类和接口。您必须先将记录处理器、记录处理器工厂和工作线程类迁移到 KCL 3.5.x+ 兼容格式,然后按照将 KCL 1.x 迁移到 KCL 3.5.x+ 的迁移步骤进行操作。

注意

GitHub 网站上的 Amazon DynamoDB Streams Kinesis Adapter 版本 2.4.x+ 支持 KCL 3.5.x+。

迁移步骤

步骤 1:迁移记录处理器

以下示例显示了为 KCL 1.x DynamoDB Streams Kinesis Adapter 实现的记录处理器:

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(); } } }
迁移 RecordProcessor 类
  1. 将接口从 com.amazonaws.services.kinesis.clientlibrary.interfaces.v2.IRecordProcessorcom.amazonaws.services.kinesis.clientlibrary.interfaces.v2.IShutdownNotificationAware 更改为 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. 更新 initializeprocessRecords 方法的导入语句:

    // 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. 使用以下新方法替换 shutdownRequested 方法:leaseLostshardEndedshutdownRequested

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

下面是记录处理器类的更新版本:

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

DynamoDB Streams Kinesis Adapter 现在使用 SDKv2 记录模型。在 SDKv2 中,复杂的 AttributeValue 对象(BSNSMLSS)从不会返回 null。使用 hasBs()hasNs()hasM()hasL()hasSs() 方法来验证这些值是否存在。

步骤 2:迁移记录处理器工厂

记录处理器工厂负责在获得租约时创建记录处理器。下面是 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); } }
迁移到 RecordProcessorFactory
  • 将已实施的接口从 com.amazonaws.services.kinesis.clientlibrary.interfaces.v2.IRecordProcessorFactory 更改为 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() {

下面是 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(); } }

步骤 3:迁移工作线程

在 KCL 版本 3.5.x+ 中,名为 Scheduler 的新类取代了 Worker 类。下面是 KCL 1.x 工作线程的示例:

final KinesisClientLibConfiguration config = new KinesisClientLibConfiguration(...) final IRecordProcessorFactory recordProcessorFactory = new RecordProcessorFactory(); final Worker worker = StreamsWorkerFactory.createDynamoDbStreamsWorker( recordProcessorFactory, workerConfig, adapterClient, amazonDynamoDB, amazonCloudWatchClient);
迁移工作程序
  1. import 类的 Worker 语句更改为 SchedulerConfigsBuilder 类的导入语句。

    // import com.amazonaws.services.kinesis.clientlibrary.lib.worker.Worker; import software.amazon.kinesis.coordinator.Scheduler; import software.amazon.kinesis.common.ConfigsBuilder;
  2. 导入 StreamTracker 并将 StreamsWorkerFactory 的导入更改为 StreamsSchedulerFactory

    import software.amazon.kinesis.processor.StreamTracker; // import software.amazon.dynamodb.streamsadapter.StreamsWorkerFactory; import software.amazon.dynamodb.streamsadapter.StreamsSchedulerFactory;
  3. 选择从中启动应用程序的位置。它可以为 TRIM_HORIZONLATEST

    import software.amazon.kinesis.common.InitialPositionInStream; import software.amazon.kinesis.common.InitialPositionInStreamExtended;
  4. 创建一个 StreamTracker 实例。

    StreamTracker streamTracker = StreamsSchedulerFactory.createSingleStreamTracker( streamArn, InitialPositionInStreamExtended.newInitialPosition(InitialPositionInStream.TRIM_HORIZON) );
  5. 创建 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. 创建 ConfigsBuilder 对象。

    import software.amazon.kinesis.common.ConfigsBuilder; ... ConfigsBuilder configsBuilder = new ConfigsBuilder( streamTracker, applicationName, adapterClient, dynamoDbAsyncClient, cloudWatchAsyncClient, UUID.randomUUID().toString(), new StreamsRecordProcessorFactory());
  7. 使用 ConfigsBuilder 创建 Scheduler,如以下示例所示:

    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 );
注意

KCL 3.5.x+ 迁移分为三个阶段:

  • 第 1 阶段 (CLIENT_VERSION_CONFIG_COMPATIBLE_WITH_2X_PHASE1):纯 KCL 1.x 兼容模式。不会向租约表中写入任何特定于迁移的元数据。可通过重新部署之前的代码安全回滚到 KCL v1。使用此阶段来验证稳定性。

  • 第 2 阶段 (CLIENT_VERSION_CONFIG_COMPATIBLE_WITH_2X):开始迁移。将 WORKER_METRIC_STATSMigration3.0 条目写入租约表。在所有工作线程都准备就绪时,KCL 自动过渡到完整的 3.x 负载均衡。支持回滚到第 1 阶段(通过 KCL 迁移工具)。无法再回滚到 KCL v1。

  • 第 3 阶段 (CLIENT_VERSION_CONFIG_3X):完整 KCL 3.x 功能。由客户显式设置,或者在配置移除时用作默认值。最终状态,无法回滚。

这些设置在 KCL v3 与 KCL v1 的 DynamoDB Streams Kinesis Adapter 之间保持兼容性,而非在 KCL v2 和 v3 之间保持兼容性。

重要

您必须从 CLIENT_VERSION_CONFIG_COMPATIBLE_WITH_2X_PHASE1(第 1 阶段)开始迁移。第 1 阶段与 KCL v1 向后兼容,并且不会向租约表写入任何特定于迁移的条目,因此您只需重新部署之前的代码,即可安全地回滚到之前的 KCL 版本。在第 1 阶段完成了全面的烘焙测试后,您可以进入第 2 阶段 (CLIENT_VERSION_CONFIG_COMPATIBLE_WITH_2X) 开始完整迁移。如果您跳过第 1 阶段并直接从第 2 阶段开始,则非租约条目会立即写入租约表,从而永久阻止回滚到 KCL v1,除非您手动清理 DynamoDB。

步骤 4:KCL 3.5.x+ 配置概述和建议

有关 KCL 1.x 之后引入的与 KCL 3.5.x+ 相关的配置的详细描述,请参阅 KCL 配置KCL 迁移客户端配置

重要

在 KCL 3.5.x+ 及更高版本中,我们建议不要直接创建 checkpointConfigcoordinatorConfigleaseManagementConfigmetricsConfigprocessorConfigretrievalConfig 的对象,而是使用 ConfigsBuilder 设置配置,来避免出现调度器初始化问题。ConfigsBuilder 提供了更灵活且更易于维护的方式来配置 KCL 应用程序。

在 KCL 3.5.x+ 中具有更新默认值的配置

billingMode

在 KCL 版本 1.x 中,billingMode 的默认值设置为 PROVISIONED。但在 KCL 版本 3.5.x+ 中,默认 billingModePAY_PER_REQUEST(按需模式)。我们建议您对租约表使用按需容量模式,以便根据您的使用情况自动调整容量。有关对租约表使用预置容量的指导,请参阅 Best practices for the lease table with provisioned capacity mode

idleTimeBetweenReadsInMillis

在 KCL 版本 1.x 中,idleTimeBetweenReadsInMillis 的默认值设置为 1000(或 1 秒)。KCL 版本 3.5.x+ 将 idleTimeBetweenReadsInMillis 的默认值设置为 1500(即 1.5 秒),但 Amazon DynamoDB Streams Kinesis Adapter 将默认值改写为 1000(即 1 秒)。

KCL 3.5.x+ 中的新配置

leaseAssignmentIntervalMillis

此配置定义在新发现的分片开始处理之前的时间间隔,计算方法为 1.5 × leaseAssignmentIntervalMillis。如果未显式配置此设置,则时间间隔默认为 1.5 × failoverTimeMillis。处理新分片包括扫描租约表并在租约表上查询全局二级索引(GSI)。降低 leaseAssignmentIntervalMillis 会增加这些扫描和查询操作的频率,从而导致 DynamoDB 成本更高。建议将此值设置为 2000(即 2 秒),以尽可能减少处理新分片的延迟。

shardConsumerDispatchPollIntervalMillis

此配置定义了分片使用者用于触发状态转换的连续轮询之间的间隔。在 KCL 版本 1.x 中,此行为由 idleTimeInMillis 参数控制,该参数未作为可配置的设置公开。在 KCL 版本 3.5.x+ 中,我们建议将此配置设置为匹配 KCL 版本 1.x 设置中用于 idleTimeInMillis 的值。

步骤 5:从 KCL 2.x 迁移到 KCL 3.5.x+

为确保平稳过渡并与最新的 Kinesis Client Library(KCL)版本兼容,请按照迁移指南的从 KCL 2.x 升级到 KCL 3.5.x+ 说明中的步骤 5-8 进行操作。

有关常见的 KCL 3.5.x+ 问题故障排除,请参阅 KCL 使用者应用程序故障排除