本教程介绍如何使用 Kafka Connect 将 Pub/Sub 中的消息注入到 Managed Service for Apache Kafka 集群中。
Kafka Connect 可管理 Kafka 集群与其他系统之间的数据移动。在本教程中,您将创建一个 Connect 集群和一个 Pub/Sub 来源连接器。Pub/Sub 源连接器会从 Pub/Sub 主题中读取消息并将其写入 Kafka 主题。
准备工作
控制台
- 登录您的 Google Cloud 账号。如果您是 Google Cloud新手,请 创建一个账号来评估我们的产品在实际场景中的表现。新客户还可获享 $300 赠金,用于运行、测试和部署工作负载。
-
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 theresourcemanager.projects.createpermission. Learn how to grant roles.
-
Verify that billing is enabled for your Google Cloud project.
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 theserviceusage.services.enablepermission. Learn how to grant roles.-
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 theresourcemanager.projects.createpermission. Learn how to grant roles.
-
Verify that billing is enabled for your Google Cloud project.
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 theserviceusage.services.enablepermission. Learn how to grant roles.-
确保您拥有项目的以下角色: Managed Kafka Cluster Editor、 Managed Kafka Connect Cluster Editor、 Managed Kafka Connector Editor、 Managed Kafka Topic Editor、 Pub/Sub Editor
检查角色
-
在 Google Cloud 控制台中,前往 IAM 页面。
转到 IAM - 选择项目。
-
在主账号列中,找到标识您或您所属群组的所有行。如需了解您属于哪些群组,请与您的管理员联系。
- 对于指定或包含您的所有行,请检查角色列以查看角色列表是否包含所需的角色。
授予角色
-
在 Google Cloud 控制台中,前往 IAM 页面。
转到 IAM - 选择项目。
- 点击 授予访问权限。
-
在新的主账号字段中,输入您的用户标识符。 这通常是 Google 账号的电子邮件地址。
- 点击选择角色,然后搜索相应角色。
- 如需授予其他角色,请点击 添加其他角色,然后添加其他各个角色。
- 点击 Save(保存)。
-
gcloud
- 登录您的 Google Cloud 账号。如果您是 Google Cloud新手,请 创建一个账号来评估我们的产品在实际场景中的表现。新客户还可获享 $300 赠金,用于运行、测试和部署工作负载。
-
安装 Google Cloud CLI。
-
如果您使用的是外部身份提供方 (IdP),则必须先使用联合身份登录 gcloud CLI。
-
如需初始化 gcloud CLI,请运行以下命令:
gcloud init -
选择或创建项目所需的角色
- 选择项目:选择项目不需要特定的 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 项目名称。
启用 Managed Kafka API:
启用 API 所需的角色
如需启用 API,您需要拥有 Service Usage Admin IAM 角色 (
roles/serviceusage.serviceUsageAdmin),该角色包含serviceusage.services.enable权限。了解如何授予角色。gcloud services enable managedkafka.googleapis.com
-
安装 Google Cloud CLI。
-
如果您使用的是外部身份提供方 (IdP),则必须先使用联合身份登录 gcloud CLI。
-
如需初始化 gcloud CLI,请运行以下命令:
gcloud init -
选择或创建项目所需的角色
- 选择项目:选择项目不需要特定的 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 项目名称。
启用 Managed Kafka API:
启用 API 所需的角色
如需启用 API,您需要拥有 Service Usage Admin IAM 角色 (
roles/serviceusage.serviceUsageAdmin),该角色包含serviceusage.services.enable权限。了解如何授予角色。gcloud services enable managedkafka.googleapis.com
-
向您的用户账号授予角色。对以下每个 IAM 角色运行以下命令一次:
roles/managedkafka.clusterEditor, roles/managedkafka.connectClusterEditor, roles/managedkafka.connectorEditor, roles/managedkafka.topicEditor, roles/pubsub.editorgcloud projects add-iam-policy-binding PROJECT_ID --member="user:USER_IDENTIFIER" --role=ROLE
替换以下内容:
PROJECT_ID:您的项目 ID。USER_IDENTIFIER:您的用户 账号的标识符。例如,myemail@example.com。ROLE:您授予用户账号的 IAM 角色。
创建 Pub/Sub 主题和订阅
在此步骤中,您将创建一个具有订阅的 Pub/Sub 主题。
控制台
前往 Pub/Sub > 主题页面。
点击 创建主题。
在主题 ID 框中,输入主题的名称。
确保选中添加默认订阅复选框。
点击创建。
gcloud
如需创建 Pub/Sub 主题,请运行
gcloud pubsub topics create命令。gcloud pubsub topics create TOPIC_ID将
TOPIC_ID替换为您的 Pub/Sub 主题的名称。如需创建对主题的订阅,请运行
gcloud pubsub subscriptions create命令:gcloud pubsub subscriptions create --topic TOPIC_ID SUBSCRIPTION_ID将
SUBSCRIPTION_ID替换为您的 Pub/Sub 订阅的名称。
如需了解如何命名 Pub/Sub 主题和订阅,请参阅主题或订阅命名指南。
创建 Managed Service for Apache Kafka 资源
在本部分中,您将创建以下 Managed Service for Apache Kafka 资源:
- 包含主题的 Kafka 集群。
- 具有 Pub/Sub 连接器的 Connect 集群。
创建 Kafka 集群
在此步骤中,您将创建一个 Managed Service for Apache Kafka 集群。创建集群最多可能需要 30 分钟。
控制台
- 前往 Managed Service for Apache Kafka > 集群页面。
- 点击 创建。
- 在集群名称字段中,输入集群的名称。
- 在区域列表中,为集群选择一个位置。
-
对于网络配置,请配置集群可访问的子网:
- 在项目部分,选择您的项目。
- 对于网络,选择 VPC 网络。
- 在子网字段中,选择相应子网。
- 点击完成。
- 点击创建。
点击创建后,集群状态为 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:集群的位置PROJECT_ID:您的项目 IDSUBNET_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
当集群准备就绪时,此命令的输出内容会包含 state:
ACTIVE 条目。如需了解详情,请参阅 监控集群创建操作。
创建 Kafka 主题
创建 Managed Service for Apache Kafka 集群后,创建 Kafka 主题。
控制台
前往 Managed Service for Apache Kafka > 集群页面。
点击集群的名称。
在集群详细信息页面中,点击 创建主题。
在主题名称框中,输入主题的名称。
点击创建。
gcloud
如需创建 Kafka 主题,请运行 managed-kafka topics create 命令。
gcloud managed-kafka topics create KAFKA_TOPIC_NAME \
--cluster=KAFKA_CLUSTER \
--location=REGION \
--partitions=10 \
--replication-factor=3
替换以下内容:
KAFKA_TOPIC_NAME:要创建的 Kafka 主题的名称KAFKA_CLUSTER:Kafka 集群的名称REGION:您创建 Kafka 集群的区域
创建 Connect 集群
在此步骤中,您将创建一个 Connect 集群。创建 Connect 集群最多可能需要 30 分钟。
在开始此步骤之前,请确保 Managed Service for Apache Kafka 集群已完全创建。
控制台
前往 Managed Service for Apache Kafka > Connect 集群页面。
点击 创建。
在 Connect 集群名称中,输入一个字符串。示例:
my-connect-cluster。对于主 Kafka 集群,选择您之前创建的 Kafka。
点击创建。
在创建集群期间,集群状态为 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:您的项目 IDSUBNET_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 服务账号授予以下 Identity and Access Management (IAM) 角色:
- Pub/Sub Subscriber
- Pub/Sub Viewer
这些角色允许连接器从 Pub/Sub 读取消息。
控制台
在 Google Cloud 控制台中,前往 IAM 页面。
选择包括 Google 提供的角色授权。
找到 Managed Kafka 服务账号对应的行,然后点击 修改主账号。
点击添加其他角色,然后选择角色 Pub/Sub Subscriber。 对 Pub/Sub 查看者角色重复此步骤。
点击保存。
如需详细了解如何授予角色,请参阅使用控制台授予 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/pubsub.subscriber
gcloud projects add-iam-policy-binding PROJECT_ID \
--member=serviceAccount:service-PROJECT_NUMBER@gcp-sa-managedkafka.iam.gserviceaccount.com \
--role=roles/pubsub.viewer
替换以下内容:
PROJECT_ID:您的项目 IDPROJECT_NUMBER:您的项目编号
如需查找您的项目编号,请使用 gcloud projects describe 命令。
创建 Pub/Sub 源连接器
在此步骤中,您将创建 Pub/Sub Source 连接器。此连接器可从 Pub/Sub 读取消息并将其写入 Kafka 主题。
控制台
前往 Managed Service for Apache Kafka > Connect 集群页面。
点击 Connect 集群的名称。
点击 Create connector。
在连接器名称中,输入一个字符串。示例:
pubsub-source。在连接器插件列表中,选择
Pub/Sub Source。对于 Cloud Pub/Sub 订阅,请选择在创建 Pub/Sub 主题时创建的默认 Pub/Sub。
对于 Kafka 主题,请选择您之前创建的 Kafka 主题。
点击创建。
gcloud
如需创建 Pub/Sub 来源连接器,请运行 gcloud managed-kafka connectors create 命令。
gcloud managed-kafka connectors create PUBSUB_CONNECTOR_NAME \
--connect-cluster=CONNECT_CLUSTER \
--location=REGION \
--configs=connector.class=com.google.pubsub.kafka.source.CloudPubSubSourceConnector,\
cps.project=PROJECT_ID,\
cps.streamingPull.enabled=true,\
cps.subscription=SUBSCRIPTION_ID,\
kafka.topic=KAFKA_TOPIC_NAME,\
key.converter=org.apache.kafka.connect.storage.StringConverter,\
tasks.max=3,\
value.converter=org.apache.kafka.connect.converters.ByteArrayConverter
替换以下内容:
PUBSUB_CONNECTOR_NAME:连接器的名称,例如pubsub-source-connectorCONNECT_CLUSTER:您的 Connect 集群的名称REGION:您在其中创建 Connect 集群的区域PROJECT_ID:您的项目 IDKAFKA_TOPIC_NAME:Kafka 主题的名称SUBSCRIPTION_ID:您的 Pub/Sub 订阅的名称
查看结果
如需查看结果,请将一些消息发布到 Pub/Sub。
控制台
在 Google Cloud 控制台中,前往 Pub/Sub > 主题页面。
在主题列表中,点击您的 Pub/Sub 主题的名称。
点击消息。
点击发布消息。
在消息数量部分,输入
10。对于消息正文,输入
{"name": "Alice", "customer_id": 1}。点击发布。
gcloud
如需将消息发布到 Pub/Sub 主题,请使用 gcloud pubsub topics publish 命令。
for run in {1..10}; do
gcloud pubsub topics publish TOPIC_ID --message='{"name": "Alice", "customer_id": 1}'
done
将 TOPIC_ID 替换为您的 Pub/Sub 主题的名称。
现在,您可以从 Kafka 主题中使用消息了。如需了解详情,请参阅使用 CLI 生成和使用消息。
清理
为避免因本教程中使用的资源导致您的 Google Cloud 账号产生费用,请删除包含这些资源的项目,或者保留项目但删除各个资源。
控制台
删除 Pub/Sub 主题。
前往 Pub/Sub > 主题页面。
选择相应主题,然后点击删除。
删除 Pub/Sub 订阅。
前往 Pub/Sub > 订阅页面。
选择使用您的主题创建的订阅,然后点击删除。
删除 Connect 集群。
前往 Managed Service for Apache Kafka > Connect 集群页面。
选择 Connect 集群,然后点击删除。
删除 Kafka 集群。
前往 Managed Service for Apache Kafka > 集群页面。
选择 Kafka 集群,然后点击删除。
gcloud
如需删除 Pub/Sub 订阅和主题,请使用
gcloud pubsub subscriptions delete和gcloud pubsub topics delete命令。gcloud pubsub subscriptions delete SUBSCRIPTION_ID gcloud pubsub topics delete TOPIC_ID如需删除 Connect 集群,请使用
gcloud managed-kafka connect-clusters delete命令。gcloud managed-kafka connect-clusters delete CONNECT_CLUSTER \ --location=REGION --async如需删除 Kafka 集群,请使用
gcloud managed-kafka clusters delete命令。gcloud managed-kafka clusters delete KAFKA_CLUSTER \ --location=REGION --async
后续步骤
- 排查 Pub/Sub 连接器问题。
- 详细了解 Pub/Sub 来源连接器。
- 详细了解 Kafka Connect。