本文档介绍如何为 Kafka Connect 创建 Cloud SQL for PostgreSQL 源连接器。
Cloud SQL for PostgreSQL 源连接器是 Debezium PostgreSQL 连接器的一个实例。 它会从 Cloud SQL for PostgreSQL数据库读取行级更改,并将其 写入 Managed Service for Apache Kafka 集群中的主题。
此连接器的使用场景包括:
- 实时监控行级数据库更改。
- 将数据库更改事件集成到事件驱动型架构中。
- 响应数据库事件,例如插入或删除行。
- 将数据库更改复制到其他系统。
准备工作
创建 Cloud SQL for PostgreSQL 源连接器之前,请确保您具备以下条件:
具有数据库的 Cloud SQL for PostgreSQL 实例。如需了解如何 创建这些资源,请参阅 使用 Google Cloud 控制台创建和查询 Cloud SQL for PostgreSQL 数据库。
一个 Connect 集群。
所需的角色和权限
如需获得创建连接器所需的权限,请让您的管理员为您授予项目的 Managed Kafka Connector Editor (roles/managedkafka.connectorEditor) IAM 角色。如需详细了解如何授予角色,请参阅管理对项目、文件夹和组织的访问权限。
此预定义角色可提供创建连接器所需的权限。如需查看所需的确切权限,请展开所需权限部分:
所需权限
创建连接器需要以下权限:
-
创建连接器:
managedkafka.connectors.create
授予从 Cloud SQL 读取数据的权限
Managed Kafka 服务帐号必须具有访问 Cloud SQL for PostgreSQL 的权限。向服务帐号授予以下 IAM 角色:
- Cloud SQL Client (
roles/cloudsql.client) - Cloud SQL Instance User (
roles/cloudsql.instanceUser)
Managed Kafka 服务帐号采用以下格式:
service-PROJECT_NUMBER@gcp-sa-managedkafka.iam.gserviceaccount.com,其中 PROJECT_NUMBER 是 Connect 集群的项目编号。
如果您的 Connect 集群与 Managed Service for Apache Kafka 集群位于不同的项目中,请参阅 在不同的项目中创建 Connect 集群。
配置数据库
创建连接器之前,您必须配置数据库复制,并允许连接器使用数据库进行身份验证。以下部分将介绍这些步骤。
启用逻辑解码
Cloud SQL for PostgreSQL 源连接器依赖于 PostgreSQL 的 逻辑解码 功能。如需在 Cloud SQL for PostgreSQL 实例中启用逻辑解码,请执行以下步骤。
控制台
前往 Cloud SQL > 实例。
点击实例的名称。
点击 修改。
展开标志和参数 。
点击添加数据库标志 。
在选择标志 列表中,选择
cloudsql.logical_decoding。对于值,选择
On。点击完成 。
点击保存 。
如需了解详情,请参阅 设置逻辑复制和解码。
配置变更数据捕获 (CDC)
在实例中启用逻辑解码后,为您要复制的表启用变更数据捕获 (CDC)。
如需为表启用 CDC,请运行 CREATE PUBLICATION SQL 语句。此
语句会创建一个 发布内容,用于定义要复制的一组表。
选项 1 。创建一个发布内容,用于复制数据库中所有表的更改。
CREATE PUBLICATION dbz_publication FOR ALL TABLES;选项 2 。为特定的一组表创建发布内容。
CREATE PUBLICATION dbz_publication FOR TABLE TABLE_LIST;将
TABLE_LIST替换为表的英文逗号分隔列表,格式为"schema_name"."table_name"。如所示,将架构和表名称用双引号括起来,可防止名称包含特殊字符或大写字母时出现语法错误。
默认情况下,连接器使用 dbz_publication 作为发布内容名称。如需使用具有不同名称的发布内容,请参阅
发布内容名称。
为 Managed Kafka 服务帐号创建用户账号
Cloud SQL for PostgreSQL 源连接器使用 IAM 数据库身份验证 连接到数据库。如需启用 IAM 数据库身份验证,请按如下方式将 Managed Kafka 服务帐号添加到 Cloud SQL 实例:
控制台
前往 Cloud SQL > 实例
点击实例的名称。
在导航窗格中,点击用户 。
点击添加用户账号 。
在添加用户账号 窗格中,选择 Cloud IAM 。
在 IAM 主账号 字段中,输入以下内容:
service-PROJECT_NUMBER@gcp-sa-managedkafka.iam.gserviceaccount.com将
PROJECT_NUMBER替换为 Connect 集群的项目编号。点击 Add (添加)。
gcloud
运行 gcloud sql users create
命令:
gcloud sql users create service-PROJECT_NUMBER@gcp-sa-managedkafka.iam \
--instance=INSTANCE_NAME \
--type=cloud_iam_service_account
替换以下内容:
PROJECT_NUMBER:Connect 集群的项目编号。INSTANCE_NAME:Cloud SQL for PostgreSQL 实例的名称。
由于数据库用户名有长度限制,因此系统会从
用户名中删除
.gserviceaccount.com后缀,因此用户名为
service-PROJECT_NUMBER@gcp-sa-managedkafka.iam。当您
运行引用 IAM 用户账号的 SQL 查询时,请指定
截断的名称。
配置用户账号
创建 IAM 用户账号后,以具有 cloudsqlsuperuser 角色的用户(例如默认的 postgres 用户)身份连接到数据库,然后运行以下 SQL 查询。
控制台
允许用户读取预写式日志。
ALTER USER "service-PROJECT_NUMBER@gcp-sa-managedkafka.iam" WITH REPLICATION;向用户授予对表的
SELECT权限。GRANT SELECT ON ALL TABLES IN SCHEMA "SCHEMA_NAME" TO "service-PROJECT_NUMBER@gcp-sa-managedkafka.iam";或者,您可以授予对单个表的
SELECT权限。如果您选择此选项,还必须将连接器的table.include.list配置属性设置为允许的表列表。 以下 SQL 查询授予对单个表的SELECT权限:GRANT SELECT ON TABLE SCHEMA_NAME.TABLE_NAME TO "service-PROJECT_NUMBER@gcp-sa-managedkafka.iam";对于每个表,授予用户访问该表的架构的权限。如果表位于默认的
public架构中,您可以跳过此步骤。GRANT USAGE ON SCHEMA SCHEMA_NAME TO "service-PROJECT_NUMBER@gcp-sa-managedkafka.iam";
配置网络
Cloud SQL for PostgreSQL 源连接器可以通过以下方式连接到 Cloud SQL 实例:
- 专用 IP
- Private Service Connect
- 公共 IP
如需详细了解这些选项,请参阅 选择连接到 Cloud SQL 的方式。 作为安全最佳实践,建议使用专用 IP 或 Private Service Connect,因为这些选项不需要连接到外部 IP 地址。
下表显示了每个选项的网络要求:
| IP 地址类型 | 要求 |
|---|---|
| 专用 IP | 为您的实例配置专用 IP。如需了解详情,请参阅 配置专用 IP。 |
| Private Service Connect |
|
| 公共 IP |
|
创建 Cloud SQL for PostgreSQL 源连接器
如需创建 Cloud SQL for PostgreSQL 源连接器,请执行以下步骤。
连接器初始化后,会执行以下操作:
- 创建数据库的初始快照。
- 为每个包含行的表创建一个 Kafka 主题。
- 对于每个数据库行,向相应主题发送一个更改事件。
连接器运行时,会继续向主题发送更改事件。如需详细了解初始快照,请参阅 Debezium 文档中的 快照 。
控制台
在 Google Cloud 控制台中,前往 Connect 集群 页面。
点击要在其中创建连接器的 Connect 集群。
点击创建连接器 。
对于连接器名称,请输入一个字符串。
如需了解如何命名连接器的指南,请参阅 Managed Service for Apache Kafka 资源命名指南。
对于连接器插件 ,选择 Cloud SQL for PostgreSQL 源 。
在实例 列表中,选择 Cloud SQL 实例。
在数据库 列表中,选择 Cloud SQL 数据库。
在主题前缀 字段中,输入要用于 Kafka 主题名称的前缀。为每个 Cloud SQL for PostgreSQL 源连接器选择唯一的前缀。
可选:在表名称字段中,输入要从中读取更改数据的表的英文逗号分隔列表,格式为
"schema_name"."table_name"。如果将此字段留空,连接器会从数据库中的所有非系统表中读取更改数据。可选:在配置 框中,添加配置属性或修改默认属性。如需了解详情,请参阅 配置连接器。
您可能需要替换以下属性的默认值:
可选:选择任务重启政策 。如需了解详情,请参阅 任务重启政策。
点击创建 。
gcloud
-
在 Google Cloud 控制台中,激活 Cloud Shell。
Cloud Shell 会话随即会在控制台 Google Cloud 的底部启动,并显示命令行提示符。Cloud Shell 是一个已安装 Google Cloud CLI 且已为当前项目设置值的 Shell 环境 。该会话可能需要几秒钟来完成初始化。
运行
gcloud managed-kafka connectors create命令:gcloud managed-kafka connectors create CONNECTOR_ID \ --location=LOCATION \ --connect-cluster=CONNECT_CLUSTER_ID \ --config-file=CONFIG_FILE替换以下内容:
CONNECTOR_ID:连接器的 ID 或名称。 如需了解如何命名连接器的指南,请参阅 Managed Service for Apache Kafka 资源命名指南。 连接器的名称是不可变的。LOCATION:您在其中创建连接器的位置。此位置必须与您创建 Connect 集群的位置相同。CONNECT_CLUSTER_ID:在其中创建连接器的 Connect 集群的 ID。CONFIG_FILE:连接器的 YAML 配置文件的路径。
以下是 Cloud SQL for PostgreSQL 源连接器的配置文件示例:
connector.class: io.debezium.connector.postgresql.PostgresConnector database.dbname: DATABASE_NAME driver.cloudSqlInstance: INSTANCE_ID driver.enableIamAuth: "true" driver.ipTypes: IP_TYPES driver.sslmode: disable key.converter: org.apache.kafka.connect.json.JsonConverter key.converter.schemas.enable: "false" plugin.name: pgoutput slot.name: SLOT_NAME table.include.list: TABLE_LIST topic.prefix: TOPIC_PREFIX value.converter: org.apache.kafka.connect.json.JsonConverter value.converter.schemas.enable: "true"替换以下内容:
INSTANCE_ID:包含数据库的 Cloud SQL 实例的 ID,格式如下:PROJECT_ID:REGION:INSTANCE_NAME
DATABASE_NAME:要从中读取数据的 Cloud SQL 数据库的名称。IP_TYPES: IP 地址类型的英文逗号分隔列表SLOT_NAME:要创建的 复制槽的名称。TABLE_LIST:要从中读取更改数据的表的英文逗号分隔列表,格式为"schema_name"."table_name"。TOPIC_PREFIX:要用于 Kafka 主题名称的前缀。
配置连接器
本部分介绍了一些您可以在连接器上设置的配置属性。如需查看完整列表,请参阅 Debezium 文档中的 Debezium 连接器(适用于 PostgreSQL) 。
IP 地址类型
driver.ipTypes 属性指定连接器用于连接到数据库的 IP 地址类型:
PRIVATE:专用 IPPSC:Private Service ConnectPUBLIC:公共 IP
driver.ipTypes 属性包含按首选
顺序排列的 IP 类型英文逗号分隔列表;例如,driver.ipTypes=PRIVATE,PUBLIC。
如需了解详情,请参阅配置网络。
发布内容名称
默认情况下,连接器会尝试从名为 dbz_publication 的发布内容进行流式传输。
如需指定其他发布内容,请将
publication.name=PUBLICATION_NAME添加到配置中,其中
PUBLICATION_NAME是发布内容名称。示例:publication.name=my_publication。
复制槽
PostgreSQL 使用复制槽来流式传输数据库表更改。 默认情况下,连接器会创建一个名为 debezium 的复制槽。如需使用其他槽名称,请设置 slot.name 属性。
如果您为同一数据库创建连接器的两个实例,则必须为每个连接器指定唯一的槽名称。
默认情况下,连接器会将
slot.drop.on.stop属性设置为false,以防止数据丢失。当您永久删除连接器时,必须手动删除连接器使用的复制槽。除非使用 slot.name 属性进行不同的配置,否则复制槽名称默认为 debezium。
我们建议您设置提醒,以监控源 PostgreSQL 数据库服务器上的 WAL 磁盘使用情况,并删除任何未使用的复制槽。
表过滤条件
默认情况下,连接器会捕获数据库中每个非系统表的更改数据。如需过滤要捕获的表,请指定以下一项或多项设置:
schema.include.list。要包含的架构的列表。schema.exclude.list。要排除的架构的列表。不能与schema.include.list搭配使用。table.include.list。要包含的表的列表。table.exclude.list。要排除的表的列表。不能与table.include.list搭配使用。
主题名称
默认情况下,连接器会使用以下命名惯例创建 Kafka 主题:topic_prefix.schema.table_name,其中 topic.prefix 是 topic.prefix 配置的值。
如需了解详情,请参阅 Debezium 文档中的主题名称。