本文档介绍了如何创建通用 PostgreSQL 源连接器。
通用 PostgreSQL 源连接器是 Debezium PostgreSQL 连接器的一个实例。 它会从 PostgreSQL 数据库读取行级更改,并将其写入 Managed Service for Apache Kafka 集群中的主题。
此连接器的使用场景包括:
- 实时监控行级数据库更改。
- 将数据库更改事件集成到事件驱动型架构中。
- 响应数据库事件,例如插入或删除行。
- 将数据库更改复制到其他系统。
- 复制或恢复 PostgreSQL 表。
准备工作
在创建通用 PostgreSQL 源连接器之前,请确保您具备以下条件:
PostgreSQL 数据库。
与 Kafka 集群关联的 Connect 集群 。
创建一个 Secret Manager Secret,用于存储数据库密码。如果您的配置使用数据库 SSL,请同时为数据库 SSL 密码创建一个 Secret。使用 Secret 配置 Connect 集群。如需了解详情,请参阅 Secret Manager 资源。
所需的角色和权限
如需获得创建连接器所需的权限,请让您的管理员为您授予项目的 Managed Kafka Connector Editor (roles/managedkafka.connectorEditor) IAM 角色。如需详细了解如何授予角色,请参阅管理对项目、文件夹和组织的访问权限。
此预定义角色可提供创建连接器所需的权限。如需查看所需的确切权限,请展开所需权限部分:
所需权限
创建连接器需要以下权限:
-
创建连接器:
managedkafka.connectors.create
授予访问 Secret Manager Secret 的权限
Managed Kafka 服务帐号需要具有查看和访问存储在 Secret Manager 中的 Secret 的权限。 向服务帐号授予以下 IAM 角色:
- Secret Manager Viewer (
roles/secretmanager.viewer) - Secret Manager Secret Accessor (
roles/secretmanager.secretAccessor)
Managed Kafka 服务帐号采用以下格式:
service-PROJECT_NUMBER@gcp-sa-managedkafka.iam.gserviceaccount.com,其中 PROJECT_NUMBER 是 Connect 集群的项目编号。
如果您的 Connect 集群与 Managed Service for Apache Kafka 集群位于不同的项目中,请参阅 在不同的项目中创建 Connect 集群。
配置 PostgreSQL 数据库
如需让连接器能够从数据库读取数据更改事件,请配置以下设置。
将服务器的
wal_level设置为logical。ALTER SYSTEM SET wal_level = logical;重启服务器以应用该设置。
为连接器创建一个数据库用户,以便向 PostgreSQL 进行身份验证。数据库用户必须是复制角色,以便在复制模式下连接到服务器。
CREATE ROLE ROLE_NAME WITH REPLICATION LOGIN PASSWORD 'ROLE_PASSWORD';替换以下内容:
ROLE_NAME:用户的名称,例如debezium_user。ROLE_PASSWORD:用户的密码。
为您要捕获的表创建 发布内容 。连接器会订阅该发布内容以接收数据更改事件。
CREATE PUBLICATION dbz_publication FOR TABLE "SCHEMA_NAME"."TABLE_NAME";替换以下内容:
SCHEMA_NAME:表的架构。TABLE_NAME:表的名称。
我们建议您将架构名称和表名称用双引号括起来,如上所示,以避免名称包含特殊字符或大写字母时出现语法错误。
或者,您可以创建一个发布内容,用于复制数据库中所有表的更改:
CREATE PUBLICATION dbz_publication FOR ALL TABLES;您可以根据连接器的
publication.autocreate.mode设置手动创建发布内容,也可以让连接器自动创建。 如需了解详情,请参阅发布模式。对于每个表,向数据库用户授予对该表的
SELECT权限。GRANT SELECT ON TABLE "SCHEMA_NAME"."TABLE_NAME" TO ROLE_NAME;或者,您可以授予对架构中所有表的选择权限:
GRANT SELECT ON ALL TABLES IN SCHEMA "SCHEMA_NAME" TO ROLE_NAME;对于每个表,向数据库用户授予对该表架构的
USAGE权限。如果该表位于默认的public架构中,您可以跳过此步骤。GRANT USAGE ON SCHEMA "SCHEMA_NAME" TO ROLE_NAME;
创建通用 PostgreSQL 源连接器
如需创建通用 PostgreSQL 源连接器,请执行以下步骤。
连接器初始化后,会执行以下操作:
- 创建数据库的初始快照。
- 为每个包含行的表创建一个 Kafka 主题。
- 对于每个数据库行,向相应主题发送一个更改事件。
连接器运行期间,会继续向主题发送更改事件。如需详细了解初始快照,请参阅 Debezium 文档中的 快照 。
控制台
在 Google Cloud 控制台中,前往 Connect 集群 页面。
点击您要在其中创建连接器的 Connect 集群。
点击创建连接器 。
对于连接器名称,请输入一个字符串。
有关如何命名连接器的指南,请参阅 Managed Service for Apache Kafka 资源命名指南。
对于连接器插件,请选择Generic PostgreSQL Source。
在数据库主机名 字段中,输入 PostgreSQL 服务器的主机名或 IP 地址。
在数据库名称 字段中,输入数据库的名称。
在数据库用户 字段中,输入副本角色的名称。连接器使用此角色向 PostgreSQL 服务器进行身份验证。
在主题前缀 字段中,输入要用于 Kafka 主题名称的前缀。
在 Secret 列表中,选择包含数据库密码的 Secret。
可选:在配置 框中,添加配置属性或修改默认属性。如需了解详情,请参阅 配置连接器。
可选:选择任务重启政策 。如需了解详情,请参阅 任务重启政策。
点击创建 。
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 配置文件的路径。
以下是通用 PostgreSQL 源连接器的配置文件示例:
connector.class: io.debezium.connector.postgresql.PostgresConnector database.dbname: DATABASE_NAME database.hostname: HOSTNAME database.password: CREDENTIALS database.user: DATABASE_USER key.converter: org.apache.kafka.connect.json.JsonConverter key.converter.schemas.enable: "false" plugin.name: pgoutput topic.prefix: TOPIC_PREFIX value.converter: org.apache.kafka.connect.json.JsonConverter value.converter.schemas.enable: "true"替换以下内容:
HOSTNAME:要从中读取数据的 PostgreSQL 数据库的主机名。DATABASE_NAME:要从中读取数据的 PostgreSQL 数据库的名称。DATABASE_USER:向数据库进行身份验证时要使用的 PostgreSQL 数据库用户。CREDENTIALS:包含数据库密码的 Secret Manager Secret 的路径。使用以下格式指定 Secret:
${directory:/var/secrets:PROJECT_ID-SECRET_NAME-SECRET_VERSION}TOPIC_PREFIX:要用于 Kafka 主题名称的前缀。
配置连接器
本部分介绍了一些您可以在连接器上设置的配置属性。如需查看完整列表,请参阅 Debezium 文档中的 Debezium 连接器(适用于 PostgreSQL) 。
密码和 SSL 密码配置
database.password 和 database.sslpassword 配置中仅支持 Secret 路径。后端希望这些配置使用以下格式:
${directory:/var/secrets:PROJECT_ID-SECRET_NAME-SECRET_VERSION}。
IP 地址类型
driver.ipTypes 属性指定连接器用于连接到数据库的 IP 地址类型:
PRIVATE:专用 IPPSC:Private Service ConnectPUBLIC:公共 IP
driver.ipTypes 属性包含一个 IP 类型逗号分隔列表,并按首选顺序排列;例如,driver.ipTypes=PRIVATE,PUBLIC。
发布模式
通用 PostgreSQL 源连接器会从数据库中的 发布内容 流式传输更改事件。您可以手动创建发布内容,也可以让连接器自动创建。
publication.autocreate.mode
设置指定连接器应如何创建发布内容以及是否应创建发布内容。
filtered。如果发布内容不存在,连接器会创建一个新发布内容,其中仅包含捕获的表。数据库用户必须对数据库具有CREATE权限,并且是所包含表的所有者。如果发布内容已存在,连接器会对其进行更改以包含捕获的表。如需更改现有发布内容,数据库用户必须是发布内容的所有者以及所包含表的所有者。
all_tables。如果发布内容不存在,连接器会使用FOR ALL TABLES参数创建一个新发布内容。数据库用户必须是超级用户。超级用户角色会绕过数据库中的所有权限检查,因此不建议向数据库用户授予
SUPERUSER。请改为手动创建发布内容,或设置publication.autocreate.mode=filtered。disabled。如果发布内容不存在,则会发生错误。连接器不会创建新的发布内容。
默认值为 all_tables。
发布内容名称
默认情况下,连接器会尝试从名为 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 文档中的主题名称。