访问 Cloud Storage 中的 Kafka 数据

如果您需要将数据从 Google Cloud Managed Service for Apache Kafka 主题加载到 Cloud Storage 存储桶,可以使用 Dataflow 模板来完成此操作。 您可以使用 Google Cloud 控制台、REST API 或 Google Cloud CLI。

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

Google Cloud 使用的产品

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

  • Dataflow:Dataflow 是一种全代管式数据处理服务。The Kafka to Cloud Storage Dataflow 模板利用 Dataflow 创建一个流水线,该流水线从您的 Kafka 主题读取数据,执行任何必要的转换,然后将其写入 Cloud Storage。Dataflow 的自动扩缩和 自我修复功能可确保流水线可靠高效地运行。
  • Cloud Storage:用作 Kafka 数据的目标位置。 您需要一个 Cloud Storage 存储桶来存储由 Dataflow 流水线传输的数据。

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

  • Google Cloud Managed Service for Apache Kafka:一项 Google Cloud 服务 ,可帮助您运行 Apache Kafka。为流水线提供源数据。 您需要一个现有的 Managed Service for Apache Kafka 集群和主题 ,其中包含您要转移到 Cloud Storage 的数据。如需详细了解 Google Cloud Managed Service for Apache Kafka 的价格,请参阅 价格指南

准备工作

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

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

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

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

  2. 已启用以下 Google Cloud API:

    • Dataflow

    • Cloud Storage

    gcloud services enable dataflow.googleapis.com storage-api.googleapis.com \
    
  3. 创建 Cloud Storage 存储桶。

    如需详细了解如何创建 Cloud Storage 存储桶, 请参阅创建存储桶

向 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 Cloud Storage Dataflow 模板

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

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

    转到集群

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

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

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

  4. 在模板中,对于 Dataflow 模板,将 模板更新为 Kafka to Cloud Storage

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

输入作业名称

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

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

为流水线选择区域端点

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

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

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

配置来源

  1. 对于来源,保留默认值 Managed Service for Apache Kafka

  2. 对于 Kafka 集群Kafka 来源身份验证模式 ,保留默认值。

  3. 对于 Kafka 主题,从可用主题列表中选择一个主题。

配置 Kafka 消息格式

Dataflow 模板支持以下三种消息格式:

  • Avro Confluent 线路格式:每条 Kafka 消息都包含一个 magic 字节, 一个架构 ID 和 Avro 二进制编码的记录。

    对于 Avro(Confluent 线路格式),您可以使用单个架构或多个架构:

    • 单个架构:所有消息都遵循单个预定义的 Avro 架构。

    • 多个架构:消息可以使用不同的架构。仅 Avro(Confluent 线路格式)支持此功能。

  • Avro(二进制编码) :消息仅包含记录的载荷,不包含任何元数据。您必须提供上传到 Cloud Storage 的 Avro 架构文件 (.avsc)。所有消息都必须遵循此单个架构。

  • JSON:记录不需要预定义的架构。不符合架构的记录会发送到死信队列(如果已配置),或者系统会记录错误消息。支持的 格式为 {"field": "value"} 格式。不支持格式 [{"name": "field", "value": "value"}]

Google Cloud Managed Service for Apache Kafka 不提供架构注册表。该模板仅支持将身份验证凭据传递给与 Confluent 线路格式兼容的架构注册表。

Avro Confluent 线路格式

如果您选择此选项作为 Kafka 消息格式,请配置以下 其他设置:

架构来源:此字段告知流水线在哪里查找架构。 请从下列选项中选择一项:

  • 架构注册表:您的架构存储在 Confluent 架构注册表中。这对于不断发展的架构和管理多个版本非常有用。 确保架构注册表可供 Managed Service for Apache Kafka 集群网络访问,并且该注册表与 Dataflow 工作器托管在同一区域中。您可以在单个架构和多个架构场景中使用架构注册表。配置以下其他设置:

    • 架构注册表连接网址:提供用于连接 到架构注册表的网址。

    • 身份验证模式:如果您的注册表需要身份验证, 请选择 OAuthTLS。否则,请选择 None

  • 单个架构文件:如果您的所有消息都遵循文件中定义的 单个固定架构,请选择此选项。

    • Cloud Storage 文件到 Avro 架构文件:用于解码主题中所有消息的 Avro 架构文件的路径。

Avro 二进制编码

如果您选择此选项作为 Kafka 消息格式,请配置以下 其他设置:

  • Cloud Storage 文件到 Avro 架构文件:用于解码主题中所有消息的 Avro 架构文件的路径。

JSON

如果您选择此选项作为 Kafka 消息格式 ,则无需进行其他 配置。

指定 Kafka 偏移量

  1. 如需避免在需要重启单个工作器或整个流水线时重新处理消息,请选择将偏移量提交到 Kafka 选项。 这样可确保流水线从中断的位置恢复处理,从而防止重复处理和潜在的数据不一致问题。

  2. 对于输入使用方群组 ID 字段,为此流水线的群组输入一个唯一名称。在大多数情况下,您希望流水线读取每条消息一次,并且可以重启。

  3. 对于默认 Kafka 启动偏移量 字段,Dataflow 流水线提供两个启动偏移量选项。请从下列选项中选择一项:

    • Earliest:从 Kafka 主题的开头处理消息。

    • 最新:从最新的可用偏移量开始处理消息。

配置目标位置

这些选项用于控制数据流水线将数据写入 Cloud Storage 的方式。

  1. 对于目标位置,输入存储桶路径,并添加 输出文件的文件名前缀。文件名前缀必须以斜杠结尾。例如, gs://test-bucket/test-prefix/

  2. 对于窗口时长,输入将数据写入 Cloud Storage 的时间窗口。 根据您的数据处理要求,选择适当的格式(Ns 表示秒,Nm 表示分钟,Nh 表示小时)。

  3. 对于要写入的文件的输出文件名前缀,您可以提供一个前缀 以添加到每个输出文件中,以便更好地组织和识别。

  4. 对于输出分片数上限,将该数字设置为零。 您可以指定写入文件时要生成的分片数。增加分片数可以提高吞吐量,但也会因混洗费用较高而导致费用增加。当您将该数字设置为零时,服务会选择最佳数量。

配置死信队列

有时,由于损坏、数据类型不兼容或架构不匹配,消息无法处理。

如需处理这些情况,请在模板中启用死信队列并提供表名称。模板会使用标准化架构创建表。

配置加密

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

配置网络

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

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

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

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

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

配置可选 Dataflow 参数

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

监控

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

通过作业指标 标签页,您可以创建自定义信息中心。对于 Kafka to Cloud Storage Dataflow 模板,我们建议设置一个作业指标 信息中心,以监控以下内容:

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

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

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

    如需了解详情,请参阅 Dataflow 数据新鲜度监控

  • 积压:等待处理的字节数。 此信息会影响自动扩缩决策。

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

问题排查

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

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

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

已知限制

  • 该模板不支持将凭据传递给架构注册表以进行身份验证。

  • 创建 Kafka to Cloud Storage Dataflow 作业时,请确保 项目设置为包含 Managed Service for Apache Kafka 集群的同一项目。 Google Cloud

Apache Kafka® 是 Apache Software Foundation 或其关联公司在美国和/或其他国家/地区的注册商标。

后续步骤