将 Kafka 消息写入 BigQuery

本教程介绍了如何使用 Kafka Connect 将消息从 Managed Service for Apache Kafka 集群写入 BigQuery。

在本教程中,您将创建一个 Connect 集群,然后使用 BigQuery sink 连接器写入现有 BigQuery 表。在此场景中,BigQuery 表定义了 Kafka 消息的架构;Kafka 消息必须与表架构匹配。如需了解详情,请参阅 BigQuery sink 连接器的架构

准备工作

控制台

  1. 登录您的 Google Cloud 账号。如果您是 Google Cloud新手,请 创建一个账号来评估我们的产品在实际场景中的表现。新客户还可获享 $300 赠金,用于运行、测试和部署工作负载。
  2. In the Google Cloud console, on the project selector page, select or create a Google Cloud project.

    Roles required to select or create a project

    • Select a project: Selecting a project doesn't require a specific IAM role—you can select any project that you've been granted a role on.
    • Create a project: To create a project, you need the Project Creator role (roles/resourcemanager.projectCreator), which contains the resourcemanager.projects.create permission. Learn how to grant roles.

    Go to project selector

  3. Verify that billing is enabled for your Google Cloud project.

  4. Enable the Managed Kafka API.

    Roles required to enable APIs

    To enable APIs, you need the Service Usage Admin IAM role (roles/serviceusage.serviceUsageAdmin), which contains the serviceusage.services.enable permission. Learn how to grant roles.

    Enable the API

  5. In the Google Cloud console, on the project selector page, select or create a Google Cloud project.

    Roles required to select or create a project

    • Select a project: Selecting a project doesn't require a specific IAM role—you can select any project that you've been granted a role on.
    • Create a project: To create a project, you need the Project Creator role (roles/resourcemanager.projectCreator), which contains the resourcemanager.projects.create permission. Learn how to grant roles.

    Go to project selector

  6. Verify that billing is enabled for your Google Cloud project.

  7. Enable the Managed Kafka API.

    Roles required to enable APIs

    To enable APIs, you need the Service Usage Admin IAM role (roles/serviceusage.serviceUsageAdmin), which contains the serviceusage.services.enable permission. Learn how to grant roles.

    Enable the API

  8. 确保您在项目上拥有以下角色: Managed Kafka Cluster EditorManaged Kafka Connect Cluster EditorManaged Kafka Connector EditorManaged Kafka Topic EditorBigQuery Data Owner

    检查角色

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

      转到 IAM
    2. 选择项目。
    3. 主账号列中,找到标识您或您所属群组的所有行。如需了解您属于哪些群组,请与您的管理员联系。

    4. 对于指定或包含您的所有行,请检查角色列以查看角色列表是否包含所需的角色。

    授予角色

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

      转到 IAM
    2. 选择项目。
    3. 点击 授予访问权限
    4. 新的主账号字段中,输入您的用户标识符。 这通常是 Google 账号的电子邮件地址。

    5. 点击选择角色,然后搜索相应角色。
    6. 如需授予其他角色,请点击 添加其他角色,然后添加其他各个角色。
    7. 点击 Save(保存)。

gcloud

  1. 登录您的 Google Cloud 账号。如果您是 Google Cloud新手,请 创建一个账号来评估我们的产品在实际场景中的表现。新客户还可获享 $300 赠金,用于运行、测试和部署工作负载。
  2. 安装 Google Cloud CLI。

  3. 如果您使用的是外部身份提供方 (IdP),则必须先使用联合身份登录 gcloud CLI

  4. 如需初始化 gcloud CLI,请运行以下命令:

    gcloud init
  5. 创建或选择 Google Cloud 项目

    选择或创建项目所需的角色

    • 选择项目:选择项目不需要特定的 IAM 角色,您可以选择已获授角色的任何项目。
    • 创建项目:如需创建项目,您需要拥有 Project Creator 角色 (roles/resourcemanager.projectCreator),该角色包含 resourcemanager.projects.create 权限。了解如何授予角色
    • 创建 Google Cloud 项目:

      gcloud projects create PROJECT_ID

      PROJECT_ID 替换为您要创建的 Google Cloud 项目的名称。

    • 选择您创建的 Google Cloud 项目:

      gcloud config set project PROJECT_ID

      PROJECT_ID 替换为您的 Google Cloud 项目名称。

  6. 验证是否已为您的 Google Cloud 项目启用结算功能

  7. 启用 Managed Kafka API:

    启用 API 所需的角色

    如需启用 API,您需要拥有 Service Usage Admin IAM 角色 (roles/serviceusage.serviceUsageAdmin),该角色包含 serviceusage.services.enable 权限。了解如何授予角色

    gcloud services enable managedkafka.googleapis.com
  8. 安装 Google Cloud CLI。

  9. 如果您使用的是外部身份提供方 (IdP),则必须先使用联合身份登录 gcloud CLI

  10. 如需初始化 gcloud CLI,请运行以下命令:

    gcloud init
  11. 创建或选择 Google Cloud 项目

    选择或创建项目所需的角色

    • 选择项目:选择项目不需要特定的 IAM 角色,您可以选择已获授角色的任何项目。
    • 创建项目:如需创建项目,您需要拥有 Project Creator 角色 (roles/resourcemanager.projectCreator),该角色包含 resourcemanager.projects.create 权限。了解如何授予角色
    • 创建 Google Cloud 项目:

      gcloud projects create PROJECT_ID

      PROJECT_ID 替换为您要创建的 Google Cloud 项目的名称。

    • 选择您创建的 Google Cloud 项目:

      gcloud config set project PROJECT_ID

      PROJECT_ID 替换为您的 Google Cloud 项目名称。

  12. 验证是否已为您的 Google Cloud 项目启用结算功能

  13. 启用 Managed Kafka API:

    启用 API 所需的角色

    如需启用 API,您需要拥有 Service Usage Admin IAM 角色 (roles/serviceusage.serviceUsageAdmin),该角色包含 serviceusage.services.enable 权限。了解如何授予角色

    gcloud services enable managedkafka.googleapis.com
  14. 向您的用户账号授予角色。对以下每个 IAM 角色运行以下命令一次: roles/managedkafka.clusterEditor, roles/managedkafka.connectClusterEditor, roles/managedkafka.connectorEditor, roles/managedkafka.topicEditor, roles/bigquery.dataOwner

    gcloud projects add-iam-policy-binding PROJECT_ID --member="user:USER_IDENTIFIER" --role=ROLE

    替换以下内容:

    • PROJECT_ID:您的项目 ID。
    • USER_IDENTIFIER:您的用户 账号的标识符。例如,myemail@example.com
    • ROLE:您授予用户账号的 IAM 角色。

创建 BigQuery 表

在此步骤中,您将创建一个具有以下架构的 BigQuery 表:

列名 数据类型
name STRING
id INTEGER

创建数据集

如需创建 BigQuery 数据集,请按以下步骤操作:

控制台

  1. 打开 BigQuery 页面。

    转到 BigQuery 页面

  2. 探索器面板中,选择您要在其中创建数据集的项目。

  3. 展开 查看操作选项,然后点击创建数据集

  4. 创建数据集页面中执行以下操作:

    • 数据集 ID 中,输入数据集的名称。

    • 位置类型字段中,为数据集选择一个地理位置。

gcloud

如需创建新数据集,请使用带有 --dataset 标志的 bq mk 命令。

bq mk --location REGION \
  --dataset PROJECT_ID:DATASET_NAME

替换以下内容:

  • PROJECT_ID:您的项目 ID
  • DATASET_NAME:数据集的名称
  • REGION:数据集的位置

如需了解详情,请参阅创建数据集

创建具有架构的表

接下来,创建一个具有架构的新 BigQuery 表:

控制台

  1. 转到 BigQuery 页面。

    转到 BigQuery 页面

  2. 探索器面板中,展开您的项目,然后选择数据集。

  3. 数据集信息部分中,点击 创建表

  4. 基于以下数据源创建表列表中,选择空表

  5. 框中,输入表的名称。

  6. 架构部分中,点击以文本形式修改

  7. 粘贴以下架构定义:

    name:STRING,
    id:INTEGER
    
  8. 点击创建表

gcloud

如需创建新表,请使用带有 --table 标志的 bq mk 命令。

bq mk --table \
  PROJECT_ID:DATASET_NAME.TABLE_NAME \
  name:STRING,id:INTEGER

替换以下内容:

  • PROJECT_ID:您的项目 ID
  • DATASET_NAME:数据集的名称
  • TABLE_NAME:要创建的表的名称。

如需了解详情,请参阅使用架构定义创建空表

默认情况下,BigQuery sink 连接器使用主题名称作为 BigQuery 表名称。您可以通过设置 topic2TableMap 配置属性来替换此行为。如需了解详情,请参阅 BigQuery Sink 连接器的工作原理

创建 Managed Service for Apache Kafka 资源

在本部分中,您将创建以下 Managed Service for Apache Kafka 资源:

  • 包含主题的 Kafka 集群。
  • 具有 BigQuery 接收器连接器的 Connect 集群。

创建 Kafka 集群

在此步骤中,您将创建一个 Managed Service for Apache Kafka 集群。创建集群最多可能需要 30 分钟。

控制台

  1. 前往 Managed Service for Apache Kafka > 集群页面。

    转到“集群”

  2. 点击 创建

  3. 集群名称字段中,输入集群的名称。

  4. 区域列表中,为集群选择一个位置。选择与 BigQuery 表相同的区域。

  5. 对于网络配置,请配置集群可访问的子网:

    1. 项目部分,选择您的项目。
    2. 对于网络,选择 VPC 网络。
    3. 子网字段中,选择相应子网。
    4. 点击完成
  6. 点击创建

在创建集群期间,集群状态为 Creating。集群创建完成后,状态为 Active

gcloud

如需创建 Kafka 集群,请运行 managed-kafka clusters create 命令。

gcloud managed-kafka clusters create KAFKA_CLUSTER \
--location=REGION \
--cpu=3 \
--memory=3GiB \
--subnets=projects/PROJECT_ID/regions/REGION/subnetworks/SUBNET_NAME \
--async

替换以下内容:

  • KAFKA_CLUSTER:Kafka 集群的名称
  • REGION:集群的位置;选择与 BigQuery 表相同的区域
  • PROJECT_ID:您的项目 ID
  • SUBNET_NAME:您要在其中创建集群的子网,例如 default

如需了解支持的位置,请参阅 Managed Service for Apache Kafka 位置

该命令会异步运行,并返回一个操作 ID:

Check operation [projects/PROJECT_ID/locations/REGION/operations/OPERATION_ID] for status.

如需跟踪创建操作的进度,请使用 gcloud managed-kafka operations describe 命令:

gcloud managed-kafka operations describe OPERATION_ID \
  --location=REGION

如需了解详情,请参阅监控集群创建操作

创建 Kafka 主题

创建 Managed Service for Apache Kafka 集群后,创建一个 Kafka 主题。为该主题指定与 BigQuery 表相同的名称。

控制台

  1. 前往 Managed Service for Apache Kafka > 集群页面。

    转到“集群”

  2. 点击集群的名称。

  3. 在集群详细信息页面中,点击 创建主题

  4. 主题名称框中,输入与您的 BigQuery 表相同的名称。

  5. 点击创建

gcloud

如需创建主题,请运行 managed-kafka topics create 命令。

gcloud managed-kafka topics create TABLE_NAME \
--cluster=KAFKA_CLUSTER \
--location=REGION \
--partitions=10 \
--replication-factor=3

替换以下内容:

  • TABLE_NAME:BigQuery 表的名称,也是主题名称
  • KAFKA_CLUSTER:Kafka 集群的名称
  • REGION:您创建 Kafka 集群的区域

创建 Connect 集群

在此步骤中,您将创建一个 Connect 集群。创建 Connect 集群最多可能需要 30 分钟。

在开始此步骤之前,请确保 Managed Service for Apache Kafka 集群已完全创建。

控制台

  1. 前往 Managed Service for Apache Kafka > Connect 集群页面。

    前往“关联集群”

  2. 点击 创建

  3. Connect 集群名称中,输入一个字符串。示例:my-connect-cluster

  4. 对于主 Kafka 集群,选择您之前创建的 Kafka。

  5. 点击创建

在创建集群期间,集群状态为 Creating。集群创建完成后,状态为 Active

gcloud

如需创建 Connect 集群,请运行 gcloud managed-kafka connect-clusters create 命令。

gcloud managed-kafka connect-clusters create CONNECT_CLUSTER \
  --location=REGION \
  --cpu=12 \
  --memory=12GiB \
  --primary-subnet=projects/PROJECT_ID/regions/REGION/subnetworks/SUBNET_NAME \
  --kafka-cluster=KAFKA_CLUSTER \
  --async

替换以下内容:

  • CONNECT_CLUSTER:Connect 集群的名称
  • REGION:您创建 Kafka 集群的区域
  • PROJECT_ID:您的项目 ID
  • SUBNET_NAME:您在其中创建 Kafka 集群的子网
  • KAFKA_CLUSTER:Kafka 集群的名称

该命令会异步运行,并返回一个操作 ID:

Check operation [projects/PROJECT_ID/locations/REGION/operations/OPERATION_ID] for status.

如需跟踪创建操作的进度,请使用 gcloud managed-kafka operations describe 命令:

gcloud managed-kafka operations describe OPERATION_ID \
  --location=REGION

如需了解详情,请参阅监控集群创建操作

授予 IAM 角色

向 Managed Kafka 服务账号授予 BigQuery Data Editor Identity and Access Management (IAM) 角色。此角色允许连接器写入 BigQuery 表。

控制台

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

    转到 IAM

  2. 选择包括 Google 提供的角色授权

  3. 找到 Managed Kafka 服务账号对应的行,然后点击 修改主账号

  4. 点击添加其他角色,然后选择 BigQuery Data Editor 角色。

  5. 点击保存

如需详细了解如何授予角色,请参阅使用控制台授予 IAM 角色

gcloud

如需向服务账号授予 IAM 角色,请运行 gcloud projects add-iam-policy-binding 命令。

gcloud projects add-iam-policy-binding PROJECT_ID \
    --member=serviceAccount:service-PROJECT_NUMBER@gcp-sa-managedkafka.iam.gserviceaccount.com \
    --role=roles/bigquery.dataEditor

替换以下内容:

  • PROJECT_ID:您的项目 ID
  • PROJECT_NUMBER:您的项目编号

如需查找您的项目编号,请使用 gcloud projects describe 命令。

创建 BigQuery 接收器连接器

在此步骤中,您将创建一个 BigQuery Sink 连接器。此连接器从一个或多个主题读取消息,并将其写入 BigQuery。

控制台

  1. 前往 Managed Service for Apache Kafka > Connect 集群页面。

    前往“关联集群”

  2. 点击 Connect 集群的名称。

  3. 点击 Create connector

  4. 连接器名称中,输入一个字符串。示例:bigquery-connector

  5. 连接器插件列表中,选择 BigQuery Sink

  6. 对于主题,选择您之前创建的主题,然后点击确定

  7. 对于数据集,请输入 BigQuery 数据集的名称,格式如下: PROJECT_ID.DATASET_NAME。 示例:my-project.dataset1

  8. 配置编辑框中,添加以下行:

    bigQueryPartitionDecorator=false
    
  9. 点击创建

gcloud

如需创建 BigQuery Sink 连接器,请运行 gcloud managed-kafka connectors create 命令。

gcloud managed-kafka connectors create CONNECTOR_NAME \
  --location=REGION \
  --connect-cluster=CONNECT_CLUSTER \
  --configs=connector.class=com.wepay.kafka.connect.bigquery.BigQuerySinkConnector,\
key.converter=org.apache.kafka.connect.storage.StringConverter,\
value.converter=org.apache.kafka.connect.json.JsonConverter,\
value.converter.schemas.enable=false,\
tasks.max=3,\
project=PROJECT_ID,\
defaultDataset=DATASET_NAME,\
topics=TABLE_NAME,\
bigQueryPartitionDecorator=false

替换以下内容:

  • CONNECTOR_NAME:连接器的名称,例如 bigquery-connector
  • CONNECT_CLUSTER:您的 Connect 集群的名称
  • REGION:您在其中创建 Connect 集群的区域
  • PROJECT_ID:您的项目 ID
  • DATASET_NAME:BigQuery 数据集的名称。
  • TABLE_NAME:Kafka 主题的名称,与 BigQuery 表的名称相同。

bigQueryPartitionDecorator 配置参数设置为 false 可防止连接器向表名称添加分区修饰器(例如 "$"yyyyMMdd")。

查看结果

如需查看结果,请向 Kafka 主题发送消息。请使用以下格式设置消息正文:

{ "name": "STRING_VALUE", "id": INTEGER_VALUE }

您可以通过多种方式向 Managed Service for Apache Kafka 发送消息,包括:

如需在 BigQuery 中查看记录,请对该表运行如下查询:

控制台

  1. 打开 BigQuery 页面。

    转到 BigQuery 页面

  2. 在查询编辑器中,运行以下查询:

    SELECT * FROM `PROJECT_ID.DATASET_NAME.TABLE_NAME`
    LIMIT 1000
    

    执行以下变量替换操作:

    • PROJECT_ID:您的 Google Cloud项目的名称
    • DATASET_NAME:BigQuery 数据集的名称
    • TABLE_NAME:BigQuery 表的名称

gcloud

使用 bq query 命令对表运行查询:

bq query --use_legacy_sql=false 'SELECT * FROM `PROJECT_ID.DATASET_NAME.TABLE_NAME`'

执行以下变量替换操作:

  • PROJECT_ID:您的 Google Cloud项目的名称
  • DATASET_NAME:BigQuery 数据集的名称
  • TABLE_NAME:BigQuery 表的名称

清理

为避免因本教程中使用的资源导致您的 Google Cloud 账号产生费用,请删除包含这些资源的项目,或者保留项目但删除各个资源。

控制台

  1. 删除 Connect 集群。

    1. 前往 Managed Service for Apache Kafka > Connect 集群页面。

      前往“关联集群”

    2. 选择 Connect 集群,然后点击删除

  2. 删除 Kafka 集群。

    1. 前往 Managed Service for Apache Kafka > 集群页面。

      转到“集群”

    2. 选择 Kafka 集群,然后点击删除

  3. 删除 BigQuery 表和数据集。

    1. 转到 BigQuery 页面。

      转到 BigQuery 页面

    2. 浏览器窗格中,展开您的项目,然后选择数据集。

    3. 展开 操作选项,然后点击删除

    4. 删除数据集对话框中,在字段中输入 delete,然后点击删除

gcloud

  1. 如需删除 Connect 集群,请使用 gcloud managed-kafka connect-clusters delete 命令。

    gcloud managed-kafka connect-clusters delete CONNECT_CLUSTER \
      --location=REGION --async
    
  2. 如需删除 Kafka 集群,请使用 gcloud managed-kafka clusters delete 命令。

    gcloud managed-kafka clusters delete KAFKA_CLUSTER \
      --location=REGION --async
    
  3. 如需同时删除 BigQuery 数据集和 BigQuery 表,请使用 bq rm 命令。

    bq rm --recursive --dataset PROJECT_ID:DATASET_NAME
    

后续步骤