创建通用 PostgreSQL 源连接器

本文档介绍了如何创建通用 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.,其中 PROJECT_NUMBER 是 Connect 集群的项目编号。

如果您的 Connect 集群与 Managed Service for Apache Kafka 集群位于不同的项目中,请参阅 在不同的项目中创建 Connect 集群

配置 PostgreSQL 数据库

如需让连接器能够从数据库读取数据更改事件,请配置以下设置。

  1. 将服务器的 wal_level 设置为 logical

    ALTER SYSTEM SET wal_level = logical;
    

    重启服务器以应用该设置。

  2. 为连接器创建一个数据库用户,以便向 PostgreSQL 进行身份验证。数据库用户必须是复制角色,以便在复制模式下连接到服务器。

    CREATE ROLE ROLE_NAME WITH REPLICATION LOGIN PASSWORD 'ROLE_PASSWORD';
    

    替换以下内容:

    • ROLE_NAME:用户的名称,例如 debezium_user
    • ROLE_PASSWORD:用户的密码。
  3. 为您要捕获的表创建 发布内容 。连接器会订阅该发布内容以接收数据更改事件。

    CREATE PUBLICATION dbz_publication FOR TABLE "SCHEMA_NAME"."TABLE_NAME";
    

    替换以下内容:

    • SCHEMA_NAME:表的架构。

    • TABLE_NAME:表的名称。

    我们建议您将架构名称和表名称用双引号括起来,如上所示,以避免名称包含特殊字符或大写字母时出现语法错误。

    或者,您可以创建一个发布内容,用于复制数据库中所有表的更改:

    CREATE PUBLICATION dbz_publication FOR ALL TABLES;
    

    您可以根据连接器的 publication.autocreate.mode 设置手动创建发布内容,也可以让连接器自动创建。 如需了解详情,请参阅发布模式

  4. 对于每个表,向数据库用户授予对该表的 SELECT 权限。

    GRANT SELECT ON TABLE "SCHEMA_NAME"."TABLE_NAME" TO ROLE_NAME;
    

    或者,您可以授予对架构中所有表的选择权限:

    GRANT SELECT ON ALL TABLES IN SCHEMA "SCHEMA_NAME" TO ROLE_NAME;
    
  5. 对于每个表,向数据库用户授予对该表架构的 USAGE 权限。如果该表位于默认的 public 架构中,您可以跳过此步骤。

    GRANT USAGE ON SCHEMA "SCHEMA_NAME" TO ROLE_NAME;
    

创建通用 PostgreSQL 源连接器

如需创建通用 PostgreSQL 源连接器,请执行以下步骤。

连接器初始化后,会执行以下操作:

  1. 创建数据库的初始快照。
  2. 为每个包含行的表创建一个 Kafka 主题。
  3. 对于每个数据库行,向相应主题发送一个更改事件。

连接器运行期间,会继续向主题发送更改事件。如需详细了解初始快照,请参阅 Debezium 文档中的 快照

控制台

  1. 在 Google Cloud 控制台中,前往 Connect 集群 页面。

    前往 Connect 集群

  2. 点击您要在其中创建连接器的 Connect 集群。

  3. 点击创建连接器

  4. 对于连接器名称,请输入一个字符串。

    有关如何命名连接器的指南,请参阅 Managed Service for Apache Kafka 资源命名指南

  5. 对于连接器插件,请选择Generic PostgreSQL Source

  6. 数据库主机名 字段中,输入 PostgreSQL 服务器的主机名或 IP 地址。

  7. 数据库名称 字段中,输入数据库的名称。

  8. 数据库用户 字段中,输入副本角色的名称。连接器使用此角色向 PostgreSQL 服务器进行身份验证。

  9. 主题前缀 字段中,输入要用于 Kafka 主题名称的前缀。

  10. Secret 列表中,选择包含数据库密码的 Secret。

  11. 可选:在配置 框中,添加配置属性或修改默认属性。如需了解详情,请参阅 配置连接器

  12. 可选:选择任务重启政策 。如需了解详情,请参阅 任务重启政策

  13. 点击创建

gcloud

  1. 在 Google Cloud 控制台中,激活 Cloud Shell。

    激活 Cloud Shell

    Cloud Shell 会话随即会在控制台 Google Cloud 的底部启动,并显示命令行提示符。Cloud Shell 是一个已安装 Google Cloud CLI 且已为当前项目设置值的 Shell 环境 。该会话可能需要几秒钟来完成初始化。

  2. 运行 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.passworddatabase.sslpassword 配置中仅支持 Secret 路径。后端希望这些配置使用以下格式: ${directory:/var/secrets:PROJECT_ID-SECRET_NAME-SECRET_VERSION}

IP 地址类型

driver.ipTypes 属性指定连接器用于连接到数据库的 IP 地址类型:

  • PRIVATE:专用 IP
  • PSC:Private Service Connect
  • PUBLIC:公共 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.prefixtopic.prefix 配置的值。

如需了解详情,请参阅 Debezium 文档中的主题名称

后续步骤