源连接器

扳手

qianmoQqianmoQ· 更新于 2026-09-28· 阅读 88 分钟· 0 次阅读

登录后可跨设备保存划线和私人笔记登录

Debezium Cloud Spanner 连接器

Debezium 的 Cloud Spanner 连接器会消费 Cloud Spanner 变更流数据,并将其输出到 Kafka 主题中。

Spanner 变更流会近实时地监视并流出 Spanner 数据库中的数据变更——插入、更新、删除。Spanner 连接器抽象了查询 Spanner 变更流的细节。使用该连接器,你无需管理变更流分区的生命周期,而这是直接使用 Spanner API 时必须处理的。

该连接器目前不支持快照功能。Kafka 连接器首次连接到 Spanner 数据库时,会从所提供的时间戳开始流出变更;如果未提供时间戳,则从当前时间戳开始流出。

概述

为了读取和处理数据库变更,连接器通过查询变更流来实现。

对于捕获到的每一行级插入、更新和删除操作,连接器都会生成一个变更事件,并为每个表在独立的 Kafka 主题中发送变更事件记录。客户端应用读取与其关注的数据库表相对应的 Kafka 主题,并可以对从这些主题接收到的每个行级事件作出响应。

连接器具有容错能力。在连接器读取变更并生成事件的过程中,它会记录每个变更流分区已处理的最后提交时间戳。如果连接器因任何原因停止(包括通信故障、网络问题或崩溃),重启后连接器会从中断处继续流出记录。

连接器的工作原理

要最佳地配置和运行 Debezium Spanner 连接器,了解连接器如何流出变更事件、如何确定 Kafka 主题名称以及如何使用元数据会很有帮助。

流出变更

Debezium Spanner 连接器会将全部时间用于从其订阅的变更流中流出变更。当某张表发生变更时,Spanner 会在数据库中写入一条相应的变更流记录,且该写入与数据变更在同一个事务中同步完成。为了扩展变更流的写入和读取能力,Spanner 会随着数据库数据一起拆分和合并变更流的内部存储。为了在数据库写入规模扩大的同时支持近实时读取变更流记录,Spanner API 被设计为可通过变更流分区并发查询变更流。参见 Spanner 变更流的分区模型。

该连接器为查询变更流提供了 Spanner API 的抽象层。通过该连接器,你无需管理变更流分区的生命周期。连接器会为你提供一个数据变更记录流,让你可以更专注于应用逻辑,而无需过多关注具体的 API 细节和动态变更流分区。

订阅变更流时,连接器需要提供项目 ID、Spanner 实例 ID、Spanner 数据库 ID 以及变更流名称。用户还可以选择性地提供起始时间戳和结束时间戳。有关连接器配置属性的详细列表,请参阅本节。

当连接器接收到变更时,它会将这些事件转换为 Debezium 的 create、update 或 delete 事件。连接器将这些变更事件以记录的形式转发给在同一进程中运行的 Kafka Connect 框架。Kafka Connect 进程会按照事件生成的相同顺序,异步地将变更事件记录写入相应的 Kafka 主题。

Kafka Connect 会定期在另一个 Kafka 主题中记录最新的 offset(偏移量)。偏移量表示 Debezium 随每个事件附带的源特定位置信息。对于 Spanner 连接器,偏移量即为变更流分区最后处理的提交时间戳。

当 Kafka Connect 正常关闭时,它会停止连接器,将所有事件记录刷新到 Kafka,并记录从每个连接器接收到的最后一个偏移量。当 Kafka Connect 重新启动时,它会读取每个连接器最后记录的偏移量,并从该偏移量处启动每个连接器。

在流式传输过程中,Spanner 还会在以下元数据主题中记录元数据信息。不建议修改这些主题的内容或配置:

  • _consumer_offsets:由 Kafka 自动创建的主题。存储 Kafka 连接器中创建的使用者的偏移量。
  • _kafka-connect-offsets:由 Kafka Connect 自动创建的主题。存储连接器的偏移量。
  • _sync_topic_spanner_connector_connectorname:由连接器自动创建的主题。存储有关变更流分区的元数据。
  • _rebalancing_topic_spanner_connector_connectorname:由连接器自动创建的主题。用于确定连接器任务的存活状态。
  • _debezium-heartbeat.connectorname:用于处理 Spanner 变更流心跳的主题。

主题名称

Spanner 连接器会将单个表上所有插入、更新和删除操作的事件写入单个 Kafka 主题。默认情况下,Kafka 主题名称为 topicPrefix.connectorName.tableName,其中:

  • topicPrefix 是 topic.prefix 连接器配置属性所指定的主题前缀。
  • connectorName 是用户为连接器指定的名称。
  • tableName 是发生该操作的数据库表的名称。

例如,假设某个连接器的配置中逻辑名称为 spanner,该连接器从一条 Spanner 变更流中捕获更改,而这条变更流所跟踪的数据库包含四个表:table1、table2、table3 和 table4。该连接器会将记录流式传输到这四个 Kafka 主题:

  • spanner.table1
  • spanner.table2
  • spanner.table3
  • spanner.table4

数据变更事件

Debezium Spanner 连接器为每个行级的 INSERT、UPDATE 和 DELETE 操作生成一个数据变更事件。每个事件都包含一个键和一个值。键和值是相互独立的文档。键和值的结构取决于被更改的表。

Debezium 和 Kafka Connect 是围绕事件消息的持续流设计的。然而,这些事件的结构可能会随时间而变化,这会给消费者带来处理上的困难。为了解决这个问题,每个事件都包含其内容的架构,或者(如果你使用的是模式注册表)包含一个架构 ID,消费者可以使用该 ID 从注册表中获取架构。这样每个事件就是自包含的。键的架构永远不会改变。请注意,值的架构是变更流自连接器启动时间以来所跟踪的表中所有列的综合结果。

以下骨架 JSON 文档展示了键文档和值文档的基本结构。不过,键和值文档的表现形式取决于你在应用中所选择配置的 Kafka Connect 转换器。只有当你配置转换器生成 schema 字段时,变更事件键或变更事件值中才会出现该字段。同样地,只有当你配置转换器生成时,事件键和事件负载才会出现。如果你使用 JSON 转换器并配置它生成架构,变更事件将具有如下结构:

// Key
{
 "schema": { (1)
   ...
  },
 "payload": { (2)
   ...
 }
}

// Value
{
 "schema": { (3)
   ...
 },
 "payload": { (4)
   ...
 }
}
项字段名称描述
1schema第一个 schema 字段是事件键的一部分。它指定了一个 Kafka Connect schema,用于描述事件键的 payload 部分包含的内容。换言之,第一个 schema 字段描述的是主键的结构。
2payload第一个 payload 字段是事件键的一部分。它具有前一个 schema 字段所描述的结构,其中包含被更改行的键值。
3schema第二个 schema 字段是事件值的一部分。它指定了一个 Kafka Connect schema,用于描述事件值的 payload 部分包含的内容。换言之,第二个 schema 描述的是被更改行的结构。通常,该 schema 中会包含嵌套的 schema。
4payload第二个 payload 字段是事件值的一部分。它具有前一个 schema 字段所描述的结构,其中包含被更改行的实际数据。

表 1. 变更事件基本内容概览

默认情况下,连接器会将变更事件记录流式传输到名称与事件来源表相同的主题。

从 Kafka 0.10 开始,Kafka 可以选择将事件键和值与消息创建时(由生产者记录)或由 Kafka 写入日志时的时间戳 一并记录。

变更事件键

对于给定的表,变更事件的键具有这样的结构:其中包含一个字段,对应于事件创建时该表主键中的每个列。所有键列都会被标记为非可选。

考虑在 business 数据库中定义的 users 表,以及该表的变更事件键示例。

示例表

CREATE TABLE Users (
  id INT64 NOT NULL,
  username STRING(MAX) NOT NULL,
  password STRING(MAX) NOT NULL,
  email STRING(MAX) NOT NULL)
PRIMARY KEY (id);

更改事件键示例

users 表在保持此定义期间的每个更改事件都具有相同的键结构,其 JSON 形式如下:

{
  "schema": { (1)
    "type": "struct",
    "name": "Users.Key", (2)
    "optional": false, (3)
    "fields": [ (4)
      {
        "type": "int64",
        "optional": "false",
        "field": "false"
      }
    ]
  },
  "payload": { (5)
      "id": "1"
  },
}
项目字段名称说明
1schema键中的 schema 部分指定一个 Kafka Connect 模式,用于描述键的 payload 部分的内容。
2Users.Key定义键载荷(payload)结构的模式名称。该模式描述了发生变更的表的主键结构。
3optional表示事件键的 payload 字段是否必须包含值。主键列始终是必需的。
4fields指定 payload 中预期出现的每个字段,包括每个字段的名称、类型以及是否为可选。
5payload包含生成此变更事件的行的键。在此示例中,该键包含一个值为 1 的 id 字段。

表 2. 变更事件键的说明

变更事件值

下面沿用之前用于展示变更事件键示例的同一张示例表:

CREATE TABLE Users (
  id INT64 NOT NULL,
  username STRING(MAX) NOT NULL,
  password STRING(MAX) NOT NULL,
  email STRING(MAX) NOT NULL)
PRIMARY KEY (id);

create 事件

以下示例展示了连接器为在 Users 表中创建数据的操作所生成的变更事件中的值部分。如果 Spanner 列被标记为非可选(non-optional),那么在所有插入行的变更(mutation)中都必须提供该值。Spanner 中的所有主键列都会被标记为非可选。请注意,即使某个非键列在 Spanner 中被标记为非可选,在模式(schema)中它仍会显示为可选。只有主键列会在模式中被标记为非可选。

{
    "schema": { (1)
        "type": "struct",
        "fields": [
            {
                "type": "struct",
                "fields": [
                    {
                        "type": "int32",
                        "optional": false,
                        "field": "id"
                    },
                    {
                        "type": "string",
                        "optional": true,
                        "field": "first_name"
                    },
                    {
                        "type": "string",
                        "optional": true,
                        "field": "last_name"
                    },
                    {
                        "type": "string",
                        "optional": true,
                        "field": "email"
                    }
                ],
                "optional": true,
                "name": "Users.Value", (2)
                "field": "before"
            },
            {
                "type": "struct",
                "fields": [
                    {
                        "type": "int32",
                        "optional": false,
                        "field": "id"
                    },
                    {
                        "type": "string",
                        "optional": false,
                        "field": "first_name"
                    },
                    {
                        "type": "string",
                        "optional": false,
                        "field": "last_name"
                    },
                    {
                        "type": "string",
                        "optional": false,
                        "field": "email"
                    }
                ],
                "optional": true,
                "name": "Users.Value",
                "field": "after"
            },
            {
                "type": "struct",
                "fields": [
                    {
                        "type": "string",
                        "optional": false,
                        "field": "version"
                    },
                    {
                        "type": "string",
                        "optional": false,
                        "field": "connector"
                    },
                    {
                        "type": "string",
                        "optional": false,
                        "field": "name"
                    },
                    {
                        "type": "int64",
                        "optional": false,
                        "field": "ts_ms"
                    },
                    {
                        "type": "int64",
                        "optional": false,
                        "field": "ts_us"
                    },
                    {
                        "type": "int64",
                        "optional": false,
                        "field": "ts_ns"
                    },
                    {
                        "type": "boolean",
                        "optional": true,
                        "default": false,
                        "field": "snapshot"
                    },
                    {
                        "type": "string",
                        "optional": false,
                        "field": "db"
                    },
                    {
                        "type": "string",
                        "optional": false,
                        "field": "sequence"
                    },
                    {
                        "type": "string",
                        "optional": false,
                        "field": "project_id"
                    },
                    {
                        "type": "string",
                        "optional": false,
                        "field": "instance_id"
                    },
                    {
                        "type": "string",
                        "optional": false,
                        "field": "database_id"
                    },
                    {
                        "type": "string",
                        "optional": false,
                        "field": "change_stream_name"
                    },
                    {
                        "type": "string",
                        "optional": true,
                        "field": "table"
                    }
                    {
                        "type": "string",
                        "optional": true,
                        "field": "server_transaction_id"
                    }
                    {
                        "type": "int64",
                        "optional": true,
                        "field": "low_watermark"
                    }
                    {
                        "type": "int64",
                        "optional": true,
                        "field": "read_at_timestamp"
                    }
                    {
                        "type": "int64",
                        "optional": true,
                        "field": "number_of_records_in_transaction"
                    }
                    {
                        "type": "string",
                        "optional": true,
                        "field": "transaction_tag"
                    }
                    {
                        "type": "boolean",
                        "optional": true,
                        "field": "system_transaction"
                    }
                    {
                        "type": "string",
                        "optional": true,
                        "field": "value_capture_type"
                    }
                    {
                        "type": "string",
                        "optional": true,
                        "field": "partition_token"
                    }
                    {
                        "type": "int32",
                        "optional": true,
                        "field": "mod_number"
                    }
                    {
                        "type": "boolean",
                        "optional": true,
                        "field": "is_last_record_in_transaction_in_partition"
                    }
                    {
                        "type": "int64",
                        "optional": true,
                        "field": "number_of_partitions_in_transaction"
                    }
                ],
                "optional": false,
                "name": "io.debezium.connector.spanner.Source", (3)
                "field": "source"
            },
            {
                "type": "string",
                "optional": false,
                "field": "op"
            },
            {
                "type": "int64",
                "optional": true,
                "field": "ts_ms"
            },
            {
                "type": "int64",
                "optional": true,
                "field": "ts_us"
            },
            {
                "type": "int64",
                "optional": true,
                "field": "ts_ns"
            }
        ],
        "optional": false,
        "name": "connector_name.Users.Envelope" (4)
    },
    "payload": { (5)
        "before": null, (6)
        "after": { (7)
            "id": 1,
            "first_name": "Anne",
            "last_name": "Kretchmar",
            "email": "[email protected]"
        },
        "source": { (8)
            "version": "3.6.3.Final",
            "connector": "spanner",
            "name": "spanner_connector",
            "ts_ms": 1670955531785,
            "ts_us": 1670955531785000,
            "ts_ns": 1670955531785000000,
            "snapshot": "false",
            "db": "database",
            "sequence": "1",
            "project_id": "project",
            "instance_id": "instance",
            "database_id": "database",
            "change_stream_name": "change_stream",
            "table": "Users",
            "server_transaction_id": "transaction_id",
            "low_watermark": 1670955471635,
            "read_at_timestamp": 1670955531791,
            "number_records_in_transaction": 2,
            "transaction_tag": "",
            "system_transaction": false,
            "value_capture_type": "OLD_AND_NEW_VALUES",
            "partition_token": "partition_token",
            "mod_number": 0,
            "is_last_record_in_transaction_in_partition": true,
            "number_of_partitions_in_transaction": 1
        },
        "op": "c", (9)
        "ts_ms": 1559033904863, (10)
        "ts_us": 1559033904863769, (10)
        "ts_ns": 1559033904863769841 (10)
    }
}

表 3. create 事件值字段说明

Item Field name Description
1 schema 值的 schema,用于描述值的负载(payload)结构。对于特定表,连接器生成的每个变更事件中,变更事件值的 schema 都是相同的。
2 name 在 schema 部分中,每个 name 字段都指定了值负载中某个字段所对应的 schema。
3 name io.debezium.connector.spanner.Source 是负载中 source 字段的 schema。该 schema 是 Spanner 连接器特有的,连接器对其生成的所有事件都使用它。
4 name connector_name.Users.Envelope 是负载整体结构的 schema,其中 connector_name 是连接器名称,customers 是表名。
5 payload 值的实际数据,即变更事件所提供的信息。
6 before 可选字段,用于指定事件发生前行的状态。当 op 字段为表示创建的 c(如本例所示)时,before 字段为 null,因为该变更事件针对的是新内容。
7 after 可选字段,用于指定事件发生后行的状态。在本例中,after 字段包含新行的 id、first_name、last_name 和 email 列的值。
8 source 必填字段,用于描述该事件的源元数据。此字段包含的信息可用于将该事件与其他事件进行比较,比较内容涉及事件的来源、事件发生的先后顺序,以及这些事件是否属于同一事务。源元数据包括:
- Debezium 版本
- 连接器类型和名称
- 包含新行的数据库和表
- 该事件是否属于快照的一部分
- 事务中该数据变更事件的记录序列号
- 项目 ID
- 实例 ID
- 数据库 ID
- 变更流名称
- 事务 ID
- 低水位线,表示提交时间戳早于该低水位线时间戳的所有记录均已由连接器流出传输
- 该变更在数据库中发生时的提交时间戳
- 连接器处理该变更的时间戳
- 源事务中的记录数
- 事务标签
- 该事务是否为系统事务
- 值捕获类型
- 用于查询该变更事件的源分区令牌
- 从 Spanner 接收到的原始数据变更事件中的修改编号
- 该数据变更事件是否为分区中事务内的最后一条记录
- 该事务中变更流分区的总数
9 op 必填字符串,用于描述导致连接器生成该事件的操作类型。在本例中,c 表示该操作创建了一行。有效值为:
- c = create(创建)
- u = update(更新)
- d = delete(删除)
10 ts_ms、ts_us、ts_ns

可选字段,显示连接器处理事件的时间。该时间基于运行 Kafka Connect 任务的 JVM 中的系统时钟。

在 source 对象中,ts_ms 表示数据库中发生变更的时间。通过比较 payload.source.ts_ms 的值与 payload.ts_ms 的值,可以确定源数据库更新与 Debezium 之间的时间延迟。

update(更新)事件

示例 Users 表中某次更新的变更事件值与该表的 create(创建)事件具有相同的架构。同样,事件值的 payload 也具有相同的结构。不过,在 update 事件中,事件值的 payload 包含不同的值。以下是连接器为 Users 表中的一次更新所生成的变更事件值的示例:

{
    "schema": { ... },
    "payload": {
        "before": { (1)
            "id": 1,
            "first_name": "Anne",
            "last_name": "Kretchmar",
            "email": "[email protected]"
        },
        "after": { (2)
            "id": 1,
            "first_name": "Anne Marie",
            "last_name": "Kretchmar",
            "email": "[email protected]"
        },
        "source": { (3)
            "version": "3.6.3.Final",
            "connector": "spanner",
            "name": "spanner_connector",
            "ts_ms": 1670955531785,
            "ts_us": 1670955531785000,
            "ts_ns": 1670955531785000000,
            "snapshot": "false",
            "db": "database",
            "sequence": "1",
            "project_id": "project",
            "instance_id": "instance",
            "database_id": "database",
            "change_stream_name": "change_stream",
            "table": "Users",
            "server_transaction_id": "transaction_id",
            "low_watermark": 1670955471635,
            "read_at_timestamp": 1670955531791,
            "number_records_in_transaction": 2,
            "transaction_tag": "",
            "system_transaction": false,
            "value_capture_type": "OLD_AND_NEW_VALUES",
            "partition_token": "partition_token",
            "mod_number": 0,
            "is_last_record_in_transaction_in_partition": true,
            "number_of_partitions_in_transaction": 1
        },
        "op": "u", (4)
        "ts_ms": 1465584025523,  (5)
        "ts_us": 1465584025523614,  (5)
        "ts_ns": 1465584025523614723  (5)
    }
}

表 4. update 事件值字段说明

序号 字段名 说明
1 before 可选字段,包含数据库提交前行中所有列的值。
2 after 可选字段,指定事件发生后行的状态。在此示例中,first_name 的值现在为 Anne Marie。
3 source 必填字段,描述事件的源元数据。source 字段的结构与 create 事件中的相同,但部分值有所不同。源元数据包括:
- Debezium 版本
- 连接器类型和名称
- 包含新行的数据库(即 keyspace)和表
- 该事件是否属于快照的一部分
- 事务中该数据变更事件的记录序列号
- 项目 ID
- 实例 ID
- 数据库 ID
- 变更流名称
- 事务 ID
- 低水位标记,表示提交时间戳早于低水位标记时间戳的所有记录均已被连接器流出
- 该变更在数据库中的提交时间戳
- 连接器处理该变更的时间
- 源事务中的记录数
- 事务标签
- 该事务是否为系统事务
- 值捕获类型
- 用于查询此变更事件的源分区令牌
- 从 Spanner 接收的原始数据变更事件中的修改编号
- 该数据变更事件是否为分区内事务中的最后一条记录
- 事务中的变更流分区总数
4 op 描述操作类型的必填字符串。在 update 事件值中,op 字段的值为 u,表示该行因更新而发生了变更。
5 ts_ms、ts_us、ts_ns 可选字段,显示连接器处理事件的时间。该时间基于运行 Kafka Connect 任务的 JVM 中的系统时钟。
在 source 对象中,ts_ms 表示该变更在数据库中发生的时间。通过比较 payload.source.ts_ms 和 payload.ts_ms 的值,可以确定源数据库更新与 Debezium 之间的延迟。

delete 事件

delete 变更事件中的值,其 schema 部分与同一表的 create 和 update 事件相同。示例 Users 表的 delete 事件中的 payload 部分如下所示:

{
    "schema": { ... },
    "payload": {
        "before": { (1)
            "id": 1,
            "first_name": "Anne Marie",
            "last_name": "Kretchmar",
            "email": "[email protected]"
        },
        "after": null, (2)
        "source": { (3)
            "version": "3.6.3.Final",
            "connector": "spanner",
            "name": "spanner_connector",
            "ts_ms": 1670955531785,
            "ts_us": 1670955531785000,
            "ts_ns": 1670955531785000000,
            "snapshot": "false",
            "db": "database",
            "sequence": "1",
            "project_id": "project",
            "instance_id": "instance",
            "database_id": "database",
            "change_stream_name": "change_stream",
            "table": "Users",
            "server_transaction_id": "transaction_id",
            "low_watermark": 1670955471635,
            "read_at_timestamp": 1670955531791,
            "number_records_in_transaction": 2,
            "transaction_tag": "",
            "system_transaction": false,
            "value_capture_type": "OLD_AND_NEW_VALUES",
            "partition_token": "partition_token",
            "mod_number": 0,
            "is_last_record_in_transaction_in_partition": true,
            "number_of_partitions_in_transaction": 1
        },
        "op": "d", (4)
        "ts_ms": 1465581902461, (5)
        "ts_us": 1465581902461425, (5)
        "ts_ns": 1465581902461425378 (5)
    }
}

表 5. delete 事件值字段说明

序号 字段名 描述
1 before 可选字段,用于指定事件发生之前行的状态。在 delete 事件值中,before 字段包含该行在随数据库提交被删除之前的值。
2 after 可选字段,用于指定事件发生之后行的状态。在 delete 事件值中,after 字段为 null,表示该行已不复存在。
3 source 必填字段,用于描述事件的源元数据。在 delete 事件值中,source 字段的结构与同一表的 create 和 update 事件相同,许多 source 字段的值也相同。在 delete 事件值中,ts_ms 和 lsn 字段的值以及其他值可能会发生变化,但 delete 事件值中的 source 字段提供的元数据是相同的:
- Debezium 版本
- 连接器类型和名称
- 包含该行的数据库(又称键空间)和表
- 该事件是否属于快照的一部分
- 事务中该数据变更事件的记录序列号
- 项目 ID
- 实例 ID
- 数据库 ID
- 变更流名称
- 事务 ID
- 低水位线,表示提交时间戳早于低水位线时间戳的所有记录均已由连接器流出
- 变更在数据库中发生时的提交时间戳
- 连接器处理该变更的时间
- 发起事务中的记录数
- 事务标签
- 该事务是否为系统事务
- 值捕获类型
- 用于查询此变更事件的起始分区令牌
- 从 Spanner 接收到的原始数据变更事件中的修改编号
- 该数据变更事件是否为该分区中事务的最后一条记录
- 事务中变更流分区的总数
4 op 必填字符串,用于描述操作类型。op 字段的值为 d,表示该行已被删除。
5 ts_ms、ts_us、ts_ns 可选字段,用于显示连接器处理该事件的时间。该时间基于运行 Kafka Connect 任务的 JVM 的系统时钟。
在 source 对象中,ts_ms 表示变更在数据库中发生的时间。通过比较 payload.source.ts_ms 与 payload.ts_ms 的值,可以确定源数据库更新与 Debezium 之间的延迟。

delete 变更事件记录为消费者提供其处理该行移除所需的信息。

Spanner 连接器的事件设计为与 Kafka 日志压缩 协同工作。日志压缩允许删除一些较旧的消息,前提是为每个键保留至少最近的一条消息。这使得 Kafka 能够回收存储空间,同时确保主题中包含完整的数据集,可用于重新加载基于键的状态。请注意,如果启用了低水位线,则不应启用压缩。

墓碑事件

当某行被删除时,delete 事件的值仍然可用于日志压缩,因为 Kafka 可以删除所有具有相同键的较早消息。但是,要让 Kafka 删除所有具有相同键的消息,消息值必须为 null。为了实现这一点,连接器会在 delete 事件之后跟随一个特殊的 tombstone(墓碑)事件,该事件具有相同的键,但值为 null。

数据类型映射

Spanner 连接器使用结构与行所在表相同的事件来表示行的变更。事件为每个列值包含一个字段。该值在事件中的表示方式取决于该列的 Spanner 数据类型。本节描述这些映射。

Spanner 类型

Spanner 数据类型字面量类型(架构类型)
BOOLEANBOOLEAN
INT64INT64
ARRAYARRAY
BYTESBYTES
STRINGSTRING
FLOAT32FLOAT
FLOAT64DOUBLE
NUMERICSTRING
TIMESTAMPSTRING
NUMERICSTRING

表 6. Spanner 数据类型的映射

设置 Spanner

检查清单

  • 请确保提供项目 ID、Spanner 实例 ID、Spanner 数据库 ID 和更改流名称。有关如何创建更改流,请参阅文档。
  • 请确保创建并配置具有适当凭据的 GCP 服务帐号。可以在连接器配置中显式提供服务帐号密钥或通过路径提供,但默认情况下连接器使用应用默认凭据 (ADC)。有关服务帐号的更多信息,请参阅此处的信息。

请参阅以下章节,了解如何配置 Debezium Spanner 连接器。

部署

安装好 Kafka 和 Kafka Connect 之后,部署 Debezium Spanner 连接器的剩余工作包括:下载连接器插件归档包,将 JAR 文件解压到 Kafka Connect 环境中,并将包含这些 JAR 文件的目录添加到 Kafka Connect 的 plugin.path 中。随后需要重启 Kafka Connect 进程,以便加载新的 JAR 文件。

如果你使用的是不可变容器,可以查看 Debezium 的容器镜像,其中提供了已预装 Spanner 连接器并可直接运行的 Kafka 和 Kafka Connect。你也可以在 Kubernetes 和 OpenShift 上运行 Debezium。

你从 quay.io 获取的 Debezium 容器镜像未经过严格测试或安全分析,仅供测试和评估之用,不建议在生产环境中使用。为降低生产部署的风险,请仅部署由可信供应商积极维护并经过潜在漏洞充分测试的容器。

连接器配置示例

以下是一个 Spanner 连接器的配置示例,该连接器连接到实例 Instance、项目 Project 中数据库 Database 里名为 changeStreamAll 的变更流。

你可以选择仅为部分模式和表生成事件。还可以按需忽略、掩码或截断敏感的、过大的或不需要的列。

{
  "name": "spanner-connector",  (1)
  "config": {
    "connector.class": "io.debezium.connector.spanner.SpannerConnector", (2)
    "gcp.spanner.change.stream": "changeStreamAll", (3)
    "gcp.spanner.project.id": "Project", (4)
    "gcp.spanner.instance.id": "Instance", (5)
    "gcp.spanner.database.id": "Database", (6)
    "gcp.spanner.credentials.json": <key.json>, (7)
    "tasks.max": 1 (8)
  }
}
1连接器注册到 Kafka Connect 服务时所使用的名称。
2此 Spanner 连接器类的名称。
3变更流的名称。
4GCP 项目 ID。
5Spanner 实例 ID。
6Spanner 数据库 ID。
7GCP 服务账号密钥 JSON。
8最大任务数。

有关这些配置中可以指定的Spanner 连接器属性完整列表,请参阅该文档。

添加连接器配置

要启动 Spanner 连接器,请创建连接器配置并将该配置添加到你的 Kafka Connect 集群。

前提条件

  • Spanner 变更流已创建并可用。
  • Spanner 连接器已安装。

操作步骤

  1. 为 Spanner 连接器创建配置。
  2. 使用 Kafka Connect REST API 将该连接器配置添加到你的 Kafka Connect 集群。

结果

连接器启动后,会开始为行级操作生成数据变更事件,并将变更事件记录流式传输到 Kafka 主题。

监控

除了 Kafka 和 Kafka Connect 内置的 JMX 指标支持外,Debezium Spanner 连接器仅提供一种类型的指标。

  • 流式传输指标提供连接器在捕获变更并流式传输变更事件记录期间的运行信息。

Debezium 监控文档详细介绍了如何通过 JMX 暴露这些指标。

自定义 MBean 名称

Debezium 连接器通过连接器的 MBean 名称来暴露指标。这些指标针对每个连接器实例,提供有关连接器快照、流式传输和架构历史记录过程运行状况的数据。

默认情况下,当你部署一个配置正确的连接器时,Debezium 会为每种不同的连接器指标生成唯一的 MBean 名称。要查看某个连接器过程的指标,你需要配置可观测性栈以监控其 MBean。但这些默认 MBean 名称依赖于连接器配置;配置的更改可能导致 MBean 名称发生变化。MBean 名称的变化会破坏连接器实例与 MBean 之间的关联,从而干扰监控活动。在这种情况下,如果要恢复监控,必须重新配置可观测性栈以使用新的 MBean 名称。

为防止因 MBean 名称变更而导致的监控中断,你可以配置自定义指标标签。通过在连接器配置中添加 custom.metric.tags 属性即可配置自定义指标。该属性接受键值对,其中每个键代表 MBean 对象名称的一个标签,对应的值则是该标签的取值。例如:k1=v1,k2=v2。Debezium 会将指定的标签追加到连接器的 MBean 名称之后。

为连接器配置 custom.metric.tags 属性后,你就可以配置可观测性栈来获取与指定标签关联的指标。此后,可观测性栈将使用指定的标签(而非可变的 MBean 名称)来唯一标识连接器。将来,如果 Debezium 重新定义了 MBean 名称的构建方式,或者连接器配置中的 topic.prefix 发生变化,指标采集也不会中断,因为指标抓取任务会使用指定的标签模式来识别连接器。

使用自定义标签的另一个好处是,你可以使用能够反映数据流水线架构的标签,从而以符合自身运维需求的方式组织指标。例如,你可以指定取值用于声明连接器活动类型、应用程序上下文或数据源的标签,如 db1-streaming-for-application-abc。如果你指定了多个键值对,所有指定的键值对都会被追加到连接器的 MBean 名称之后。

以下示例说明了标签如何修改默认的 MBean 名称。

示例 1. 自定义标签如何修改连接器 MBean 名称

默认情况下,Spanner 连接器为流式传输指标使用以下 MBean 名称:

debezium.Cloud Spanner:type=connector-metrics,context=streaming,server=<topic.prefix>

如果你将 custom.metric.tags 的值设置为 database=salesdb-streaming,table=inventory,Debezium 会生成如下的自定义 MBean 名称:

debezium.Cloud Spanner:type=connector-metrics,context=streaming,server=<topic.prefix>,database=salesdb-streaming,table=inventory

流式处理指标

MBean 为 debezium.Cloud Spanner:type=connector-metrics,context=streaming,server=<topic.prefix>。

AttributesTypeDescription
MilliSecondsSinceLastEventlong连接器读取并处理最近一次事件以来经过的毫秒数。
TotalNumberOfEventsSeenlong自上次启动或重置以来,该连接器看到的事件总数。
NumberOfEventsFilteredlong被连接器上配置的包含/排除列表过滤规则所过滤掉的事件数量。
QueueTotalCapacityint用于在流处理器与主 Kafka Connect 循环之间传递事件的队列长度。
QueueRemainingCapacityint用于在流处理器与主 Kafka Connect 循环之间传递事件的队列的剩余容量。
Connectedboolean指示连接器当前是否已连接到数据库服务器的标志。
MilliSecondsBehindSourcelong最后一个更改事件的时间戳与连接器处理该事件之间相差的毫秒数。该值会包含运行数据库服务器和连接器的机器之间时钟的任何差异。
NumberOfCommittedTransactionslong已处理且已提交的事务数量。
MaxQueueSizeInByteslong用于在流处理器与主 Kafka Connect 循环之间传递事件的队列的最大缓冲字节数。
CurrentQueueSizeInByteslong用于在流处理器与主 Kafka Connect 循环之间传递事件的队列的当前缓冲字节数。
LowWatermarklong连接器任务的当前低水位标记。低水位标记描述的是时间 T,连接器保证在此时已流式传输所有时间戳 < T 的事件。
MilliSecondsLowWatermarklong连接器任务的当前低水位标记(以毫秒为单位)。低水位标记描述的是时间 T,连接器保证在此时已流式传输所有时间戳 < T 的事件。
MilliSecondsLowWatermarkLaglong低水位标记落后于当前时间的毫秒数。低水位标记描述的是时间 T,连接器保证在此时已流式传输所有时间戳 < T 的事件。
LatencyLowWatermarkLag<variant>MilliSecondslong低水位标记落后于当前时间的毫秒数分布。该变体将包含 P50、P95、P99、平均值、最小值、最大值的计算。
LatencySpanner<variant>MilliSecondslongSpanner 提交时间戳到连接器读取延迟的分布。该变体将包含 P50、P95、P99、平均值、最小值、最大值的计算。
LatencyReadToEmit<variant>MilliSecondslongSpanner 读取时间戳到连接器发出延迟的分布。该变体将包含 P50、P95、P99、平均值、最小值、最大值的计算。
LatencyCommitToEmit<Variant>MilliSecondslongSpanner 提交时间戳到连接器发出延迟的分布。该变体将包含 P50、P95、P99、平均值、最小值、最大值的计算。
LatencyCommitToPublish<variant>MilliSecondslongSpanner 提交时间戳到 Kafka 发布时间戳的延迟分布。该变体将包含 P50、P95、P99、平均值、最小值、最大值的计算。
LatencyEmitToPublish<variant>MilliSecondslong连接器发出时间戳到 Kafka 发布时间戳的延迟分布。该变体将包含 P50、P95、P99、平均值、最小值、最大值的计算。
SpannerEventQueueCapacitylongSpanner 事件队列的总容量。该队列表示 StreamEventQueue 的总容量,这是一个 Spanner 特有的队列,用于存储从变更流查询接收到的元素。
RemainingSpannerEventQueueCapacitylongSpanner 事件队列的剩余容量。
TaskStateChangeEventQueueCapacitylong任务状态更改事件队列的总容量。该队列表示 TaskStateChangeEventQueue 的总容量,这是一个 Spanner 特有的队列,用于存储连接器中发生的事件。
RemainingTaskStateChangeEventQueueCapacitylong任务状态更改事件队列的剩余容量。
NumberOfChangeStreamPartitionsDetectedlong当前任务检测到的分区总数。
NumberOfChangeStreamQueriesIssuedlong当前任务发出的变更流查询总数。
NumberOfActiveChangeStreamQuerieslong当前任务检测到的活动变更流查询数量。

连接器配置属性

Debezium Spanner 连接器提供了许多配置属性,你可以通过它们调整连接器以满足应用的需要。许多属性都有默认值。属性信息组织如下:

除非提供了默认值,否则下列配置属性均为必填。

表 7. 连接器必填配置属性

属性 默认值 描述

name

无默认值

连接器的唯一名称。使用相同名称再次注册会失败。所有 Kafka Connect 连接器都要求此属性。

connector.class

无默认值

连接器对应的 Java 类名。对于 Spanner 连接器,始终使用 io.debezium.connector.spanner.SpannerConnector 这个值。

tasks.max

1

应为该连接器创建的最大任务数。如果启用了 offset.storage.per.task 模式,Spanner 连接器可以使用超过 1 个任务。

gcp.spanner.project.id

无默认值

GCP 项目 ID

gcp.spanner.instance.id

无默认值

Spanner 实例 ID

gcp.spanner.database.id

无默认值

Spanner 数据库 ID

gcp.spanner.change.stream

无默认值

Spanner 变更流

gcp.spanner.credentials.path

无默认值

GCP 服务账号密钥 JSON 文件的路径。

gcp.spanner.credentials.json

无默认值

GCP 服务账号密钥 JSON 内容。如果未提供 gcp.spanner.credentials.path 且不存在应用默认凭据,则必须提供此属性。

schema.name.adjustment.mode

none

指定为兼容连接器使用的消息转换器,应如何调整架构名称。可选设置:

  • none 不进行任何调整。
  • avro 将 Avro 类型名称中无法使用的字符替换为下划线。
  • avro_unicode 将下划线或 Avro 类型名称中无法使用的字符替换为对应的 unicode,例如 _uxxxx。注意:_ 是 Java 中类似于反斜杠的转义序列

field.name.adjustment.mode

none

指定字段名称应如何调整,以便与连接器所使用的消息转换器兼容。可选设置:

  • none 不进行任何调整。
  • avro 将 Avro 类型名称中无法使用的字符替换为下划线。
  • avro_unicode 将下划线或 Avro 类型名称中无法使用的字符替换为对应的 unicode,例如 _uxxxx。注意:_ 是 Java 中类似于反斜杠的转义序列

更多详细信息请参阅 Avro 命名。

以下高级配置属性的默认值适用于大多数情况,因此很少需要在连接器配置中指定。

表 8. 高级连接器配置属性

属性 默认值 说明

gcp.spanner.low-watermark.enabled

false

是否为连接器启用低水位标记(low watermark)。

gcp.spanner.low-watermark.update-period.ms

1000 ms

低水位标记的更新间隔。

heartbeat.interval.ms

300000

Spanner 心跳间隔。

gcp.spanner.start.time

current time

连接器开始时间。

gcp.spanner.end.time

indefinite end time

连接器结束时间。

gcp.spanner.stream.event.queue.capacity

10000

Spanner 事件队列容量。如果在连接器运行期间,流事件队列的剩余容量接近零,请增大此容量。

connector.spanner.task.state.change.event.queue.capacity

1000

任务状态变更事件队列容量。如果在连接器运行期间,任务状态变更事件队列的剩余容量接近零,请增大此容量。

connector.spanner.max.missed.heartbeats

5

变更流查询在抛出异常之前允许丢失的最大心跳数

scaler.monitor.enabled

false

是否启用任务自动扩缩容

connector.spanner.sync.topic

sync_topic_spanner_connector<connectorname>

同步主题(Sync topic)的名称。同步主题是连接器的内部主题,用于存储任务之间的通信。

connector.spanner.sync.poll.duration

500 ms

同步主题的轮询时长。

connector.spanner.sync.request.timeout.ms

5000 ms

向同步主题发出请求的超时时间。

connector.spanner.sync.delivery.timeout.ms

15000 ms

向同步主题发布消息的超时时间。

connector.spanner.sync.commit.offsets.timeout.ms

5000 ms

提交同步主题偏移量的超时时间。

connector.spanner.sync.commit.offsets.interval.ms

60000 ms

提交同步主题偏移量的时间间隔。

connector.spanner.sync.publisher.wait.timeout

5 ms

向同步主题发布消息的时间间隔。

connector.spanner.rebalancing.topic

rebalancing_topic_spanner_connector<connectorname>

再均衡主题(rebalancing topic)的名称。再均衡主题是连接器的内部主题,用于判定任务是否存活。

connector.spanner.rebalancing.poll.duration

5000

再均衡主题的轮询时长。

connector.spanner.rebalancing.commit.offsets.timeout

5000

提交再均衡主题偏移量的超时时间。

connector.spanner.rebalancing.commit.offsets.interval.ms

60000 ms

提交同步主题偏移量的时间间隔。

connector.spanner.rebalancing.task.waiting.timeout

1000 ms

任务在处理重新平衡事件之前等待的时间长度。

custom.metric.tags

无默认值

定义用于自定义 MBean 对象名称的标签,这些标签会添加元数据以提供上下文信息。请指定以逗号分隔的键值对列表。每个键代表 MBean 对象名称的一个标签,相应的值代表该键的取值,例如 k1=v1,k2=v2

连接器会将指定的标签附加到基础 MBean 对象名称之后。标签有助于你组织和分类指标数据。你可以定义标签来标识特定的应用实例、环境、区域、版本等。有关更多信息,请参阅 自定义 MBean 名称。

custom.sanitize.pattern

.*secret$|.*password$|.*sasl\.jaas\.config$|.*basic\.auth\.user\.info$|.*registry\.auth\.client-secret

一个可选的正则表达式,用于指定屏蔽敏感配置键的自定义模式。

默认情况下,Debezium 会屏蔽与预定义模式匹配的配置键的值,这些模式针对密码和身份验证令牌等已知的常见敏感属性。若要自定义连接器屏蔽配置键值的方式,请将此属性设置为能够匹配你想要屏蔽的特定配置键值的正则表达式。

例如,若要屏蔽包含字符串 api.key 或 token 的键,请将此属性设置为以下值:

"custom.sanitize.pattern": ".*api\\.key.*\|.*token.*"

如果为连接器设置此属性,Debezium 会在显示、日志记录或响应 API 调用时,对其返回的与该配置键值匹配的实例进行脱敏处理。原始值将被一串星号替换,例如 ********。

你为 custom.sanitize.pattern 属性指定的模式会覆盖 Debezium 默认的脱敏模式;指定的值不会扩展或增强默认模式。

errors.max.retries

-1

指定当操作出现可重试错误(例如连接错误)后,连接器应如何响应。
请设置以下选项之一:

-1

无限制。无论之前失败了多少次,连接器总是自动重启并重试该操作。

0

已禁用。连接器立即失败,且永不重试该操作。需要人工干预才能重启连接器。

> 0

连接器会自动重启,直到达到指定的最大重试次数。在下一次失败后,连接器将停止,需要人工干预才能重启。

extended.headers.enabled

true

此属性指定 Debezium 是否为其发出的消息添加带有前缀 __debezium.context. 的上下文标头。

OpenLineage 集成需要这些标头,它们提供的元数据使下游处理系统能够跟踪并识别变更事件的来源。

此属性会添加以下标头:

__debezium.context.connectorLogicalName

Debezium 连接器的逻辑名称。

__debezium.context.taskId

连接器任务的唯一标识符。

__debezium.context.connectorName

Debezium 连接器的名称。

有关更完整的高级配置列表,请参阅 Github 代码。

直通连接器配置属性

连接器还支持直通配置属性,这些属性用于创建 Kafka 生产者和消费者。

请务必查阅 Kafka 文档,了解 Kafka 生产者和消费者的所有配置属性。Spanner 连接器使用的是新版消费者配置属性。

出现问题时的行为

Debezium 是一个分布式系统,它捕获多个上游数据库中的所有更改;它绝不会遗漏或丢失事件。当系统正常运行或受到精心管理时,Debezium 会为每个更改事件记录提供 exactly once(精确一次)的交付保证。

如果确实发生了故障,系统不会丢失任何事件。但是,在从故障恢复的过程中,可能会重复某些更改事件。在这种异常情况下,Debezium 与 Kafka 一样,会为更改事件提供 at least once(至少一次)的交付保证。

本节余下部分将描述 Debezium 如何处理各种故障和问题。

配置和启动错误

在以下情况下,连接器在尝试启动时会失败,在日志中报告错误/异常并停止运行:

  • 连接器的配置无效。
  • 连接器无法使用指定的连接参数成功连接到 Spanner。

在这些情况下,错误消息中包含有关问题的详细信息,可能还会提供建议的解决方法。更正配置或解决 Spanner 的问题后,重新启动连接器。

Spanner 变得不可用

当连接器正在运行时,Spanner 可能会由于各种原因变得不可用。连接器将继续运行,并在 Spanner 再次可用后能够继续流式传输事件。

Kafka Connect 进程正常停止

假设 Kafka Connect 以分布式模式运行,其中一个 Kafka Connect 进程被正常停止。在关闭该进程之前,Kafka Connect 会将该进程的连接器任务迁移到同一组中的另一个 Kafka Connect 进程。新的连接器任务从之前任务停止的位置精确地继续处理。在连接器任务正常停止并在新进程上重新启动期间,处理会有短暂的延迟。

Kafka Connect 进程崩溃

如果 Kafka Connector 进程意外停止,它正在运行的任何连接器任务都会终止,且不会记录其最近处理的偏移量。当 Kafka Connect 以分布式模式运行时,Kafka Connect 会在其他进程上重新启动这些连接器任务。但是,Spanner 连接器会从先前进程 已记录 的最后一个偏移量处恢复。这意味着新的替代任务可能会重新生成崩溃前刚刚处理过的部分更改事件。重复事件的数量取决于偏移量刷新周期以及崩溃前的数据更改量。

由于在故障恢复过程中某些事件可能会重复,消费者应始终预期会出现一些重复事件。

在每个更改事件记录中,Debezium 连接器会插入有关事件来源的特定源信息,例如原始分区令牌、提交时间戳、事务 ID、记录序列和修改编号。消费者可以使用这些标识符进行去重。

Kafka 变得不可用

连接器生成变更事件时,Kafka Connect 框架会使用 Kafka 生产者 API 将这些事件记录到 Kafka 中。Kafka Connect 会按照你在其配置中指定的频率,定期记录这些变更事件中出现的最新偏移量。如果 Kafka 代理不可用,正在运行连接器的 Kafka Connect 进程会反复尝试重新连接 Kafka 代理。换句话说,连接器任务会暂停,直到能够重新建立连接为止,届时连接器将从上次中断的位置精确恢复。

连接器停止一段时间

如果连接器被正常停止,数据库仍可继续使用。连接器重新启动后,它会从中断处继续流式传输变更。也就是说,它会为连接器停止期间发生的所有数据库变更生成变更事件记录。

配置得当的 Kafka 集群能够处理海量吞吐。Kafka Connect 是按照 Kafka 的最佳实践编写的,在资源充足的情况下,Kafka Connect 连接器同样可以处理极大量的数据库变更事件。因此,连接器停止一段时间后,Debezium 连接器重启时,很可能能够追上其停止期间发生的数据库变更。追上的速度取决于 Kafka 的能力和性能,以及 Spanner 中数据变更的数量。

请注意,当前连接器只能回溯约一小时的时间。Kafka 连接器会在其启动时间戳处读取 information schema 以获取架构信息。默认情况下,Spanner 无法在早于版本保留期的读取时间戳上读取 information schema,而版本保留期默认为一小时。如果你想让连接器从早于一小时前的时间点开始运行,就需要增加数据库的版本保留期。

限制

评论

登录后参与评论

正在加载评论…