View a markdown version of this page

使用 Kinesis 客户端库 (KCL) 处理 Amazon Keyspaces 流 - Amazon Keyspaces(Apache Cassandra 兼容)

本文属于机器翻译版本。若本译文内容与英语原文存在差异,则一律以英文原文为准。

使用 Kinesis 客户端库 (KCL) 处理 Amazon Keyspaces 流

本主题介绍如何使用 Kinesis 客户端库 (KCL) 使用和处理来自亚马逊密钥空间变更数据捕获 (CDC) 流的数据。

与其直接使用 Amazon Keyspaces Streams API,不如使用 Kinesis 客户端库 (KCL) 提供许多好处,例如:

  • 内置分片谱系跟踪和迭代器处理。

  • 工作人员之间的自动负载平衡。

  • 容错能力和从工作人员故障中恢复。

  • 检查点以跟踪处理进度。

  • 适应流容量的变化。

  • 简化了处理 CDC 记录的分布式计算。

下一节概述了使用 Kinesis 客户端库 (KCL) 处理数据流的原因和方法,并提供了使用 KCL 处理亚马逊密钥空间 CDC 流的示例。

有关定价的信息,请参阅 Amazon Keyspaces(Apache Cassandra 兼容)定价

什么是 Kinesis Client Library?

Kinesis 客户端库 (KCL) 是一个独立的 Java 软件库,旨在简化消费和处理来自流的数据的过程。KCL 可处理许多与分布式计算相关的复杂任务,让您在处理流数据时专注于实现业务逻辑。KCL 管理的活动包括跨多个工作线程进行负载平衡、响应工作器故障、检查已处理记录以及响应流中分片数量的变化。

要处理 Amazon Keyspaces CDC 流,您可以使用 KCL 中的设计模式来处理流分片和流记录。KCL 提供低级 Kinesis Data Streams API 之上的有用抽象来简化编码。有关 KCL 的更多信息,请参阅 Amazon Kinesis Data Streams 开发人员指南中的使用 KCL 开发消费者。

要使用 KCL 编写应用程序,请使用 Amazon Keyspaces Streams Kinesis 适配器。Kinesis 适配器实现了 Kinesis 数据流接口,因此您可以使用 KCL 来消费和处理来自亚马逊 Keyspaces 流的记录。有关如何设置和安装亚马逊 Keyspaces 流 Kinesis 适配器的说明,请访问存储库。 GitHub

下图显示了这些库如何相互交互。

处理亚马逊 Keyspaces CDC 流记录时,客户端应用程序与 Kinesis Data Streams、KCL、Amazon Keyspaces Streams Kinesis 适配器以及亚马逊 Keyspaces API 之间的交互。

KCL 经常更新,以纳入新版底层库、安全改进和错误修复。建议使用最新版本的 KCL,以避免出现已知问题,并从所有最新的改进中受益。要查找最新的 KCL 版本,请参阅 KCL 存储库 GitHub 。

KCL 概念

在使用 KCL 实现消费者应用程序之前,您应该了解以下概念:

KCL 消费者应用程序

KCL 消费者应用程序是处理来自亚马逊 Keyspaces CDC 流的数据的程序。KCL 充当您的消费者应用程序代码和亚马逊 Keyspaces CDC 流之间的中介。

工人

工作程序是您的 KCL 消费者应用程序的执行单元,用于处理来自亚马逊 Keyspaces CDC 流的数据。您的应用程序可以运行分布在多个实例上的多个工作程序。

录制处理器

记录处理器是应用程序中的逻辑,用于处理来自 Amazon Keyspaces CDC 流中分片的数据。记录处理器由工作人员为其管理的每个分片进行实例化。

租赁

租约代表分片的处理责任。工作人员使用租约来协调哪个工作人员正在处理哪个分片。KCL 将租赁数据存储在亚马逊 DynamoDB 的表中。

检查点

检查点是记录处理器成功处理记录的分片中位置的记录。Checkpointing 使您的应用程序能够在工作人员出现故障时从中断的地方恢复处理。

有了亚马逊Keyspaces Kinesis适配器,您就可以开始使用KCL接口进行开发,API调用可以无缝地指向亚马逊Keyspaces流终端节点。有关可用端点的列表,请参阅如何访问 Amazon Keyspaces 中的 CDC 流终端节点

应用程序启动后,调用 KCL 来实例化工作进程。您必须向工作程序提供应用程序的配置信息,例如流描述符和 AWS 凭据,以及您提供的记录处理器类的名称。在记录处理器中运行代码时,工作进程执行以下任务:

  • 连接到流

  • 枚举流中的分片

  • 协调与其他工作程序的分片关联(如果有)

  • 为其管理的每个分片实例化记录处理器

  • 从流中提取记录

  • 将记录推送到对应的记录处理器

  • 对已处理记录进行检查点操作

  • 在工作程序实例计数更改时均衡分片与工作程序的关联

  • 在分片被拆分时平衡分片与工作程序的关联