在 BigQuery 中访问 Kafka 数据

如果您需要将 AVRO 或 JSON 消息从 Google Cloud Managed Service for Apache Kafka 主题加载到 BigQuery 表中,可以使用 Dataflow 模板来完成此操作。本文档可帮助您使用 Google Cloud 控制台基于此模板设置流水线。您还可以使用 REST API 或 Google Cloud CLI 设置此流水线。

使用的Google Cloud 款产品

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

  • BigQuery:BigQuery 是 Google Cloud的无服务器数据仓库。在此解决方案中,BigQuery 用作 Kafka 数据的目标。BigQuery Streaming Write API 用于将数据写入 BigQuery。然后,您可以使用 BigQuery 通过 SQL 查询分析这些数据,使用 BigQuery ML 构建机器学习模型,并为商业智能应用提供支持。
  • Dataflow:Dataflow 是一种全代管式数据处理服务。Kafka to BigQuery Dataflow 利用 Dataflow 创建一个流水线,该流水线可从 Kafka 主题读取数据、执行任何必要的转换,并将数据写入 BigQuery。Dataflow 的自动扩缩和自我修复功能可确保流水线可靠高效地运行。
  • Cloud Storage:虽然 Cloud Storage 不直接参与核心数据传输,但如果您使用 Avro 二进制编码来编码 Kafka 消息,则可以使用 Cloud Storage 来存储 Avro 架构文件。这样,Dataflow 流水线便可以正确解读数据结构。

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

  • Google Cloud Managed Service for Apache Kafka:可帮助您运行 Apache Kafka 的 Google Cloud 服务。借助 Managed Service for Apache Kafka,您可以专注于构建事件驱动型系统和流式数据流水线,而无需费心管理基础设施。如需详细了解 Google Cloud Managed Service for Apache Kafka 定价,请参阅价格指南

准备工作

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

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

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

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

  2. 启用以下 Google Cloud API:

    • Dataflow

    • BigQuery

    • Cloud Storage

    gcloud services enable dataflow.googleapis.com bigquery.googleapis.com \
    storage.googleapis.com
    
  3. 创建 BigQuery 数据集和表。确保表的架构与 Kafka 输入主题的架构匹配。

    如果在同一主题中使用多个架构并写入多个表,则无需在配置流水线之前创建表。

    如需详细了解如何创建 BigQuery 数据集和表,请参阅创建数据集创建具有架构定义的空表

向 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 Worker 服务账号的电子邮件地址。

启动 Kafka to BigQuery Dataflow 模板

您可以在控制台的集群详情页面中启动 Kafka 到 BigQuery Dataflow 模板。

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

    前往“集群”

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

  2. 如需查看集群详情页面,请点击集群名称。
  3. 在集群详情页面中,对于主题,点击任意主题的 BigQuery Export 图标。

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

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

输入作业名称

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

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

为流水线选择区域端点

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

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

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

选择 Dataflow 模板

对于 Dataflow 模板,保留默认值 Kafka to BigQuery

配置来源

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

  2. 系统已根据您选择的主题预先选择 Kafka 集群Kafka 主题Kafka 来源身份验证模式的默认值。保留这些值。

配置 Kafka 消息格式

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

  • Avro Confluent 传输格式:每条 Kafka 消息都包含一个魔数、一个架构 ID 和一个 Avro 二进制编码记录。默认情况下,消息键不会加载到 BigQuery 中。

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

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

    • 多个架构:消息可以利用不同的架构。此功能仅适用于 Avro(Confluent 线格式)。

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

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

Avro Confluent 有线格式

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

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

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

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

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

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

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

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

Avro 二进制编码

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

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

JSON

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

配置为将 Kafka 键持久保存到 BigQuery 表

如果您启用将 Kafka 消息键持久保存到 BigQuery 表选项,还必须在目标表中添加一个名为 _key 且类型为 BYTES 的字段。虽然键可能表示结构化数据,但模板会将其视为字节数组。

指定偏移量选项

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

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

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

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

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

配置目标位置

这些选项用于控制数据流水线如何将数据写入 BigQuery。

表名称策略字段用于确定如何将数据整理到 BigQuery 表中。您有以下两种选择:

  • 单个表名称(默认):Kafka 主题中的所有数据都会进入单个 BigQuery 表。这是最简单的选项,如果您的 Kafka 数据具有相同的架构,建议选择此选项。

    对于 BigQuery 输出表,请指定将存储 Kafka 主题中所有数据的表名称。

  • 动态表名称:根据 Kafka 消息的架构创建多个 BigQuery 表。如果您的 Kafka 消息具有不同的架构,我们建议您选择此选项。只有当流水线使用 Avro Confluent Wire 或架构注册表来处理消息结构时,此选项才有效。

    Kafka 消息中的每个唯一架构都会在 BigQuery 中获得自己的表。 系统会根据架构名称自动生成表名称。

    如果您选择动态表名,请配置以下其他字段:

    • BigQuery 输出项目:BigQuery 数据集所在的 Google Cloud 项目。

    • BigQuery 输出数据集:项目内用于创建表的相应数据集。

    • BigQuery 表命名前缀(可选):为自动生成的表名称添加前缀,以便更好地组织。

配置 BigQuery Storage Write API

对于 BigQuery Storage Write API 的流数量,请从 0(默认值)开始。此选项可让流水线自动确定最佳数据流数量。

此设置用于控制将数据写入 BigQuery 时使用的并行流数量。较高的数值可以提高写入速度,尤其是在处理大量数据时。最佳数量取决于您的数据和 BigQuery 设置。

配置触发频率

对于 BigQuery Write API 的触发频率(以秒为单位)选项,请从 5 秒(默认值)开始。对于大多数用例,这都是一个不错的平衡点。

此设置决定了数据写入 BigQuery 的频率。

数字越小,写入频率越高。这可以缩短延迟时间,但可能会增加 BigQuery 费用。数值越大,写入频率越低。这种方式的成本效益更高,但可能会增加延迟时间。

配置死信队列

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

为处理这些情况,请在模板中启用死信队列并提供表名称。该模板会使用标准化架构创建表。错误消息会写入单独的 BigQuery 表。

配置加密

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

配置网络

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

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

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

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

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

配置可选的 Dataflow 参数

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

监控

Kafka 到 BigQuery 的 Dataflow 模板提供了一种监控体验,可让您在控制台中探索日志、指标和错误。这套监控工具是 Dataflow 用户界面的一部分。

借助作业指标标签页,您可以创建自定义信息中心。对于“Kafka 到 BigQuery”模板,我们建议设置一个作业指标信息中心,用于监控以下内容:

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

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

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

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

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

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

问题排查

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

以下是两种常见情况以及各自的问题排查指南:

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

已知限制

  • 该模板不支持将凭据传递给您的 Schema Registry 以进行身份验证。

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

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

后续步骤