将 Kafka Connect 与架构注册表搭配使用

本教程介绍如何将 Kafka Connect架构注册表集成, 以便连接器可以使用注册的架构来序列化和反序列化 消息。

在开始学习本教程之前,请完成 使用架构注册表生成 Avro 消息中的步骤。完成这些步骤后,您将拥有本教程所需的以下组件:

  • 一个架构注册表,其中包含已注册的 Avro 架构。
  • 一个 Kafka 集群,其中包含一个包含 Avro 消息的主题。

概览

架构定义了消息或数据载荷的结构。架构注册表提供了一个集中位置来存储架构。通过使用架构注册表,您可以确保数据从生产者到消费者的编码和解码保持一致。

本教程演示了以下场景:

  1. Java 生产者客户端将 Apache Avro 消息写入 Managed Service for Apache Kafka 集群。Avro 消息架构存储在架构注册表中。

  2. BigQuery 接收器 连接器从 Kafka 读取传入消息,并将其写入 BigQuery。该连接器会自动创建一个与 Avro 架构兼容的 BigQuery 表。

准备工作

  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 Editor、 和 BigQuery Data Owner

    检查角色

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

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

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

    授予角色

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

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

    5. 点击选择角色,然后搜索该角色。
    6. 如需授予其他角色,请点击 添加其他角色 ,然后添加其他各个角色。
    7. 点击 Save (保存)。
  9. 完成使用架构注册表生成 Avro 消息中的步骤。

创建 Connect 集群

请按照以下步骤创建 Connect 集群。创建 Connect 集群最多可能需要 30 分钟。

控制台

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

    前往 Connect 集群

  2. 点击 创建

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

  4. 对于 Kafka 主集群 ,选择 Kafka 集群。

  5. 点击创建

创建 Connect 集群时,集群状态为 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 Service Account 对应的行,然后点击 修改主账号

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

  5. 点击 Save (保存)。

如需详细了解如何授予角色,请参阅 使用控制台授予 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 命令。

默认情况下,Connect 集群有权访问同一项目中的架构注册表。如果您在与架构注册表不同的项目中创建 Connect 集群,则必须向 Managed Kafka 服务代理授予架构注册表项目中的 Managed Kafka Client (roles/managedkafka.client) 角色。如需了解详情,请参阅 在不同项目中创建 Connect 集群

创建 BigQuery 数据集

在此步骤中,您将创建一个 数据集 来保存 BigQuery 表。BigQuery 接收器连接器会自动创建该表。

如需创建数据集,请执行以下步骤。

控制台

  1. 打开 BigQuery 页面。

    转到 BigQuery 页面

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

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

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

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

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

gcloud

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

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

替换以下内容:

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

创建 BigQuery 接收器连接器

在此步骤中,您将创建一个 BigQuery 接收器连接器, 该连接器从 Kafka 读取数据并写入 BigQuery。您将连接器配置为使用存储在架构注册表中的 Avro 架构反序列化 Kafka 消息。

控制台

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

    前往 Connect 集群

  2. 点击 Connect 集群的名称。

  3. 点击 创建连接器

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

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

  6. 对于主题 ,选择名为 newUsers 的 Kafka 主题。此主题由 Java 生产者客户端 创建。

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

  8. 配置 编辑框中,将现有配置替换为以下配置:

    tasks.max=3
    key.converter=org.apache.kafka.connect.storage.StringConverter
    value.converter=io.confluent.connect.avro.AvroConverter
    value.converter.bearer.auth.credentials.source=GCP
    value.converter.schema.registry.url=https://managedkafka.googleapis.com/v1/projects/PROJECT_ID/locations/REGISTRY_REGION/schemaRegistries/REGISTRY_NAME
    

    替换以下内容:

    • PROJECT_ID:您的项目 ID
    • REGISTRY_REGION:您在其中创建架构注册表的区域
    • REGISTRY_NAME:架构注册表的名称
  9. 点击创建

gcloud

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


export REGISTRY_PATH=projects/PROJECT_ID/locations/REGISTRY_REGION/schemaRegistries/REGISTRY_NAME

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=io.confluent.connect.avro.AvroConverter,\
value.converter.bearer.auth.credentials.source=GCP,\
value.converter.schema.registry.url=https://managedkafka.googleapis.com/v1/$REGISTRY_PATH,\
tasks.max=3,\
project=PROJECT_ID,\
defaultDataset=DATASET_NAME,\
topics=newUsers

替换以下内容:

  • REGISTRY_REGION:您在其中创建架构注册表的区域
  • REGISTRY_NAME:架构注册表的名称
  • CONNECTOR_NAME:连接器的名称,例如 bigquery-connector
  • CONNECT_CLUSTER:Connect 集群的名称
  • REGION:您在其中创建 Connect 集群的区域
  • PROJECT_ID:您的项目 ID
  • DATASET_NAME:BigQuery 数据集的名称

查看结果

如需在 BigQuery 中查看结果,请按如下方式对表运行查询。将第一批消息写入表可能需要几分钟时间。

控制台

  1. 打开 BigQuery 页面。

    转到 BigQuery 页面

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

    SELECT * FROM `PROJECT_ID.DATASET_NAME.newUsers`
    

    执行以下变量替换操作:

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

gcloud

使用 bq query command 对表运行查询:

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

执行以下变量替换操作:

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

清理

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

控制台

  1. 删除 Connect 集群。

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

      前往 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
    

后续步骤