在 Google Cloud Managed Service for Apache Kafka 中访问 Kafka 数据

如果您需要将数据从开源 Apache Kafka 主题移至 Managed Service for Apache Kafka,可以使用 Dataflow 模板来完成此操作。 您可以使用 Google Cloud 控制台、REST API 或 Google Cloud CLI。 Kafka 主题中的消息必须使用 AVRO 或 JSON 架构和编码进行发布。

本文档可帮助您使用 控制台配置 Google Cloud **Kafka to Kafka Dataflow** 模板。

Google Cloud 使用的产品

Kafka to Kafka Dataflow 模板 使用以下付费 Google Cloud 产品。 您可使用价格计算器来生成 基于您的预计使用情况的费用估算。

  • Dataflow:Dataflow 是一种全代管式数据处理服务。The Kafka to Kafka Dataflow 模板利用 Dataflow 创建一个流水线,该流水线从外部 Kafka 主题读取数据,并将其写入 Managed Service for Apache Kafka 主题。 Dataflow 的自动扩缩和 自我修复功能可确保您的流水线可靠高效地运行。
  • Cloud Storage:Cloud Storage 提供 对象存储。仅当您对您的 外部 Kafka 来源使用 TLS 身份验证时,才使用此服务。您将包含 TLS 证书的密钥库和信任库 文件存储在 Cloud Storage 中,并且 Dataflow 流水线从 Cloud Storage 访问这些文件。
  • Secret Manager:Secret Manager 是一种 Secret 和凭据管理服务。仅当您为外部 Kafka 来源选择 TLS 或 SASL_PLAIN 身份验证时,才使用此产品。它会存储您的 身份验证凭据,例如用户名、密码和证书。 这些凭据用于安全地连接到外部 Kafka 集群。

此外,该解决方案还使用 Managed Service for Apache Kafka。

  • Google Cloud Managed Service for Apache Kafka:一项 Google Cloud 服务 ,可帮助您运行 Apache Kafka。在此解决方案中,Managed Service for Apache Kafka 用作数据的目标位置。您可以将数据从外部 Kafka 主题流式传输到 Managed Service for Apache Kafka 集群中的主题。如需详细了解 Managed Service for Apache Kafka 定价, 请参阅 定价指南

准备工作

在启动 Kafka to Kafka Dataflow 模板之前,请确保您已完成以下操作:

  1. 创建 Managed Service for Apache Kafka 集群和主题。

    创建集群和主题的一种方法是按照 Managed Service for Apache Kafka 快速入门中的说明进行操作。

    如果您的主题包含 Avro 记录,如需了解其他资源 要求,请参阅 指定消息格式

  2. 启用以下 Google Cloud API:

    • Dataflow

    • Secret Manager(仅当您对外部 Kafka 来源使用 SASL_PLAIN 或 TLS 身份验证时)

    • Cloud Storage(仅当您对外部 Kafka 来源使用 TLS 身份验证时)

    gcloud services enable dataflow.googleapis.com \
    secretmanager.googleapis.com storage.googleapis.com
    

准备 SASL_PLAIN 身份验证资源

如果您计划对外部 Kafka 来源使用 SASL_PLAIN 身份验证凭据,请按照本部分中的说明进行操作。

  1. 在 Google Cloud 控制台中,前往 Secret Manager。

    前往 Secret Manager

  2. 创建两个 Secret:一个用于 Kafka 用户名,另一个用于 Kafka 密码。

    为每个 Secret 提供一个描述性名称,例如“kafka-username”和“kafka-password”。

如需详细了解如何创建 Secret(包括所需的 角色和权限),请参阅创建 Secret

准备 TLS 身份验证资源

如果您计划对外部 Kafka 来源使用 TLS 身份验证凭据,请按照本部分中的说明进行操作。

  1. 在 Google Cloud 控制台中,前往 Secret Manager。

    前往 Secret Manager

  2. 确保您的密钥库文件(JKS 格式)包含向 Kafka 集群验证客户端所需的 TLS 证书和私钥。

  3. 确保您的信任库文件(JKS 格式)包含验证 Kafka 代理身份所需的受信任证书。

  4. 将密钥库和信任库文件上传到 Cloud Storage。

    如需详细了解如何将文件上传到 Cloud Storage, 请参阅从文件系统上传对象

  5. 为 TLS 配置创建三个 Secret:

    • 信任库密码:存储用于访问信任库文件的密码。

    • 密钥库密码:存储用于访问密钥库文件的密码。

    • 私钥密码:存储用于访问密钥库文件中的私钥的密码。

    为每个 Secret 提供一个描述性名称。

如需详细了解如何创建 Secret(包括所需的 角色和权限),请参阅创建 Secret

向 Dataflow 工作器服务帐号授予 Managed Kafka 客户端角色

如需将 Dataflow 作业连接到 Managed Service for Apache Kafka,您需要向 Dataflow 工作器服务帐号授予特定权限。此服务帐号是 Dataflow 作业中所有工作器虚拟机的身份,并且从这些虚拟机发出的任何请求都使用此账号。

如需允许访问 Kafka 资源,您必须向 Dataflow 工作器服务帐号授予 roles/managedkafka.client 角色。此角色包含建立连接所需的 managedkafka.clusters.connect 权限。

如需详细了解工作器服务帐号,请参阅 Google Cloud 上流水线的安全性和权限

如需向 Dataflow 服务帐号授予 Managed Kafka 客户端角色,请按照以下步骤操作:

控制台

  1. 在 Google Cloud 控制台中,前往 IAM 页面。
    前往 IAM
  2. 检查项目是否设置为 Managed Service for Apache Kafka 客户端将访问的使用方项目。
  3. 点击授予访问权限
  4. 在新页面中,对于添加主账号,输入您正在使用的 Dataflow 工作器服务帐号的电子邮件 地址。
  5. 对于分配角色 ,选择 Managed Kafka 客户端 角色。
  6. 点击保存

gcloud CLI

  1. 在 Google Cloud 控制台中,激活 Cloud Shell。

    激活 Cloud Shell

    Cloud Shell 会话随即会在控制台的底部启动,并显示命令行提示符。 Google Cloud Cloud Shell 是一个已安装 Google Cloud CLI 且已为当前项目设置值的 Shell 环境 。该会话可能需要几秒钟来完成初始化。

  2. 运行 gcloud projects add-iam-policy-binding 命令:

    gcloud projects add-iam-policy-binding PROJECT_ID \
      --member serviceAccount:SERVICE_ACCOUNT_EMAIL \
      --role roles/managedkafka.client

    替换以下内容:

    • PROJECT_ID 是项目 ID。

    • SERVICE_ACCOUNT_EMAIL 是 Dataflow 工作器服务帐号的电子邮件地址。

启动 Kafka to Kafka Dataflow 模板

您可以从控制台的集群详情页面启动 Kafka to Kafka Dataflow 模板。

  1. 在 Google Cloud 控制台中,前往集群页面。

    前往集群

    系统会列出您在项目中创建的集群。

  2. 如需查看集群详情页面,请点击集群名称。
  3. 在集群详情页面中,点击导入数据

    系统会打开使用“Kafka to Kafka”模板创建 Dataflow 作业 页面。

根据以下部分中包含的信息配置模板中的字段。

输入作业名称

对于作业名称 字段,输入 Dataflow 作业的名称。

该名称在项目中所有当前运行的作业中必须是唯一的。

为流水线选择区域端点

对于区域端点 字段,将区域端点设置为 Kafka 集群的位置,以最大限度地减少跨区域数据传输费用。

Dataflow 工作器可以独立于 Kafka 集群的区域运行。但是,如果您在 Kafka 集群区域之外启动工作器,则会产生区域间出站流量费用。

如需查看集群的位置,请按照 列出 Managed Service for Apache Kafka 集群中的步骤操作。

配置来源

  1. 对于来源,选择自行管理的 Kafka 或外部 Kafka

  2. 对于 Kafka 引导服务器,输入外部 Kafka 集群引导服务器的地址。

  3. 对于来源 Kafka 引导服务器和主题,输入 外部 Kafka 引导服务器和主题地址。

  4. 对于 Kafka 来源身份验证模式 ,从提供的选项中选择您的身份验证模式 。对于外部 Kafka,您有三个选项:

    • SASL_PLAIN

    • TLS

对于 SASL_PLAIN Kafka 来源身份验证

在继续本部分之前,请参阅 准备 SASL_PLAIN 身份验证资源

如果您选择 SASL_PLAIN 作为 Kafka 来源身份验证,请填写以下其他字段。

  1. 对于 Kafka SASL/PLAIN 用户名的 Secret 版本 ID,输入 用户名的版本 ID。

  2. 对于 Kafka SASL/PLAIN 密码的 Secret 版本 ID,输入 密码的版本 ID。

对于 TLS Kafka 来源身份验证

在继续本部分之前,请参阅 准备 TLS 身份验证资源

如果您选择 TLS 作为 Kafka 来源身份验证,请填写以下其他字段。

  1. 对于密钥库的位置,输入 您的 Java 密钥库 (JKS) 文件的 Cloud Storage 路径。此文件包含向 Kafka 集群验证流水线所需的 TLS 证书和私钥。例如:gs://your-bucket/keystore.jks

  2. 对于信任库文件位置 字段,输入 Java 信任库 (JKS) 文件的 Cloud Storage 路径。此文件包含用于验证 Kafka 代理身份的受信任证书。

  3. 对于信任库密码的 Secret 版本 ID,提供 存储用于访问信任库 JKS 文件的密码的 Secret Manager Secret ID。

  4. 对于密钥库密码的 Secret 版本 ID,提供 存储用于访问 JKS 文件的密码的 Secret Manager Secret ID 。

  5. 对于私钥密码的 Secret 版本 ID 字段,如果密钥库中的私钥有单独的密码,请输入包含该密码的 Secret Manager Secret ID。

将偏移量提交到 Kafka

此设置用于控制流水线如何跟踪外部 Kafka 来源中的哪些消息已处理。启用此设置有助于确保流水线可靠地处理消息,并避免在重启时出错。

Kafka 使用偏移量跟踪消息。您可以将这些偏移量视为书签,用于记住流水线在每个消息流中的位置。启用此设置后,流水线会定期将书签保存在 Kafka 中。如果流水线因错误、更新或计划维护而重启,则可以使用已保存的书签从中断的位置继续。这样可以防止流水线遗漏消息或意外地处理同一消息两次。

  1. 启用将偏移量提交到 Kafka

  2. 提供消费者组 ID 。 这是用于向 Kafka 标识流水线的名称。

  3. 对于 Kafka 默认起始偏移量,请选择以下选项之一:

    • 最早:从 Kafka 主题的开头开始处理消息。

    • 最新:从最新消息开始处理消息。

配置目标位置

  1. 对于目标位置,选择 Managed Service for Apache Kafka

  2. 对于 Kafka 集群 ,选择您已创建的集群。

  3. 您可以选择现有主题,也可以创建一个新主题。

  4. 对于 Kafka 目标身份验证方法 ,选择 应用默认凭证

配置加密

默认情况下,所有静态数据和传输中的数据都由 Google-owned and Google-managed encryption key进行加密。如果您有客户管理的加密密钥 (CMEK),则可以选择自己的密钥。 如需详细了解如何配置 CMEK,请参阅 配置消息加密

配置网络

您必须在 Dataflow 模板中指定集群的网络和子网。您可以在模板的可选参数 部分中为 Dataflow 工作器定义网络。

默认情况下,Kafka to Kafka Dataflow 模板会在项目的默认网络中预配 Dataflow 工作器。 如需让外部 Kafka 集群通过 Dataflow 向 Managed Service for Apache Kafka 发送数据,请确保 Dataflow 工作器可以访问集群的网络。

我们建议,如果 Kafka 集群未连接到项目默认网络中的子网,请为 Kafka 集群使用项目的默认网络。

如需详细了解如何使用 Dataflow 流水线设置网络,请参阅以下内容:

如果您在配置 Dataflow 网络时遇到问题, 请参阅 Dataflow 网络问题排查指南

配置可选 Dataflow 参数

只有在了解配置对 Dataflow 工作器的影响时,才配置可选 Dataflow 参数。设置不正确可能会影响性能或费用。如需详细了解每个选项,请参阅 可选参数

监控

Kafka to Managed Service for Apache Kafka 的 Dataflow 模板提供监控体验,让您可以在控制台中探索日志、指标和错误。此监控工具套件作为 Dataflow 用户界面的一部分提供。

您可以在作业指标 标签页中创建自定义信息中心。对于 Kafka to Managed Service for Apache Kafka 模板,我们建议设置一个作业指标信息中心,用于监控以下内容:

  • 吞吐量:在任意时间点处理的数据量。这有助于监控作业中的数据传输并发现潜在的性能问题。

    如需了解详情,请参阅 Dataflow 吞吐量监控

  • 数据新鲜度:数据元素的时间戳与流水线处理事件的时间之间的差值(以秒为单位)。这有助于发现性能和数据源瓶颈或频繁重试。

    如需了解详情,请参阅以下内容: Dataflow 新鲜度监控

  • 积压数据:等待处理的字节数。 此信息可为 自动扩缩 决策提供依据。

如需详细了解 Dataflow 监控, 请参阅 Dataflow 监控文档

问题排查

如果 Dataflow 流水线遇到性能问题,Dataflow 会提供一套全面的问题排查和诊断工具。

以下是两个常见场景及其相应的问题排查指南:

如需大致了解如何调试 Dataflow 流水线,请参阅 排查和调试 Dataflow 流水线