创建 Cloud SQL for PostgreSQL 源连接器

本文档介绍如何为 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 源连接器之前,请确保您具备以下条件:

所需的角色和权限

如需获得创建连接器所需的权限,请让您的管理员为您授予项目的 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 实例中启用逻辑解码,请执行以下步骤。

控制台

  1. 前往 Cloud SQL > 实例

    转到实例

  2. 点击实例的名称。

  3. 点击 修改

  4. 展开标志和参数

  5. 点击添加数据库标志

  6. 选择标志 列表中,选择 cloudsql.logical_decoding

  7. 对于,选择On

  8. 点击完成

  9. 点击保存

如需了解详情,请参阅 设置逻辑复制和解码

配置变更数据捕获 (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 实例:

控制台

  1. 前往 Cloud SQL > 实例

    转到实例

  2. 点击实例的名称。

  3. 在导航窗格中,点击用户

  4. 点击添加用户账号

  5. 添加用户账号 窗格中,选择 Cloud IAM

  6. IAM 主账号 字段中,输入以下内容:

    service-PROJECT_NUMBER@gcp-sa-managedkafka.iam.gserviceaccount.com
    

    PROJECT_NUMBER 替换为 Connect 集群的项目编号。

  7. 点击 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 查询。

控制台

  1. 允许用户读取预写式日志。

    ALTER USER "service-PROJECT_NUMBER@gcp-sa-managedkafka.iam" WITH REPLICATION;
    
  2. 向用户授予对表的 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";
    
  3. 对于每个表,授予用户访问该表的架构的权限。如果表位于默认的 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
  1. 为您的实例配置 Private Service Connect 并获取 Private Service Connect 端点的 DNS 名称。如需了解详情,请参阅 使用 Private Service Connect 连接到实例
  2. 将端点的 DNS 名称添加到 Connect 集群的 可解析 DNS 网域。如需了解详情,请参阅 更新 Connect 集群
公共 IP
  1. 为您的实例配置公共 IP。如需了解详情,请参阅 配置公共 IP
  2. 设置 Public NAT,以使 Connect 集群工作器能够 与互联网通信。如需了解详情,请参阅 设置 Public NAT。创建 Cloud NAT 网关时,请指定包含 Connect 集群 主子网的 VPC 网络。

创建 Cloud SQL for PostgreSQL 源连接器

如需创建 Cloud SQL for PostgreSQL 源连接器,请执行以下步骤。

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

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

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

控制台

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

    前往 Connect 集群

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

  3. 点击创建连接器

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

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

  5. 对于连接器插件 ,选择 Cloud SQL for PostgreSQL 源

  6. 实例 列表中,选择 Cloud SQL 实例。

  7. 数据库 列表中,选择 Cloud SQL 数据库。

  8. 主题前缀 字段中,输入要用于 Kafka 主题名称的前缀。为每个 Cloud SQL for PostgreSQL 源连接器选择唯一的前缀。

  9. 可选:在表名称字段中,输入要从中读取更改数据的表的英文逗号分隔列表,格式为"schema_name"."table_name"。如果将此字段留空,连接器会从数据库中的所有非系统表中读取更改数据。

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

    您可能需要替换以下属性的默认值:

    • driver.ipTypes:此属性必须与 网络配置的 Cloud SQL 实例匹配。请参阅 IP 地址类型

    • slot.name:如果您为同一数据库创建连接器的多个实例,请为每个连接器指定唯一的值。请参阅 复制槽

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

  12. 点击创建

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 配置文件的路径。

    以下是 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_TYPESIP 地址类型的英文逗号分隔列表

    • SLOT_NAME:要创建的 复制槽的名称。

    • TABLE_LIST:要从中读取更改数据的表的英文逗号分隔列表,格式为 "schema_name"."table_name"

    • TOPIC_PREFIX:要用于 Kafka 主题名称的前缀。

配置连接器

本部分介绍了一些您可以在连接器上设置的配置属性。如需查看完整列表,请参阅 Debezium 文档中的 Debezium 连接器(适用于 PostgreSQL)

IP 地址类型

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

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

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

后续步骤