本教程将展示如何使用 MirrorMaker 2.0 源连接器在 Managed Service for Apache Kafka 集群之间复制数据。
集群可以位于同一项目中,也可以位于不同项目中,并且可以位于同一区域,也可以位于不同区域。在本教程中,您将设置跨区域但在同一项目内的复制。不过,其他组合的步骤是相同的。
架构
在此场景中,有三个集群:
正在复制数据的 Kafka 集群。此集群称为源集群,因为它是数据的来源。
写入复制数据的 Kafka 集群。此集群称为目标集群。
一个 Connect 集群,可用于创建和管理 MirrorMaker 2.0 源连接器。
Connect 集群具有主 Kafka 集群,即与 Connect 集群关联的 Kafka 集群。为了最大限度地缩短写入操作的延迟时间,建议将目标集群指定为主集群,并将 Connect 集群与目标集群放在同一区域中。
下图展示了这些组件:

数据复制由 MirrorMaker 2.0 源连接器执行。在此场景中,还有两个 MirrorMaker 2.0 连接器是可选的:
MirrorMaker 2.0 检查点连接器:通过确保目标集群上的使用方可以从与源集群相同的点恢复处理,实现无缝故障切换。
MirrorMaker 2.0 检测信号连接器:在源 Kafka 集群上定期生成检测信号消息,以便您监控数据复制的健康状况和状态。
本教程不使用这些可选连接器。如需了解详情,请参阅何时使用 MirrorMaker 2.0。
准备工作
控制台
- 登录您的 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
检查角色
-
在 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.connectorEditorgcloud projects add-iam-policy-binding PROJECT_ID --member="user:USER_IDENTIFIER" --role=ROLE
替换以下内容:
PROJECT_ID:您的项目 ID。USER_IDENTIFIER:您的用户 账号的标识符。例如,myemail@example.com。ROLE:您授予用户账号的 IAM 角色。
创建源 Kafka 集群
在此步骤中,您将创建一个 Managed Service for Apache Kafka 集群。此集群是源集群,其中包含要复制的数据。
控制台
前往 Managed Service for Apache Kafka > 集群页面。
点击 创建。
在集群名称字段中,输入集群的名称。
在区域列表中,为集群选择一个位置。
对于网络配置,请配置集群可访问的子网。
- 在项目部分,选择您的项目。
- 对于网络,选择项目中的 VPC 网络。
- 对于子网,请选择一个子网。
- 点击完成。
点击创建。
在创建集群期间,集群状态为 Creating。集群创建完成后,状态为 Active。
gcloud
如需创建源 Kafka 集群,请运行 managed-kafka clusters create 命令。
gcloud managed-kafka clusters create SOURCE_KAFKA_CLUSTER \
--location=SOURCE_REGION \
--cpu=3 \
--memory=3GiB \
--subnets=projects/PROJECT_ID/regions/SOURCE_REGION/subnetworks/SOURCE_SUBNET \
--async
替换以下内容:
SOURCE_KAFKA_CLUSTER:Kafka 集群的名称SOURCE_REGION:集群的位置如需了解支持的位置,请参阅 Managed Service for Apache Kafka 位置。
PROJECT_ID:您的项目 IDSOURCE_SUBNET:您要部署集群的子网;例如default
该命令会异步运行,并返回一个操作 ID:
Check operation [projects/PROJECT_ID/locations/SOURCE_SUBNET/operations/OPERATION_ID] for status.
如需跟踪创建操作的进度,请使用 gcloud managed-kafka operations describe 命令:
gcloud managed-kafka operations describe OPERATION_ID \
--location=SOURCE_REGION
如需了解详情,请参阅监控集群创建操作。
创建 Kafka 主题
当源集群准备就绪后,按如下方式创建主题:
控制台
前往 Managed Service for Apache Kafka > 集群页面。
点击集群的名称。
在集群详细信息页面中,点击 创建主题。
在主题名称框中,输入主题的名称。
点击创建。
gcloud
如需创建 Kafka 主题,请运行 managed-kafka topics create 命令。
gcloud managed-kafka topics create TOPIC_NAME \
--cluster=SOURCE_KAFKA_CLUSTER \
--location=SOURCE_REGION \
--partitions=10 \
--replication-factor=3
替换以下内容:
TOPIC_NAME:要创建的 Kafka 主题的名称SOURCE_KAFKA_CLUSTER:源集群的名称SOURCE_REGION:您创建源集群的区域
创建目标 Kafka 集群
在此步骤中,您将创建第二个 Managed Service for Apache Kafka 集群。此集群是 目标集群,MirrorMaker 2.0 会将数据复制到此集群。
控制台
前往 Managed Service for Apache Kafka > 集群页面。
点击 创建。
在集群名称字段中,输入集群的名称。
在区域列表中,为集群选择一个位置。选择与源集群所在区域不同的区域。
对于网络配置,请配置目标集群可访问的子网。子网必须与源集群的子网位于同一 VPC 网络中。
- 在项目部分,选择您的项目。
- 对于网络,请选择与源集群的子网相同的 VPC 网络。
- 在子网字段中,选择相应子网。
- 点击完成。
点击创建。
在创建集群期间,集群状态为 Creating。集群创建完成后,状态为 Active。
gcloud
如需创建目标 Kafka 集群,请运行 managed-kafka clusters create 命令。
gcloud managed-kafka clusters create TARGET_KAFKA_CLUSTER \
--location=TARGET_REGION \
--cpu=3 \
--memory=3GiB \
--subnets=projects/PROJECT_ID/regions/TARGET_REGION/subnetworks/TARGET_SUBNET \
--async
替换以下内容:
TARGET_KAFKA_CLUSTER:Kafka 集群的名称TARGET_REGION:集群的位置;选择与源集群所在区域不同的区域。如需了解支持的位置,请参阅 Managed Service for Apache Kafka 位置。
PROJECT_ID:您的项目 IDTARGET_SUBNET:您要部署集群的子网;例如default选择与源集群位于同一 VPC 中的子网。
创建 Connect 集群
在此步骤中,您将创建一个 Connect 集群。创建 Connect 集群通常需要 20-30 分钟。
在开始执行此步骤之前,请确保上一步中的目标 Kafka 集群已完全创建。
控制台
前往 Managed Service for Apache Kafka > Connect 集群页面。
点击 创建。
在 Connect 集群名称中,输入一个字符串。示例:
my-connect-cluster。对于主 Kafka 集群,请选择您在上一步中创建的目标 Kafka 集群。(请勿选择源 Kafka 集群。)
对于位置、网络配置和工作器子网,您可以选择使用默认值,也可以根据自己的具体需求进行自定义。
如需使关联集群能够解析源集群的 DNS 网域,请执行以下步骤:
展开可解析的 DNS 域名。
点击添加 DNS 网域。
在 Kafka 集群列表中,选择源 Kafka 集群。
点击创建。
在创建集群期间,集群状态为 Creating。集群创建完成后,状态为 Active。
gcloud
如需创建 Connect 集群,请运行 gcloud managed-kafka connect-clusters create 命令。
gcloud managed-kafka connect-clusters create CONNECT_CLUSTER \
--location=TARGET_REGION \
--cpu=12 \
--memory=12GiB \
--primary-subnet=projects/PROJECT_ID/regions/TARGET_REGION/subnetworks/TARGET_SUBNET \
--kafka-cluster=TARGET_KAFKA_CLUSTER \
--dns-name=SOURCE_KAFKA_CLUSTER.SOURCE_REGION.managedkafka.PROJECT_ID.cloud.goog. \
--async
替换以下内容:
CONNECT_CLUSTER:Connect 集群的名称TARGET_REGION:您创建目标 Kafka 集群的区域PROJECT_ID:您的项目 IDTARGET_SUBNET:您在其中创建目标 Kafka 集群的子网TARGET_KAFKA_CLUSTER:目标 Kafka 集群的名称SOURCE_KAFKA_CLUSTER:源 Kafka 集群的名称SOURCE_REGION:您创建源 Kafka 集群的区域
该命令会异步运行,并返回一个操作 ID:
Check operation [projects/PROJECT_ID/locations/TARGET_REGION/operations/OPERATION_ID] for status.
如需跟踪创建操作的进度,请使用 gcloud managed-kafka operations describe 命令:
gcloud managed-kafka operations describe OPERATION_ID \
--location=TARGET_REGION
如需了解详情,请参阅监控集群创建操作。
创建 MirrorMaker 2.0 来源连接器
在此步骤中,您将创建 MirrorMaker 2.0 源连接器。 此连接器会将消息从源 Kafka 集群复制到目标 Kafka 集群。
控制台
前往 Managed Service for Apache Kafka > Connect 集群页面。
点击 Connect 集群的名称。
点击 Create connector。
在连接器名称中,输入一个字符串。示例:
mm2-connector。在连接器插件列表中,选择
MirrorMaker 2.0 Source。选择将 Kafka 主集群用作目标集群。
对于源集群,请选择 Managed Service for Apache Kafka 集群。
在 Kafka 集群列表中,选择源集群。
在以逗号分隔的主题名称或主题正则表达式框中,输入要复制的 Kafka 主题的名称。
点击创建。
gcloud
如需创建 MirrorMaker 2.0 源连接器,请运行 gcloud managed-kafka connectors create 命令。
gcloud managed-kafka connectors create CONNECTOR_NAME \
--location=TARGET_REGION \
--connect-cluster=CONNECT_CLUSTER \
--configs=connector.class=org.apache.kafka.connect.mirror.MirrorSourceConnector,\
source.cluster.alias=source,\
source.cluster.bootstrap.servers=bootstrap.SOURCE_KAFKA_CLUSTER.SOURCE_REGION.managedkafka.PROJECT_ID.cloud.goog:9092,\
target.cluster.alias=target,\
target.cluster.bootstrap.servers=bootstrap.TARGET_KAFKA_CLUSTER.TARGET_REGION.managedkafka.PROJECT_ID.cloud.goog:9092,\
tasks.max=3,\
topics=TOPIC_NAME
替换以下内容:
CONNECTOR_NAME:连接器的名称,例如mm2-connectorTARGET_REGION:您创建 Connect 集群和目标 Kafka 集群的区域CONNECT_CLUSTER:您的 Connect 集群的名称SOURCE_KAFKA_CLUSTER:源 Kafka 集群的名称SOURCE_REGION:您创建源 Kafka 集群的区域PROJECT_ID:您的项目 IDTARGET_KAFKA_CLUSTER:目标 Kafka 集群的名称TOPIC_NAME:要复制的主题的名称。此参数还可以指定逗号分隔列表形式的主题名称或正则表达式。
MirrorMaker 2.0 源连接器会在目标集群中创建一个名为 "source.TOPIC_NAME" 的新主题,其中 TOPIC_NAME 是源集群中主题的名称。
查看结果
如需验证消息是否正在复制,您可以使用 Kafka 命令行工具。如需了解如何设置 Kafka CLI,请参阅使用 CLI 生成和使用消息文档中的设置客户端机器。
例如,如需向源集群发送消息,请在命令行中输入以下内容:
export BOOTSTRAP=bootstrap.SOURCE_KAFKA_CLUSTER.SOURCE_REGION.managedkafka.PROJECT_ID.cloud.goog:9092
for msg in {1..10}; do
echo "message $msg"
done | kafka-console-producer.sh --topic TOPIC_NAME \
--bootstrap-server $BOOTSTRAP --producer.config client.properties
如需从目标集群读取重复的消息,请在命令行中输入以下内容:
export BOOTSTRAP=bootstrap.TARGET_KAFKA_CLUSTER.TARGET_REGION.managedkafka.PROJECT_ID.cloud.goog:9092
kafka-console-consumer.sh --topic source.TOPIC_NAME --from-beginning \
--bootstrap-server $BOOTSTRAP --consumer.config client.properties
输出如下所示:
message 1
message 2
message 3
message 4
[...]
清理
为避免因本教程中使用的资源导致您的 Google Cloud 账号产生费用,请删除包含这些资源的项目,或者保留项目但删除各个资源。
控制台
gcloud
如需删除 Connect 集群,请使用
gcloud managed-kafka connect-clusters delete命令。gcloud managed-kafka connect-clusters delete CONNECT_CLUSTER \ --location=TARGET_REGION --async如需删除源 Kafka 集群,请使用
gcloud managed-kafka clusters delete命令。gcloud managed-kafka clusters delete SOURCE_KAFKA_CLUSTER \ --location=SOURCE_REGION --async重复上一步,以删除目标 Kafka 集群。
gcloud managed-kafka clusters delete TARGET_KAFKA_CLUSTER \ --location=TARGET_REGION --async
后续步骤
- 排查 MirrorMaker 2.0 连接器问题
- 详细了解 MirrorMaker 2.0 连接器。
- 详细了解 Kafka Connect。