为 Managed Service for Apache Kafka 集群生成合成数据
了解如何使用 Google Cloud 控制台为 Google Cloud Managed Service for Apache Kafka 集群生成合成测试数据。
本指南使用 Dataflow 流式数据生成器模板自动将示例游戏遥测数据发布到 Managed Service for Apache Kafka 主题。流式数据生成器是一个 Dataflow 模板,可根据指定架构以可配置的速率生成合成测试记录。生成合成数据后,您无需安装本地 Kafka 客户端或编写自定义生产者代码,即可观察集群活动、测试负载处理并验证监控指标。如需详细了解该模板,请参阅 Dataflow 流式数据生成器模板。
准备工作
在开始本教程之前,请创建一个新的 Managed Service for Apache Kafka 集群。如果您已有集群,则可以跳过此步骤。如需了解创建集群所需的角色和权限,请参阅创建和查看集群。如果您按照该指南操作,请仅完成创建集群 部分,然后再返回本指南。
如何创建集群
控制台
- 前往 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。如需了解详情,请参阅
监控
集群创建操作。
所需的角色
如需获得为集群生成合成数据所需的权限,请让管理员向您授予项目的以下 IAM 角色:
- Dataflow Developer (
roles/dataflow.developer) - Project IAM Admin (
roles/resourcemanager.projectIamAdmin)
如需详细了解如何授予角色,请参阅管理对项目、文件夹和组织的访问权限。
这些预定义角色包含 为集群生成合成数据所需的权限。如需查看所需的确切权限,请展开所需权限部分:
所需权限
为集群生成合成数据需要以下权限:
-
dataflow.jobs.create -
dataflow.jobs.get -
managedkafka.clusters.get -
managedkafka.topics.get -
managedkafka.topics.create -
managedkafka.topics.publish -
resourcemanager.projects.setIamPolicy
为确保 Compute Engine 默认服务账号 具有运行 Dataflow 作业所需的权限,请让管理员向 Compute Engine 默认服务账号 授予项目的以下 IAM 角色:
如果未能向正确的主账号授予这些角色,可能会导致权限错误。- Dataflow Worker (
roles/dataflow.worker) - Managed Kafka Client (
roles/managedkafka.client)
如需详细了解如何授予角色,请参阅管理对项目、文件夹和组织的访问权限。
您的管理员也可以通过自定义角色或其他预定义角色为Compute Engine 默认服务账号授予所需的权限。
如果您拥有授予 IAM 角色的权限,控制台会在合成数据生成过程中提示您授予所需的角色。Google Cloud 如果您没有授予角色的权限,控制台会显示一条消息,要求您请求管理员授予所需的权限。
生成合成数据
如需创建并启动为 Kafka 主题生成合成数据的 Dataflow 作业,请执行以下操作:
在 Google Cloud 控制台中,前往 Managed Service for Apache Kafka > 集群 页面。
点击集群的名称,例如
test-cluster。选择来源 标签页。
在来源 页面上的生成合成数据 卡片中,点击创建 Dataflow 作业 。随即会打开生成数据 窗格。
在生成数据 窗格中,从 Kafka 主题 下拉列表中选择一个主题,例如
test-topic。如果您没有主题,请创建一个:- 在 Kafka 主题 下拉列表中,点击创建主题 。随即会打开创建主题 窗格。
- 在主题名称 字段中,输入
test-topic。 - 保留分区数 (
3) 和复制因子 (3) 的默认值。 - 点击创建 。
在输出速率 (QPS) 字段中,输入您希望生成器生成的每秒查询次数速率,例如
100。这样,您就可以测试集群如何处理不同的负载。如果出现一条警告,指出您的 Dataflow 服务账号缺少必要的权限,请点击授予 以分配以下角色:
- Dataflow Worker (
roles/dataflow.worker) - Managed Kafka Client (
roles/managedkafka.client)
- Dataflow Worker (
在生成数据 窗格中,点击创建 以启动 Dataflow 作业。
系统会显示一条通知,指出 Dataflow 作业已创建。
在通知中,点击查看作业 以打开 Dataflow 作业详情 页面,您可以在其中观察作业图、状态和执行指标。
查看集群指标
Dataflow 作业启动后,观察流入集群的合成数据:
在
test-cluster的集群详情 页面上,点击监控 标签页。查看字节速率 和按生成吞吐量排名的前 5 个主题 图表,以验证是否正在向您的主题积极生成数据。
查看消息
使用以下方法之一验证是否正在向您的主题发布合成消息。
在 Kafka 命令行工具中查看
如需在客户端虚拟机上使用 Kafka CLI 工具直接使用集群中的消息,请执行以下操作:
使用 SSH 连接到客户端虚拟机。如果您尚未设置客户端虚拟机,请参阅 创建客户端虚拟机。
从 Google Cloud 控制台 中获取集群的引导服务器地址,并将其设置为客户端虚拟机上的环境变量:
在 Google Cloud 控制台中,前往 Managed Service for Apache Kafka > 集群 页面。
点击集群的名称,例如
test-cluster。在集群详情 页面上,点击配置 。
复制引导网址 下列出的值。
在客户端虚拟机上,设置环境变量:
```sh export BOOTSTRAP="BOOTSTRAP_URL" ```将
BOOTSTRAP_URL替换为您复制的引导地址。运行
kafka-console-consumer.sh命令以读取消息:kafka-console-consumer.sh \ --bootstrap-server $BOOTSTRAP \ --topic TOPIC_ID \ --from-beginning \ --consumer.config client.properties将 TOPIC_ID 替换为主题名称,例如
test-topic。控制台会显示流式传输的合成游戏数据记录,因为它们会被使用。
按 Ctrl+C 停止使用消息。
在 BigQuery 中查看
如需将数据从 Kafka 主题流式传输到 BigQuery 并查看记录,请执行以下操作:
由于合成数据是原始 JSON,因此您必须先在 BigQuery 中手动创建目标表,然后再创建连接器。如需了解如何创建表,请参阅使用架构定义创建空表。在数据集中创建一个名为
test-topic的表,并使用以下架构:[ {"name": "eventId", "type": "STRING"}, {"name": "eventTimestamp", "type": "INTEGER"}, {"name": "ipv4", "type": "STRING"}, {"name": "ipv6", "type": "STRING"}, {"name": "country", "type": "STRING"}, {"name": "username", "type": "STRING"}, {"name": "quest", "type": "STRING"}, {"name": "score", "type": "INTEGER"}, {"name": "completed", "type": "BOOLEAN"} ]在 Connect 集群中创建一个 BigQuery 接收器连接器 ,以将消息从主题流式传输到 BigQuery 表。配置连接器时,请使用以下示例属性,并将
PROJECT_ID替换为您的项目 ID:bigQueryPartitionDecorator=false connector.class=com.wepay.kafka.connect.bigquery.BigQuerySinkConnector defaultDataset=test_dataset key.converter=org.apache.kafka.connect.storage.StringConverter project=PROJECT_ID tasks.max=3 topics=test-topic value.converter=org.apache.kafka.connect.json.JsonConverter value.converter.schemas.enable=false连接器开始流式传输数据后,前往BigQuery页面 中的 Google Cloud 控制台。
在浏览器 面板中,展开您的项目 ID 并选择您的数据集
test_dataset。点击表名称
test-topic。点击预览 标签页以查看流式传输的综合记录。或者,点击编写新查询 并运行以下 SQL 查询:
SELECT * FROM `PROJECT_ID.DATASET_ID.TABLE_ID` LIMIT 10;替换以下内容:
- PROJECT_ID:您的项目 ID
- DATASET_ID:您的数据集 ID,例如
test_dataset - TABLE_ID:您的表 ID,例如
test-topic
点击运行 以在查询结果 窗格中查看示例记录。
注意: 请勿使用
SELECT COUNT(*)查询来验证记录。 由于连接器使用 BigQuery 流式插入 API,因此数据最初会写入流式缓冲区。虽然使用SELECT *可以立即看到数据,但行数可能需要几分钟才能更新。
清理
为避免因本页中使用的资源导致您的 Google Cloud 账号产生费用,请按照以下步骤操作。
在 Google Cloud 控制台中,前往 Dataflow 作业 页面。
点击为您的主题创建的作业的名称。
点击停止 。
选择取消 ,然后点击停止作业 。
可选:如果您不再需要 Kafka 集群,请前往 Managed Service for Apache Kafka 集群 页面,选择
test-cluster,然后点击删除 。