源连接器

Vitess

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

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

Debezium Vitess 连接器

Debezium 的 Vitess 连接器捕获 Vitess 键空间(keyspace)中分片的行级变更。有关与此连接器兼容的 Vitess 版本信息,请参阅 Debezium 发布概览。

连接器首次连接到 Vitess 集群时,会对键空间进行一致性快照。快照完成后,连接器会持续捕获提交到 Vitess 键空间的行级变更,从而对数据库内容执行插入、更新或删除操作。连接器生成数据变更事件记录,并将其流式传输到 Kafka 主题中。对于每张表,默认行为是连接器将该表生成的所有事件流式传输到该表专属的 Kafka 主题中。应用和服务随后可以从生成的主题中消费数据变更事件记录。

概述

Vitess 的 VStream 功能是在 4.0 版本中引入的。它是一项变更事件订阅服务,提供与底层 MySQL 分片的 MySQL 二进制日志等效的信息。用户可以订阅一个键空间中的多个分片,这使其成为为下游 CDC 流程提供数据的便捷工具。

为了读取和处理数据库变更,Vitess 连接器会订阅 VTGate 的 VStream gRPC 服务。VTGate 是一个轻量级的无状态 gRPC 服务器,是 Vitess 集群架构的一部分。

连接器为用户提供了灵活性,可以选择订阅 MASTER 节点或 REPLICA 节点来获取变更事件。

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

Vitess 底层的 MySQL 实现会基于可配置的时间周期清理二进制日志。由于 binlog 的内容可能不完整,连接器需要另一种机制来确保其捕获特定数据库的完整内容。因此,当连接器首次连接到数据库时,会对该数据库执行一致性快照。连接器完成快照后,会从快照的确切位置继续流式传输变更。通过这种方式,连接器以所有数据的一致视图作为起点,并且不会遗漏快照期间所做发生的任何变更。

该连接器具备故障容错能力。当连接器读取变更并生成事件时,会记录每个事件对应的 Vitess 全局事务 ID(VGTID)位置。如果连接器因任何原因停止运行(包括通信故障、网络问题或崩溃),在连接器重启后,它会从上次存储的最后一条变更事件位置继续通过 VStream 读取。此行为不适用于快照。如果连接器在快照过程中停止,重启后连接器不会从上次中断处继续执行快照。我们稍后会讨论出现问题时连接器的行为。

连接器的工作原理

要优化配置和运行 Debezium Vitess 连接器,有必要了解该连接器如何执行快照、流式传输变更事件、确定 Kafka 主题名称以及使用元数据。

快照

通常,MySQL 服务器不会被配置为在二进制日志中保留完整的数据库历史。因此,连接器无法从二进制日志中读取整个数据库历史。基于此原因,连接器首次启动时会对数据库执行初始一致性快照。你可以通过将 snapshot.mode 连接器配置属性设置为 initial 以外的值来更改此行为。此快照功能基于 7.0 版本引入的 VStream Copy 构建。

+ 从空 GTID 启动时,可以通过设置 snapshot.include.collection.list 连接器配置属性来指定以逗号分隔的表模式列表,从而控制在快照阶段复制哪些表。如果未设置该属性或该属性为空,则在快照阶段会复制所有表。

预计将在未来版本中提供对失败快照的自动重试功能。

流式传输变更

Vitess 连接器将全部时间用于从其所订阅的 VTGate 的 VStream gRPC 服务流式传输变更。客户端会以底层 MySQL 服务器 binlog 中提交变更时的特定位置(称为 VGTID)从 VStream 接收变更。

Vitess 中的 VGTID 相当于 MySQL 中的 GTID,它描述了变更事件在 VStream 中发生的位置。通常,一个 VGTID 包含多个分片 GTID,每个分片 GTID 是一个 (Keyspace, Shard, GTID) 元组,用于描述给定分片的 GTID 位置。

订阅 VStream 服务时,连接器需要提供一个 VGTID 和一个 Tablet 类型(例如 MASTER、REPLICA)。VGTID 描述了 VStream 应从哪个位置开始发送变更事件;Tablet 类型描述了我们从每个分片中的哪个底层 MySQL 实例(主库或副本)读取变更事件。

连接器首次连接到 Vitess 集群时,会从名为 VTCtld 的 Vitess 组件获取当前的 VGTID,并将该 VGTID 提供给 VStream。

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

Kafka Connect 会定期将最新的 offset(偏移量)记录到另一个 Kafka 主题中。该偏移量表示 Debezium 随每个事件一起包含的、与源相关的位点信息。对于 Vitess 连接器而言,每条变更事件中记录的 VGTID 就是偏移量。

当 Kafka Connect 正常关闭时,它会停止连接器,将所有事件记录刷新到 Kafka,并记录从每个连接器接收到的最后一个偏移量。当 Kafka Connect 重新启动时,它会读取每个连接器最后记录的偏移量,并从各自的最后记录偏移量处启动每个连接器。连接器重新启动时,会向 VStream 发送请求,从该位置之后开始发送事件。

主题名称

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

  • topicPrefix 是由 topic.prefix 连接器配置属性指定的主题前缀。
  • keyspaceName 是发生该操作的键空间(即数据库)的名称。
  • tableName 是发生该操作的数据库表的名称。

例如,假设 fulfillment 是某个连接器配置中的逻辑服务器名,该连接器用于捕获某个 Vitess 安装中的变更,该安装包含一个名为 commerce 的键空间,其中有四张表:products、products_on_hand、customers 和 orders。无论该键空间有多少个分片,连接器都会将记录流式传输到这四个 Kafka 主题:

  • fulfillment.commerce.products
  • fulfillment.commerce.products_on_hand
  • fulfillment.commerce.customers
  • fulfillment.commerce.orders

事务元数据

Debezium 可以生成表示事务边界的事件,并用于丰富数据变更事件消息。

Debezium 接收事务元数据的限制

Debezium 只对部署连接器之后发生的事务注册并接收元数据。部署连接器之前发生的事务,其元数据不可用。

Debezium 会为每个事务中的 BEGIN 和 END 分隔符生成事务边界事件。事务边界事件包含以下字段:

status

BEGIN 或 END。

id

唯一事务标识符的字符串表示形式。

ts_ms

事务边界事件(BEGIN 或 END 事件)在数据源处发生的时间。如果数据源没有向 Debezium 提供事件时间,则该字段表示 Debezium 处理该事件的时间。注意:ts_ms 的单位是毫秒,但由于 MySQL 的限制,其精度只能达到秒级,因为 MySQL 只能提供秒级精度的 binlog 时间戳。

event_count(针对 END 事件)

该事务产生的事件总数。

data_collections(针对 END 事件)

由 data_collection 和 event_count 元素对组成的数组,用于表示连接器针对源自某个数据集合的变更所发出的事件数量。

示例

{
  "status": "BEGIN",
  "id": "[{\"keyspace\":\"test_unsharded_keyspace\",\"shard\":\"0\",\"gtid\":\"MySQL56/e03ece6c-4c04-11ec-8e20-0242ac110004:1-37\"}]",
  "ts_ms": 1486500577000,
  "event_count": null,
  "data_collections": null
}

{
  "status": "END",
  "id": "[{\"keyspace\":\"test_unsharded_keyspace\",\"shard\":\"0\",\"gtid\":\"MySQL56/e03ece6c-4c04-11ec-8e20-0242ac110004:1-37\"}]",
  "ts_ms": 1486500577000,
  "event_count": 1,
  "data_collections": [
    {
      "data_collection": "test_unsharded_keyspace.my_seq",
      "event_count": 1
    }
  ]
}

除非通过 topic.transaction 选项被覆盖,否则连接器会将事务事件发送到 <topic.prefix>.transaction 主题。

变更数据事件的元数据增强

启用事务元数据后,数据消息的 Envelope 会新增一个 transaction 字段进行增强。该字段以复合字段的形式提供每个事件的相关信息:

  • id —— 唯一事务标识符的字符串表示
  • total_order —— 该事件在事务生成的所有事件中的绝对位置
  • data_collection_order —— 该事件在事务发出的所有事件中,针对单个数据集合的位置

下面是一个消息示例:

{
  "before": null,
  "after": {
    "pk": "2",
    "aa": "1"
  },
  "source": {
...
  },
  "op": "c",
  "ts_ms": 1637988245467,
  "ts_us": 1637988245467841,
  "ts_ns": 1637988245467841698,
  "transaction": {
    "id": "[{\"keyspace\":\"test_unsharded_keyspace\",\"shard\":\"0\",\"gtid\":\"MySQL56/e03ece6c-4c04-11ec-8e20-0242ac110004:1-68\"}]",
    "total_order": 1,
    "data_collection_order": 1
  }
}

有序事务元数据

你可以配置 Debezium,使其在数据变更事件记录中包含额外的元数据。当重新分区或其他干扰可能导致数据被乱序消费时,这类补充元数据可以帮助下游消费者以正确的顺序处理消息。

变更数据富化

要配置连接器以输出经过富化的数据变更事件记录,请设置 transaction.metadata.factory 属性。当该属性设置为 VitessOrderedTransactionMetadataFactory 时,连接器会在消息的 Envelope 中包含一个 transaction 字段。transaction 字段会添加提供事务发生顺序信息的元数据。每条消息都会添加以下字段:

transaction_epoch

一个非递减的值,表示该事务排名所属的纪元(epoch)。

transaction_rank

在纪元内非递减的值,表示事务的顺序。

还有第三个字段与事件排序相关:

total_order

表示事件在事务生成的所有事件中的绝对位置。标准事务元数据默认包含该字段。

以下示例说明了如何使用这些字段来确定事件顺序。

假设 Debezium 为发生在同一分片(shard)中且共享相同主键的两个事件发出变更事件记录。如果这些事件被发送到的 Kafka 主题发生了重新分区,那么消费者对这两个事件的读取顺序就无法保证。如果 Debezium 被配置为提供富化的事务元数据,那么从该主题消费的应用程序可以应用以下逻辑来决定两个事件中哪一个应当被应用(较新的事件),哪一个应当被丢弃:

  1. 如果 transaction_epoch 的值不相等,则返回 transaction_epoch 值较大的事件;否则继续。
  2. 如果 transaction_rank 的值不相等,则返回 transaction_rank 值较大的事件;否则继续。
  3. 返回 total_order 值较大的事件。

如果两个事件的 total_order 值都不比对方大,则说明这两个事件属于同一个事务。由于 total_order 字段表示事件在事务内的顺序,因此该值较大的事件就是最新的事件。

以下示例展示了一个带有有序事务元数据的数据变更事件:

{
  "before": null,
  "after": {
    "pk": "2",
    "aa": "1"
  },
  "source": {
...
  },
  "op": "c",
  "ts_ms": 1637988245467,
  "ts_us": 1637988245467841,
  "ts_ns": 1637988245467841698,
  "transaction": {
    "id": "[{\"keyspace\":\"test_unsharded_keyspace\",\"shard\":\"0\",\"gtid\":\"MySQL56/e03ece6c-4c04-11ec-8e20-0242ac110004:1-68\"}]",
    "total_order": 1,
    "data_collection_order": 1,
    "transaction_rank": 68,
    "transaction_epoch": 0
  }
}

连接器代次(Connector Generation)

Debezium Vitess 连接器提供了 connector generation 属性,用于存储事务排序语义的代次编号。

某些连接器配置的更改会影响连接器计算 transaction-rank 字段的方式。做出此类更改后,必须递增 vitess.connector.generation 属性的值。递增连接器代次后,连接器便能获得必要的信息,从而以正确的顺序从不同的分片(shard)向 Kafka 流式传输变更事件。

例如,当你设置 vitess.transaction.chunk.size.bytes 以启用事务分块时,连接器会根据上一个事务的 VGTID 而非当前事务的 VGTID 来计算 transaction_rank。之所以会出现这种行为,是因为对于分块事务,VStream 只有在提交事务的最后一个分块之后,才会发送该事务最终已提交的 VGTID。因此,在启用事务分块之前和之后产生的事件会使用不同的机制来计算 rank。启用事务分块后,连接器计算 rank 所使用的机制与禁用分块时不同。如果你没有递增连接器代次(epoch),处理这些由不同排序机制产生的事件的消费者可能会错误地排列事件顺序。

连接器如何应用代次变更

连接器启动时,会将偏移量(offset)中存储的代次值与 vitess.connector.generation 属性中配置的值进行比较。如果两个代次值不同——无论配置的代次值高于还是低于存储值——连接器都会在开始流式传输之前,将所有分片的 transaction_epoch 递增 1。在配置变更后递增事务代次,连接器随后发出的事件所关联的代次将晚于变更前发出的事件。代次作为一个边界标记,表明连接器配置发生了变化,使下游事件处理框架能够可靠地区分跨越该转换边界发生的变更事件。由于代次是有序事务元数据比较逻辑中的主要排序键,无论 transaction_rank 值如何,消费者都能正确地将变更后的事件识别为更新的事件。

每当做出以下任何一种影响事务排序语义的配置更改时,都应递增 vitess.connector.generation:

  • 启用或禁用事务分块(vitess.transaction.chunk.size.bytes)。
  • 任何其他会改变 transaction_rank 计算方式的更改。

高效的事务元数据

如果启用连接器以提供事务元数据,它会生成明显更多的数据。连接器不仅会向事务主题发送额外的消息,而且它发送到数据变更主题的消息也会更大,因为这些消息包含一个事务元数据块。数据量增加是由以下因素造成的:

  • VGTID 被存储了两次,一次作为 source.vgtid,另一次作为 transaction.id。在包含大量分片的 keyspace 中,这些 VGTID 可能相当大
  • 在分片环境中,VGTID 通常包含每个分片的 VGTID。在拥有大量分片的 keyspace 中,VGTID 字段中的数据量可能相当大。
  • 连接器会为每个事务边界事件发送事务主题消息。通常,包含大量分片的 keyspace 往往会产生大量的事务边界事件。

为了使 Vitess 连接器能够在不显著增加产出数据量的情况下编码事务元数据,Debezium 提供了多个单消息转换(SMT)。以下 SMT 旨在减少 Vitess 连接器输出事件中的数据量:

下面的示例展示了使用上述转换的 Vitess 连接器配置的片段:

}
 [...]
 "provide.transaction.metadata": "true",
 "transaction.metadata.factory": "io.debezium.connector.vitess.pipeline.txmetadata.VitessOrderedTransactionMetadataFactory",
 "transforms": "filterTransactionTopicRecords,removeField,useLocalVgtid",
 "transforms.filterTransactionTopicRecords.type": "io.debezium.connector.vitess.transforms.FilterTransactionTopicRecords",
 "transforms.removeField.type": "io.debezium.connector.vitess.transforms.RemoveField",
 "transforms.removeField.field_names": "transaction.id",
 "transforms.useLocalVgtid.type": "io.debezium.connector.vitess.transforms.UseLocalVgtid",
 [...]
}

数据变更事件

Debezium Vitess 连接器为每个行级的 INSERT、UPDATE 和 DELETE 操作生成一个数据变更事件。每个事件都包含一个键和一个值。键和值的结构取决于发生变更的表。

Debezium 和 Kafka Connect 是围绕事件消息的连续流设计的。然而,这些事件的结构可能会随时间发生变化,这对消费者来说可能难以处理。为了解决这个问题,每个事件都包含其内容的 schema,或者(如果你使用的是 schema 注册表)包含一个 schema ID,消费者可以使用它从注册表中获取 schema。这样每个事件都是自包含的。

下面的 JSON 骨架展示了变更事件的四个基本部分。不过,如何表示变更事件中的这四个部分,取决于你在应用中配置的 Kafka Connect 转换器。只有当你配置转换器生成 schema 字段时,变更事件中才会包含该字段。同样,只有当你配置转换器生成事件键和事件负载时,变更事件中才会包含它们。如果你使用 JSON 转换器,并且配置它生成所有四个基本的变更事件部分,那么变更事件将具有如下结构:

{
 "schema": { (1)
   ...
  },
 "payload": { (2)
   ...
 },
 "schema": { (3)
   ...
 },
 "payload": { (4)
   ...
 },
}

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

序号字段名说明
1schema第一个 schema 字段是事件键的一部分。它指定一个 Kafka Connect schema,用于描述事件键的 payload 部分包含的内容。换句话说,第一个 schema 字段描述的是发生变更的表的主键结构;如果该表没有主键,则描述其第一个单列唯一键的结构。不支持多列唯一键。
可以通过设置 message.key.columns 连接器配置属性 来覆盖表的主键。在这种情况下,第一个 schema 字段描述的是由该属性标识的键的结构。
2payload第一个 payload 字段是事件键的一部分。它的结构由前面的 schema 字段描述,其中包含发生变更的那一行的键。
3schema第二个 schema 字段是事件值的一部分。它指定一个 Kafka Connect schema,用于描述事件值的 payload 部分包含的内容。换句话说,第二个 schema 描述的是发生变更的那一行的结构。通常,该 schema 中会包含嵌套的 schema。
4payload第二个 payload 字段是事件值的一部分。它的结构由前面的 schema 字段描述,其中包含发生变更的那一行的实际数据。
默认情况下,连接器会将变更事件记录流式传输到名称与事件来源表同名的主题。
从 Kafka 0.10 开始,Kafka 可以选择性地将事件键和事件值与消息创建(由生产者记录)或由 Kafka 写入日志时的时间戳 一并记录下来。

Vitess 连接器确保所有 Kafka Connect schema 名称都符合 Avro schema 名称格式。这意味着逻辑服务器名称必须以拉丁字母或下划线开头,即 a-z、A-Z 或 _。逻辑服务器名称中的其余每个字符,以及 schema 和表名中的每个字符,都必须是拉丁字母、数字或下划线,即 a-z、A-Z、0-9 或 _。如果存在无效字符,则会将其替换为下划线字符。

如果逻辑服务器名称、schema 名称或表名中包含无效字符,而区分这些名称的唯一字符恰好都是无效字符并因此被替换为下划线,就可能导致意外的冲突。

目前该连接器不允许以 @ 前缀命名列。例如,age 是合法的列名,而 @age 不是。这是因为 Vitess vstreamer 存在一个 bug,会发送列名被匿名化的事件(例如,列名 age 被匿名化为 @1)。目前没有简单的方法来区分合法的带 @ 前缀的列名和 Vitess 的这个 bug。更多讨论请参见此处。

更改事件键

对于给定的表,更改事件的键具有一种结构,该结构包含事件创建时表中主键的每个列所对应的一个字段。

请看在 commerce 键空间中定义的 customers 表,以及该表的更改事件键示例。

示例表

CREATE TABLE customers (
  id INT NOT NULL,
  first_name VARCHAR(255) NOT NULL,
  last_name VARCHAR(255) NOT NULL,
  email VARCHAR(255) NOT NULL,
  PRIMARY KEY(id)
);

示例变更事件键

如果连接器配置属性 topic.prefix 的值为 Vitess_server,那么在 customers 表具有当前定义期间,该表的每个变更事件都具有相同的键结构,其 JSON 形式如下所示:

{
  "schema": { (1)
    "type": "struct",
    "name": "Vitess_server.commerce.customers.Key", (2)
    "optional": false, (3)
    "fields": [ (4)
          {
              "name": "id",
              "index": "0",
              "schema": {
                  "type": "INT32",
                  "optional": "false"
              }
          }
      ]
  },
  "payload": { (5)
      "id": "1"
  },
}

表 2. 变更事件键的说明

序号字段名称说明
1schema键的 schema 部分指定了一个 Kafka Connect schema,用来描述键的 payload 部分的内容。
2Vitess_server.commerce.customers.Key定义键的 payload 结构的 schema 名称。该 schema 描述了被更改表的主键结构。键 schema 的名称格式为 connector-name.keyspace-name.table-name.Key。在本示例中:

- Vitess_server 是生成此事件的连接器名称。
- commerce 是包含被更改表的 keyspace。
- customers 是被更新的表。
3optional表示事件键是否必须在其 payload 字段中包含值。在本示例中,键的 payload 中必须包含值。当表没有主键时,键的 payload 字段中的值是可选的。
4fields指定 payload 中预期的每个字段,包括每个字段的名称、索引和 schema。
5payload包含生成此变更事件的行的键。在本示例中,该键包含一个值为 1 的 id 字段。
尽管 column.exclude.list 和 column.include.list 连接器配置属性允许你仅捕获表列的子集,但主键或唯一键中的所有列始终会包含在事件的键中。
如果表没有主键,则变更事件的键为 null。没有主键约束的表中的行无法被唯一标识。

变更事件值

变更事件中的值比键稍微复杂一些。与键一样,值也包含 schema 部分和 payload 部分。schema 部分包含描述 payload 部分 Envelope 结构的 schema,包括其嵌套字段。创建、更新或删除数据的操作所产生的变更事件,其值 payload 都具有信封(envelope)结构。

以用于展示变更事件键示例的同一个示例表为例:

CREATE TABLE customers (
  id INT NOT NULL,
  first_name VARCHAR(255) NOT NULL,
  last_name VARCHAR(255) NOT NULL,
  email VARCHAR(255) NOT NULL,
  PRIMARY KEY(id)
);

UPDATE 和 DELETE 操作所产生的事件中包含表中所有列的前值。

create 事件

以下示例展示了连接器为在 customers 表中创建数据的操作所生成的变更事件的值部分:

{
    "schema": { (1)
        "type": "struct",
        "fields": [
            {
                "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": "Vitess_server.commerce.customers.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": "Vitess_server.commerce.customers.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": "string",
                        "optional": true,
                        "name": "io.debezium.data.Enum",
                        "version": 1,
                        "parameters": {
                            "allowed": "true,last,false,incremental"
                        },
                        "default": "false",
                        "field": "snapshot"
                    },
                    {
                        "type": "string",
                        "optional": false,
                        "field": "db"
                    },
                    {
                        "type": "string",
                        "optional": true,
                        "field": "sequence"
                    },
                    {
                        "type": "string",
                        "optional": false,
                        "field": "keyspace"
                    },
                    {
                        "type": "string",
                        "optional": false,
                        "field": "table"
                    },
                    {
                        "type": "string",
                        "optional": false,
                        "field": "shard"
                    },
                    {
                        "type": "int64",
                        "optional": true,
                        "field": "vgtid"
                    }
                ],
                "optional": false,
                "name": "io.debezium.connector.vitess.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": "Vitess_server.commerce.customers.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": "vitess",
            "name": "my_sharded_connector",
            "ts_ms": 1559033904000,
            "ts_us": 1559033904000000,
            "ts_ns": 1559033904000000000,
            "snapshot": "false",
            "db": "",
            "sequence": null,
            "keyspace": "commerce",
            "table": "customers",
            "shard": "-80",
            "vgtid": "[{\"keyspace\":\"commerce\",\"shard\":\"80-\",\"gtid\":\"MariaDB/0-54610504-47\"},{\"keyspace\":\"commerce\",\"shard\":\"-80\",\"gtid\":\"MariaDB/0-1592148-45\"}]"
        },
        "op": "c", (9)
        "ts_ms": 1559033904863, (10)
        "ts_us": 1559033904863497, (10)
        "ts_ns": 1559033904863497147 (10)
    }
}

表 3. create 事件值字段说明

序号 字段名 描述
1 schema 值的 schema,用于描述值的有效负载(payload)的结构。对于特定的表,连接器生成的每个变更事件中,其值的 schema 都是相同的。
2 name 在 schema 部分中,每个 name 字段都指定了值的有效负载中某个字段所对应的 schema。Vitess_server.commerce.customers.Value 是有效负载中 before 和 after 字段所对应的 schema。该 schema 是 customers 表特有的。before 和 after 字段的 schema 名称采用 logicalName.keyspaceName.tableName.Value 的形式,从而确保该名称在数据库中是唯一的。这意味着,当使用 Avro 转换器时,每个逻辑源中每张表生成的 Avro schema 都有各自独立的演进历史。
3 name io.debezium.connector.vitess.Source 是有效负载中 source 字段所对应的 schema。该 schema 是 Vitess 连接器特有的,连接器为其生成的所有事件都使用它。
4 name Vitess_server.commerce.customers.Envelope 是有效负载整体结构所对应的 schema,其中 Vitess_server 是连接器名称,commerce 是 key space,customers 是表名。
5 payload 值的实际数据,即变更事件所提供 information。粗看起来,事件的 JSON 表示比它所描述的行要大得多。这是因为 JSON 表示中必须包含消息的 schema 和有效负载部分。不过,通过使用 Avro 转换器,可以显著减小连接器传输到 Kafka 主题的消息体积。
6 before 一个可选字段,用于指定事件发生前行的状态。当 op 字段为表示创建的 c 时(如本例所示),由于该变更事件针对的是新增内容,因此 before 字段为 null。
7 after 一个可选字段,用于指定事件发生后行的状态。在本例中,after 字段包含新行的 id、first_name、last_name 和 email 列的值。
8 source 描述事件源元数据的必填字段。该字段包含的信息可用于将此事件与其他事件进行比较,包括事件的来源、事件发生的先后顺序,以及这些事件是否属于同一事务。源元数据包括:- Debezium 版本
- 连接器类型和名称
- 包含新行的数据库(也称为 key space)、表和分片(shard)
- 该事件是否属于快照(始终为 false)
- 操作的 VGTID
- 变更在数据库中发生的时间戳
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 之间的时间延迟。注意:ts_ms 的单位是毫秒,但由于 MySQL 的限制(其 binlog 时间戳只能精确到秒),它实际上只具有秒级精度。

update 事件

示例 customers 表中更新操作的变更事件值,其 schema 与该表的 create 事件相同。同样,事件值的 payload 也具有相同的结构。不过,update 事件的 payload 中包含的值有所不同。下面是连接器针对 customers 表中的更新操作所生成事件的变更事件值示例:

{
    "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": "vitess",
            "name": "my_sharded_connector",
            "ts_ms": 1559033904000,
            "ts_us": 1559033904000000,
            "ts_ns": 1559033904000000000,
            "snapshot": "false",
            "db": "",
            "sequence": null,
            "keyspace": "commerce",
            "table": "customers",
            "shard": "-80",
            "vgtid": "[{\"keyspace\":\"commerce\",\"shard\":\"80-\",\"gtid\":\"MariaDB/0-54610504-47\"},{\"keyspace\":\"commerce\",\"shard\":\"-80\",\"gtid\":\"MariaDB/0-1592148-46\"}]"
        },
        "op": "u", (4)
        "ts_ms": 1465584025523,  (5)
        "ts_us": 1465584025523763,  (5)
        "ts_ns": 1465584025523763547  (5)
    }
}

表 4. update 事件值字段说明

编号 字段名称 说明
1 before 可选字段,包含数据库提交之前该行所有列的值。
2 after 可选字段,指定事件发生后该行的状态。在此示例中,first_name 的值现在是 Anne Marie。
3 source 必填字段,描述该事件的源元数据。source 字段的结构与 create 事件中的相同,但部分值有所不同。源元数据包括:

- Debezium 版本
- 连接器类型和名称
- 包含新行的数据库(也称为键空间)、表和分片
- 该事件是否属于快照的一部分(始终为 false)
- 操作的 VGTID
- 数据库中发生变更的时间戳
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 之间的延迟。注意:ts_ms 的单位是毫秒,但由于 MySQL 的限制,其精度仅到秒级,因为 MySQL 只能提供秒级的 binlog 时间戳。
更新某行主键的列会改变该行键的值。当键发生变化时,Debezium 会输出三个事件:一个带有该行旧键的 DELETE 事件和一个墓碑事件,随后是一个带有该行新键的事件。详情见下一节。

主键更新

更改行主键字段的 UPDATE 操作称为主键变更。对于主键变更,连接器不会发出 UPDATE 事件记录,而是针对旧键发出 DELETE 事件记录,针对新(已更新的)键发出 CREATE 事件记录。这些事件具有通常的结构和内容,并且每个事件还包含一个与主键变更相关的消息头:

  • DELETE 事件记录带有 __debezium.newkey 消息头。该头的值是被更新行的新主键。
  • CREATE 事件记录带有 __debezium.oldkey 消息头。该头的值是被更新行之前(旧)的主键。

delete 事件

delete 变更事件中的值与同一表的 create 和 update 事件具有相同的 schema 部分。针对示例 customers 表的 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": "vitess",
            "name": "my_sharded_connector",
            "ts_ms": 1559033904000,
            "ts_us": 1559033904000000,
            "ts_ns": 1559033904000000000,
            "snapshot": "false",
            "db": "",
            "sequence": null,
            "keyspace": "commerce",
            "table": "customers",
            "shard": "-80",
            "vgtid": "[{\"keyspace\":\"commerce\",\"shard\":\"80-\",\"gtid\":\"MariaDB/0-54610504-47\"},{\"keyspace\":\"commerce\",\"shard\":\"-80\",\"gtid\":\"MariaDB/0-1592148-47\"}]"
        },
        "op": "d", (4)
        "ts_ms": 1465581902461, (5)
        "ts_us": 1465581902461324, (5)
        "ts_ns": 1465581902461324871 (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 版本
  • 连接器类型和名称
  • 包含新行的数据库(也称为 keyspace)、表和分片
  • 该事件是否属于快照的一部分(始终为 false)
  • 操作的 VGTID
  • 在数据库中做出更改的时间戳

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 之间的延迟。注意:ts_ms 的单位是毫秒,但由于 MySQL 的限制(其只能提供秒级精度的 binlog 时间戳),该值仅具有秒级精度。

delete 变更事件记录为使用者提供了处理该行删除所需的信息。

Vitess 连接器的事件设计为可与 Kafka 日志压缩配合使用。日志压缩允许删除一些较旧的消息,前提是为每个键保留至少最新的一条消息。这样既能回收 Kafka 的存储空间,又能确保主题包含完整的数据集,可用于重新加载基于键的状态。

墓碑事件

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

数据类型映射

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

如果默认的数据类型转换不能满足你的需求,你可以为连接器创建自定义转换器。

基本类型

下表描述了连接器如何将基本的 Vitess 数据类型映射为事件字段中的字面类型(literal type)和语义类型(semantic type)。

  • 字面类型描述该值如何使用 Kafka Connect 模式类型进行字面表示:INT8、INT16、INT32、INT64、FLOAT32、FLOAT64、BOOLEAN、STRING、BYTES、ARRAY、MAP 和 STRUCT。
  • 语义类型描述 Kafka Connect 模式如何通过字段的 Kafka Connect 模式名称来体现该字段的含义。

表 6. Vitess 基本数据类型的映射

Vitess 数据类型 字面类型(模式类型) 语义类型(模式名称)及说明

BOOLEAN, BOOL

INT16

n/a

BIT(1)

暂不支持

n/a

BIT(>1)

暂不支持

n/a

TINYINT

INT16

n/a

SMALLINT[(M)]

INT16

n/a

MEDIUMINT[(M)]

INT32

n/a

INT, INTEGER[(M)]

INT32

n/a

BIGINT[(M)]

INT64

n/a

REAL[(M,D)]

FLOAT64

n/a

FLOAT[(M,D)]

FLOAT64

n/a

DOUBLE[(M,D)]

FLOAT64

n/a

CHAR[(M)]

STRING

n/a

VARCHAR[(M)]

STRING

n/a

BINARY[(M)]

BYTES

n/a

VARBINARY[(M)]

BYTES

n/a

TINYBLOB

BYTES

n/a

TINYTEXT

STRING

n/a

BLOB

BYTES

n/a

TEXT

STRING

n/a

MEDIUMBLOB

BYTES

n/a

MEDIUMTEXT

STRING

n/a

LONGBLOB

BYTES

n/a

LONGTEXT

STRING

n/a

JSON

STRING

io.debezium.data.Json
包含 JSON 文档、数组或标量的字符串表示形式。

ENUM

STRING

io.debezium.data.Enum
allowed 模式参数包含以逗号分隔的允许值列表。

SET

STRING

io.debezium.data.EnumSet
allowed 模式参数包含以逗号分隔的允许值列表。

YEAR[(2|4)]

INT32

io.debezium.time.Year

TIMESTAMP[(M)]

STRING

n/a
基于 UTC,采用 yyyy-MM-dd HH:mm:ss.SSS 格式,精度为微秒。MySQL 允许 M 的取值范围为 0-6。

NUMERIC[(M[,D])]

STRING

n/a

DECIMAL[(M[,D])]

STRING

n/a

GEOMETRY, LINESTRING, POLYGON, MULTIPOINT, MULTILINESTRING, MULTIPOLYGON, GEOMETRYCOLLECTION

暂不支持

n/a

时间类型

Vitess 时间类型取决于 time.precision.mode 连接器配置属性的值。TIMESTAMP 数据类型是个例外:除非将 time.precision.mode 设置为 connect,否则它始终表示为 io.debezium.time.ZonedTimestamp。

不带时区的时间值

DATETIME 类型表示本地日期和时间,例如 "2018-01-13 09:48:27"。如上例所示,该类型不包含时区信息。此类型的列会根据列的精度,使用 UTC 转换为纪元毫秒或微秒。TIMESTAMP 类型表示不带时区信息的时间戳。写入数据时,MySQL 会将 TIMESTAMP 类型从服务器或会话的时区转换为 UTC 格式。读取值时,数据库会将 UTC 格式转换回服务器或会话的当前时区。例如:

  • DATETIME 类型的值 2018-06-20 06:37:03 转换为 1529476623000。
  • TIMESTAMP 类型的值 2018-06-20 06:37:03 转换为 2018-06-20T13:37:03Z。

此类列会根据服务器或会话的时区,转换为 UTC 下等效的 io.debezium.time.ZonedTimestamp。默认情况下,Debezium 会向服务器查询时区。如果查询失败,你必须通过在 JDBC 连接字符串中设置 connectionTimeZone 选项来显式指定时区。例如,如果数据库的时区(无论是全局设置,还是通过 connectionTimeZone 选项为连接器配置的时区)为 "America/Los_Angeles",那么 TIMESTAMP 值 "2018-06-20 06:37:03" 将由值为 "2018-06-20T13:37:03Z" 的 ZonedTimestamp 表示。

运行 Kafka Connect 和 Debezium 的 JVM 的时区不会影响这些转换。

有关影响时间值的属性的更多信息,请参阅 连接器配置属性。

time.precision.mode=adaptive_time_microseconds(默认)

Vitess 连接器根据列的数据类型定义来确定字面类型和语义类型,使事件精确地表示数据库中的值。所有时间字段均以微秒为单位。只有在 00:00:00.000000 到 23:59:59.999999 范围内的正 TIME 字段值才能被正确捕获。

表 7. time.precision.mode=adaptive_time_microseconds 时的映射

Vitess 类型字面类型语义类型
DATEINT32io.debezium.time.Date —— 表示自 UNIX 纪元以来经过的天数。
TIME[(M)]INT64io.debezium.time.MicroTime —— 以微秒表示的时间值,不包含时区信息。MySQL 允许 M 的取值范围为 0-6。
DATETIME, DATETIME(0), DATETIME(1), DATETIME(2), DATETIME(3)INT64io.debezium.time.Timestamp —— 表示自 UNIX 纪元以来经过的毫秒数,不包含时区信息。
DATETIME(4), DATETIME(5), DATETIME(6)INT64io.debezium.time.MicroTimestamp —— 表示自 UNIX 纪元以来经过的微秒数,不包含时区信息。

time.precision.mode=connect

Vitess 连接器使用 Kafka Connect 中定义的逻辑类型。这种方式不如默认方式精确,如果数据库列的小数秒精度值大于 3,事件的精度可能会降低。该连接器能够处理范围从 00:00:00.000 到 23:59:59.999 的值。只有在确定表中的 TIME 值从不超过支持的范围时,才应设置 time.precision.mode=connect。预计 connect 设置将在 Debezium 的未来版本中被移除。

表 8. time.precision.mode=connect 时的映射关系

Vitess 数据类型字面类型语义类型
DATEINT32org.apache.kafka.connect.data.Date
表示自 UNIX 纪元以来经过的天数。
TIME[(M)]INT64org.apache.kafka.connect.data.Time
表示自午夜起以微秒为单位的时间值,不包含时区信息。
DATETIME[(M)]INT64org.apache.kafka.connect.data.Timestamp
表示自 UNIX 纪元以来经过的毫秒数,不包含时区信息。
TIMESTAMP[(M)]INT64org.apache.kafka.connect.data.Timestamp
表示自 UNIX 纪元以来经过的毫秒数。MySQL 在写入 TIMESTAMP 值时会将其转换为 UTC。

time.precision.mode=isostring

Vitess 连接器对所有时间类型使用字符串表示形式。

Vitess 数据类型字面类型语义类型
DATESTRINGn/a
TIME[(M)]STRINGn/a
DATETIME[(M)]STRINGn/a

表 9. time.precision.mode=connect 时的映射关系

设置 Vitess

Debezium 在与 Vitess 一起使用时不需要任何特殊配置。请按照 通过 Docker 本地安装 指南或 Kubernetes 的 Vitess Operator 指南中的标准说明安装 Vitess。

检查清单

  • 确保从安装 Vitess 连接器的机器可以访问 VTGate 主机及其 gRPC 端口(默认为 15991)
  • 确保从安装 Vitess 连接器的机器可以访问 VTCtld 主机及其 gRPC 端口(默认为 15999)

gRPC 身份验证

由于 Vitess 连接器从 VTGate VStream gRPC 服务器读取变更事件,因此它不需要直接连接到 MySQL 实例。所以,不需要特殊的数据库用户和权限。目前,Vitess 连接器仅支持对 VTGate gRPC 服务器进行非认证访问。

部署

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

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

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

连接器配置示例

以下是一个 Vitess 连接器的配置示例,该连接器连接到位于 192.168.99.100 上端口为 15991 的 Vitess 服务器(VTGate 的 VStream),其逻辑名称为 fullfillment。它还会连接到位于 192.168.99.101 上端口为 15999 的 VTCtld 服务器,以获取初始 VGTID。通常,你会使用 .json 文件并利用连接器提供的配置属性来配置 Debezium Vitess 连接器。

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

{
  "name": "inventory-connector",  (1)
  "config": {
    "connector.class": "io.debezium.connector.vitess.VitessConnector", (2)
    "database.hostname": "192.168.99.100", (3)
    "database.port": "15991", (4)
    "database.user": "vitess", (5)
    "database.password": "vitess_password", (6)
    "vitess.keyspace": "commerce", (7)
    "vitess.tablet.type": "MASTER", (8)
    "vitess.vtctld.host": "192.168.99.101", (9)
    "vitess.vtctld.port": "15999", (10)
    "vitess.vtctld.user": "vitess", (11)
    "vitess.vtctld.password": "vitess_password", (12)
    "topic.prefix": "fullfillment", (13)
    "tasks.max": 1 (14)
  }
}
1连接器在注册到 Kafka Connect 服务时的名称。
2此 Vitess 连接器类的名称。
3Vitess(VTGate 的 VStream)服务器的地址。
4Vitess(VTGate 的 VStream)服务器的端口号。
5Vitess 数据库服务器(VTGate gRPC)的用户名。
6Vitess 数据库服务器(VTGate gRPC)的密码。
7keyspce(也称数据库)的名称。由于未指定分片,连接器会从该 keyspace 中所有分片读取变更事件。
8读取变更事件所针对的 MySQL 实例类型(MASTER 或 REPLICA)。
9VTCtld 服务器的地址。
10VTCtld 服务器的端口。
11VTCtld 服务器(VTCtld gRPC)的用户名。
12VTCtld 数据库服务器(VTCtld gRPC)的密码。
13Vitess 集群的主题前缀,它构成一个命名空间,并用于连接器写入的所有 Kafka 主题名称、Kafka Connect 模式名称,以及使用 Avro 转换器时对应 Avro 模式的命名空间。
14任意时刻只能有一个任务在运行。

有关可在这些配置中指定的完整的 Vitess 连接器属性列表。

你可以使用 POST 命令将此配置发送到正在运行的 Kafka Connect 服务。该服务会记录配置,并启动连接器任务,该任务连接到 Vitess 数据库,并将变更事件记录流式传输到 Kafka 主题。

offset-storage-per-task 模式的连接器配置示例

当你的 Vitess 安装规模较大,需要一个以上的连接器任务来处理变更日志时,可以使用 offset-storage-per-task 特性启动多个连接器任务,让每个任务处理 Vitess 分片(shard)的子集。每个任务会将自己的偏移量(它所跟踪的各分片的 vgtid)持久化到 Kafka 偏移量主题中属于自己的分区空间里。

下面是同一个 Vitess 连接器的示例,该连接器连接到 Vitess(VTGate 的 VStream)服务器,但额外增加了三个参数以启用 offset-storage-per-task 模式。

{
  "name": "inventory-connector",
  "config": {
    "connector.class": "io.debezium.connector.vitess.VitessConnector",
    "database.hostname": "192.168.99.100",
    "database.port": "15991",
    "database.user": "vitess",
    "database.password": "vitess_password",
    "topic.prefix": "fullfillment",
    "vitess.keyspace": "commerce",
    "vitess.tablet.type": "MASTER",
    "vitess.vtctld.host": "192.168.99.101",
    "vitess.vtctld.port": "15999",
    "vitess.vtctld.user": "vitess",
    "vitess.vtctld.password": "vitess_password",
    "vitess.offset.storage.per.task": true, (1)
    "vitess.offset.storage.task.key.gen": 1, (2)
    "vitess.prev.num.tasks": 1, (3)
    "tasks.max": 2 (4)
  }
}
1指定我们要启用 offset-storage-per-task 功能
2指定当前任务并行度的代号(generation)为 1
3指定上一代任务并行度中的任务数为 1
4指定为当前任务并行度启动两个任务

任务到 Vitess 分片的分配基于简单的轮询算法。在本例中启动两个连接器任务,并假设有 4 个 Vitess 分片(-40、40-80、80-c0、c0-),task0 将处理分片(-40、80-c0),task1 将处理分片(40-80、c0-)。

需要三个配置参数的原因是确保每个连接器任务保存的偏移量(offset)不会相互冲突,并自动处理上一代任务并行度的偏移量迁移。为了确保 Kafka 偏移量主题中的分区键不会冲突,我们为每个连接器任务使用如下分区名称方案:taskId_numTasks_gen。因此,对于当前启动两个任务且代号为 1 的示例,task0 将以分区键 task0_2_1 将其偏移量写入 Kafka 的偏移量主题,task1 将使用分区键 task1_2_1。gen 配置参数用于区分不同代(generation)生成的分区键(代对应每次任务并行度的变更)。

当任务并行度发生变化时(例如,你希望启动 4 个连接器任务而不是 2 个,以处理来自 Vitess 的更大数据流量),你需要指定 tasks.max=4、vitess.offset.storage.task.key.gen=2、vitess.prev.num.tasks=2,该任务并行度代的偏移量分区将为:task0_4_2、task1_4_2、task2_4_2、task3_4_2。一旦连接器重启,它会检测到当前这 4 个分区键没有保存过偏移量,于是会从上一代键中保存的偏移量(task0_2_1 和 task1_2_1)自动执行偏移量迁移。对于当前 4 个 Vitess 分片(-40、40-80、80-c0、c0-)的示例,task0 将处理分片 (-40),task1 处理 (40-80),task2 处理 (80-c0),task3 处理 (c0-)。上一代并行度(使用 2 个任务,每个任务处理 2 个分片)中这 4 个分片的偏移量将自动迁移到当前使用 4 个任务(每个任务处理一个分片)的这一代。

请注意,在启用 offset-storage-per-task 功能之前保存到 Kafka 偏移量主题中的偏移量,其任务并行度代号默认为 0,因此在偏移量迁移期间会进行一次特殊的偏移量查找。所以,如果你的 Vitess 连接器在未启用 offset-storage-per-task 功能的情况下运行了一段时间,现在想启用该功能,请指定 vitess.offset.storage.task.key.gen=1 和 vitess.prev.num.tasks=1,以协助偏移量自动迁移。

请注意,vitess.prev.num.tasks 必须与上一代任务并行度中实际启动的任务数相匹配。连接器任务的数量通常与你指定的 tasks.max 配置参数相同,但在极少数情况下,当 tasks.max 大于 Vitess 分片数量时,连接器只会启动 the_number_of_tasks = the_number_of_vitess_shards。这种罕见情况很可能本身就是配置错误。

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

你可以使用 POST 命令将此配置发送到正在运行的 Kafka Connect 服务。该服务会记录该配置,并启动连接器任务,该任务连接到 Vitess 数据库并将变更事件记录流式传输到 Kafka 主题。

添加连接器配置

要开始运行 Vitess 连接器,请创建连接器配置,并将该配置添加到你的 Kafka Connect 集群。

前提条件

  • 安装 Vitess 连接器的机器可以访问 VTGate 主机及其 gRPC 端口(默认为 15991)
  • 安装 Vitess 连接器的机器可以访问 VTCtld 主机及其 gRPC 端口(默认为 15999)
  • 已安装 Vitess 连接器。

操作步骤

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

结果

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

监控

Debezium Vitess 连接器提供了一种额外的指标类型,是对 Kafka 和 Kafka Connect 内置的 JMX 指标支持的补充。

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

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 名称

默认情况下,Vitess 连接器为流式处理指标使用以下 MBean 名称:

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

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

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

流式传输指标

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

属性类型描述
MilliSecondsSinceLastEventlong连接器读取并处理最近一个事件以来经过的毫秒数。
TotalNumberOfEventsSeenlong自上次启动或重置以来,该连接器看到的事件总数。
NumberOfEventsFilteredlong被连接器上配置的包含/排除列表过滤规则过滤掉的事件数量。
QueueTotalCapacityint用于在流处理器与主 Kafka Connect 循环之间传递事件的队列长度。
QueueRemainingCapacityint用于在流处理器与主 Kafka Connect 循环之间传递事件的队列的剩余容量。
Connectedboolean表示连接器当前是否已连接到数据库服务器的标志。
MilliSecondsBehindSourcelong最后一个变更事件的时间戳与连接器处理该事件之间的毫秒数。该值会将数据库服务器和连接器所在机器之间的时钟差异考虑在内。
NumberOfCommittedTransactionslong已处理且已提交的事务数量。
MaxQueueSizeInByteslong用于在流处理器与主 Kafka Connect 循环之间传递事件的队列的最大缓冲区字节数。
CurrentQueueSizeInByteslong用于在流处理器与主 Kafka Connect 循环之间传递事件的队列当前缓冲区字节数。

连接器配置属性

Debezium Vitess 连接器提供了许多配置属性,你可以利用它们来实现应用所需的连接器行为。许多属性都有默认值。属性相关信息组织如下:

除非有默认值可用,否则以下配置属性均为必填项。

表 10. 必需的连接器配置属性

属性 默认值 描述

name

无默认值

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

connector.class

无默认值

连接器的 Java 类名。对于 Vitess 连接器,始终使用 io.debezium.connector.vitess.VitessConnector 这个值。

tasks.max

1

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

database.hostname

无默认值

Vitess 数据库服务器(VTGate)的 IP 地址或主机名。

database.port

15991

Vitess 数据库服务器(VTGate)的整数端口号。

vitess.keyspace

要从中流式传输变更的 keyspace 名称。

vitess.shard

n/a

要从中流式传输变更的可选分片(shard)名称。如果未配置,对于未分片的 keyspace,连接器会从其唯一的分片流式传输变更;对于已分片的 keyspace,连接器会从该 keyspace 中的所有分片流式传输变更。我们建议不要配置该属性,以便从 keyspace 中的所有分片流式传输变更,因为这种方式对重分片(reshard)操作有更好的支持。如果配置了该属性,例如 -80,连接器将从 -80 分片流式传输变更。

vitess.gtid

current

要从中流式传输的可选 GTID 位置。此属性必须与 vitess.shard 一起设置。如果未配置,连接器将从给定分片的最新位置流式传输变更。

vitess.stop_on_reshard

false

控制 Vitess 标志 stop_on_reshard。

true - 重分片(reshard)操作结束后,流将被停止。

false - 重分片操作结束后,流将自动迁移到新的分片上。

如果设置为 true,建议同时在配置中设置 vitess.gtid。

vitess.stream_keyspace_heartbeats

false

控制 Vitess 标志 StreamKeyspaceHeartbeats。

+ true - 流将接收来自 _vt.heartbeat 表的事件(tablet 也必须启用该功能)。

+ false - 流将不会接收来自 _vt.heartbeat 表的事件。

如果设置为 true,你可能还需要将 <keyspace>.heartbeat 添加到表包含列表中,以便输出这些事件。

vitess.database.user

无

Vitess 数据库服务器(VTGate)的可选用户名。若未配置,则使用未经身份验证的 VTGate gRPC。

vitess.database.password

无

Vitess 数据库服务器(VTGate)的可选密码。若未配置,则使用未经身份验证的 VTGate gRPC。

vitess.tablet.type

MASTER

用于流式传输变更的 Tablet(即 MySQL)类型:

MASTER 表示从主 MySQL 实例进行流式传输

REPLICA 表示从副本从属 MySQL 实例进行流式传输

RDONLY 表示从只读从属 MySQL 实例进行流式传输。

topic.prefix

无默认值

主题前缀,为 Debezium 捕获变更所在的特定 Vitess 数据库服务器或集群提供命名空间。数据库服务器逻辑名称只能使用字母数字字符、连字符、点和下划线。该前缀在所有其他连接器中应当是唯一的,因为它是接收来自该连接器记录的所有 Kafka 主题的主题名称前缀。

+

请勿更改此属性的值。如果更改了该名称值,重启之后,连接器将不再继续向原有主题发出事件,而是将后续事件发送到基于新值命名的主题。

table.include.list

无默认值

一个可选的、以逗号分隔的正则表达式列表,用于匹配你想要捕获其变更的表的完全限定标识符。任何未包含在 table.include.list 中的表,其变更都不会被捕获。每个标识符的形式为 keyspace.tableName。默认情况下,连接器会捕获其变更正在被捕获的每个 schema 中所有非系统表的变更。请勿同时设置 table.exclude.list 属性。

table.exclude.list

无默认值

一个可选的、以逗号分隔的正则表达式列表,用于匹配你不希望捕获其变更的表的完全限定标识符。任何未包含在 table.exclude.list 中的表,其变更都会被捕获。每个标识符的形式为 keyspace.tableName。请勿同时设置 table.include.list 属性。

column.include.list

无默认值

一个可选的、以逗号分隔的正则表达式列表,用于匹配应包含在变更事件记录值中的列的完全限定名称。列的完全限定名称形式为 keyspace.tableName.columnName。请勿同时设置 column.exclude.list 属性。

column.exclude.list

无默认值

一个可选的、以逗号分隔的正则表达式列表,用于匹配应从变更事件记录值中排除的列的完全限定名称。列的完全限定名称形式为 keyspace.tableName.columnName。请勿同时设置 column.include.list 属性。

column.truncate.to.length.chars

不适用

一个可选的、以逗号分隔的正则表达式列表,用于匹配基于字符的列的完全限定名称。如果你想在一组列中的数据超过属性名称中 length 所指定的字符数时截断这些数据,可以设置此属性。将 length 设置为正整数值,例如 column.truncate.to.20.chars。

列的完全限定名称遵循以下格式:databaseName.tableName.columnName。为了匹配列的名称,Debezium 会将你指定的正则表达式作为锚定正则表达式来应用。也就是说,指定的表达式会与列的完整名称字符串进行匹配;该表达式不会匹配列名中可能出现的子字符串。

你可以在单个配置中指定多个具有不同长度的属性。

tombstones.on.delete

true

控制 delete 事件之后是否跟一个墓碑事件。

true - 删除操作由一个 delete 事件及其后的墓碑事件表示。

false - 只发出一个 delete 事件。

源记录被删除后,发出墓碑事件(默认行为)可以让 Kafka 在主题启用了日志压缩时,彻底删除与被删除行的键相关的所有事件。

vitess.offset.storage.per.task

false

指定是否开启按任务存储偏移量的模式,以便启动多个连接器任务并按任务分区持久化偏移量。

true - 开启按任务存储偏移量的模式。

false - 不使用按任务存储偏移量的模式。

如果开启了按任务存储偏移量的模式,你还需要指定 vitess.offset.storage.task.key.gen 和 vitess.prev.num.tasks 参数。

vitess.offset.storage.task.key.gen

-1

指定当 vitess.offset.storage.per.task 开启时的任务并行度代次编号。当你决定更改连接器任务的并行度时(启动更多或更少的连接器任务),应当增加该代次编号。

vitess.prev.num.tasks

-1

指定当 vitess.offset.storage.per.task 开启时,上一代任务并行度所使用的连接器任务数量。

message.key.columns

空字符串

由分号分隔的表列表,每个表都带有与表列名匹配的正则表达式。连接器会将匹配列中的值映射为其发送到 Kafka 主题的变更事件记录中的键字段。当表没有主键,或者你希望按照非主键字段在 Kafka 主题中对变更事件记录进行排序时,此配置非常有用。

各项之间用分号分隔。在完全限定的表名与其正则表达式之间插入冒号。格式如下:

keyspace-name.table-name:_regexp_;…​

例如,

keyspaceA.table_a:regex_1;keyspaceA.table_b:regex_2;keyspaceA.table_c:regex_3

如果 table_a 有一个 id 列,且 regex_1 为 ^i(匹配任何以 i 开头的列),连接器会将 table_a 的 id 列中的值映射为其发送到 Kafka 的变更事件中的键字段。

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 命名。

snapshot.mode

initial

指定连接器启动时执行快照的条件。将该属性设置为以下值之一:

initial

当连接器启动时,如果在偏移主题(offsets topic)中未检测到值,则对数据库执行快照。

snapshot.include.collection.list

none

一个可选的、以逗号分隔的正则表达式列表,用于匹配使用 VStream Copy 时需要包含在快照中的表的全限定名。

连接器使用此属性来确定在快照阶段(即从空 GTID 启动时的 VStream)应复制哪些表。表名必须以 keyspace.table 格式指定。可以使用正则表达式匹配多个表,例如:test_keyspace.numeric_.*,test_keyspace.string_.*

如果未设置此属性或其值为空,则在快照阶段复制所有表。

此属性仅在 snapshot.mode 设置为 initial,或连接器从空 GTID 恢复时适用。

当从非空 GTID(例如 current 或有效的 binlog 位置)启动时,无论此设置如何,都会跳过复制阶段。

time.precision.mode

adaptive_time_microseconds

可以设置以下选项,以确定 Debezium 如何表示时间、日期和时间戳值的精度:

adaptive_time_microseconds

(默认)根据数据库列类型,值以毫秒、微秒或纳秒的精度表示;但 TIME 类型字段始终以微秒捕获。

connect

时间和时间戳值始终使用 Kafka Connect 为 Time、Date 和 Timestamp 定义的默认格式表示,无论数据库列的精度如何,这些格式都采用毫秒精度。

isostring

时间类型以字符串形式表示。对于某些无法用数字表示的时间值(例如 0000-00-00),会发送该值的原始字符串表示形式。

bigint.unsigned.handling.mode.mode

string

指定变更事件中应如何表示 BIGINT UNSIGNED 列。
将该属性设置为以下值之一:

string::
使用 Java 的 string 表示值

long::
使用 Java 的 long 表示值,这种方式可能无法提供所需的精度,但在消费者中使用起来要方便得多。

precise::
将值表示为精确值(Java 的 BigDecimal)。这种方式是精确的,但使用起来较为困难。

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

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

属性 默认值 描述

converters

无默认值

列出连接器可以使用的自定义转换器实例的符号名称(以逗号分隔)。例如,

isbn

必须设置 converters 属性,才能让连接器使用自定义转换器。

对于为连接器配置的每个转换器,还必须添加一个 .type 属性,用于指定实现转换器接口的类的完全限定名称。.type 属性使用以下格式:

<converterSymbolicName>.type

例如,

isbn.type: io.debezium.test.IsbnConverter

若要进一步控制已配置转换器的行为,可以添加一个或多个配置参数,以便向转换器传递值。要将任意额外的配置参数与某个转换器关联起来,请在参数名称前加上该转换器的符号名作为前缀。例如:

isbn.schema.name: io.debezium.vitess.type.Isbn

event.processing.failure.handling.mode

fail

指定连接器在处理事件时遇到异常应如何应对:

fail 会抛出异常,指出问题事件的偏移量,并导致连接器停止。

warn 会记录问题事件的偏移量,跳过该事件并继续处理。

skip 会跳过问题事件并继续处理。

max.queue.size

20240

一个正整数值,用于指定阻塞队列可容纳的最大记录数。Debezium 从数据库读取流式事件时,会先将事件放入阻塞队列,然后再写入 Kafka。当连接器接收消息的速度快于写入 Kafka 的速度,或 Kafka 变得不可用时,阻塞队列可以对从数据库读取变更事件提供背压。连接器定期记录偏移量时,队列中尚未被写出的事件会被忽略。请始终将 max.queue.size 的值设置为大于 max.batch.size 的值。

max.batch.size

2048

一个正整数值,用于指定连接器处理的每批事件的最大数量。

max.queue.size.in.bytes

0

一个长整型数值,用于指定阻塞队列的最大容量(以字节为单位)。默认情况下,阻塞队列没有容量限制。若要指定队列可占用的字节数,请将此属性设置为一个正的长整型值。
如果同时设置了 max.queue.size,那么当队列的大小达到其中任一属性所设定的限制时,写入队列的操作就会被阻塞。例如,如果你设置 max.queue.size=1000 和 max.queue.size.in.bytes=5000,那么在队列包含 1000 条记录之后,或者队列中记录的总容量达到 5000 字节之后,写入队列的操作就会被阻塞。

poll.interval.ms

500

一个正整数值,用于指定连接器在开始处理一批事件之前,等待新变更事件出现的毫秒数。默认值为 500 毫秒,即 0.5 秒。

skipped.operations

t

一个以逗号分隔的操作类型列表,表示你在流式传输过程中希望连接器跳过的操作。你可以将连接器配置为跳过以下类型的操作:

  • c(插入/创建)
  • u(更新)
  • d(删除)
  • t(清空)

由于 Debezium Vitess 连接器从不向数据变更主题发送 truncate 事件,因此设置默认的 t 选项与将该属性设置为 none 效果相同。也就是说,连接器会流式传输所有 insert、update 和 delete 操作。

statistics.metrics.enabled

true

指定连接器是否为流式传输指标收集高级统计指标,例如分位数。设置为 true 时,连接器会收集以下统计数据:

  • 最小值
  • 最大值
  • 平均值
  • P50(中位数)百分位
  • P95 百分位
  • P99 百分位
目前仅 MilliSecondsBehindSource 指标支持收集分位数。

统计数据使用一种概率数据结构(DDSketch)计算得出,可提供相对精度为 1% 的近似分位数值。设置为 false 时,连接器不会收集分位数,分位数 JMX 指标将返回 null 值。连接器仍会继续收集最小值、最大值和平均值。禁用分位数收集可略微减少内存开销。有关更多信息,请参阅流式传输指标。

provide.transaction.metadata

false

决定连接器是否生成带有事务边界的事件,并使用事务元数据充实变更事件信封。如果希望连接器执行此操作,请指定 true。详情请参阅事务元数据。

transaction.metadata.factory

io.debezium.pipeline.txmetadata.DefaultTransactionMetadataFactory

指定连接器用于跟踪事务上下文以及构建表示事务的数据结构和模式的类。io.debezium.connector.vitess.pipeline.txmetadata.VitessOrderedTransactionMetadataFactory 提供额外的事务元数据,可帮助使用者确定两个事件的正确顺序,无论它们的消费顺序如何。有关更多信息,请参阅有序事务元数据。

vitess.connector.generation

0

用于事务排序语义的代次编号。当连接器启动时,它会将此值与 Kafka Connect 偏移中存储的代次进行比较。如果两者不同,连接器会在开始流式传输之前递增所有分片的 transaction_epoch。在您进行了会影响 transaction_rank 计算方式的配置更改(例如启用或禁用事务分块)之后,请递增此值。此属性仅在 transaction.metadata.factory 设置为 io.debezium.connector.vitess.pipeline.txmetadata.VitessOrderedTransactionMetadataFactory 时才会生效。有关更多信息,请参阅 连接器代次。

vitess.keepalive.interval.ms

Long.MAX_VALUE

控制 VStream 定期 gPRC keepalive ping 之间的间隔。默认值为 Long.MAX_VALUE(禁用)。

vitess.grpc.headers

无默认值

指定以逗号分隔的 gRPC 请求头列表。默认为空。格式为:

key1:value1,key2:value2,…​

例如:

x-envoy-upstream-rq-timeout-ms:0,x-envoy-max-retries:2

vitess.grpc.max_inbound_message_size

无默认值

指定通道上允许接收的最大消息大小(以字节为单位)。

默认值为 4MiB

column.propagate.source.type

无

一个可选的、以逗号分隔的正则表达式列表,用于匹配列的完全限定名称。匹配的列的原始类型和长度会作为参数添加到所生成变更事件记录中相应字段的 schema 中。这些 schema 参数:

__debezium.source.column.type

分别用于传播可变宽度类型的原始类型名称和长度。这对于在接收端数据库中正确设置相应列的大小非常有用。列的完全限定名称形式如下:

keyspaceName.tableName.columnName

datatype.propagate.source.type

无

一个可选的、以逗号分隔的正则表达式列表,用于匹配列的数据库特定数据类型名称。匹配的列的原始类型和长度会作为参数添加到所生成变更事件记录中相应字段的 schema 中。这些 schema 参数:

__debezium.source.column.type

分别用于传播可变宽度类型的原始类型名称和长度。这对于在接收端数据库中正确设置相应列的大小非常有用。列的完全限定名称形式如下:

keyspaceName.tableName.columnName

有关 Vitess 特定数据类型名称的列表,请参阅 Vitess 连接器如何映射数据类型。

topic.naming.strategy

io.debezium.schema.SchemaTopicNamingStrategy

用于确定数据变更、模式变更、事务、心跳事件等主题名称的 TopicNamingStrategy 类的名称,默认为 SchemaTopicNamingStrategy。

override.data.change.topic.prefix

无默认值

指定当 topic.naming.strategy 设置为 io.debezium.connector.vitess.TableTopicNamingStrategy 时,连接器用于创建数据变更主题名称的前缀。连接器将使用此属性的值,而不是所指定的 topic.prefix。

topic.delimiter

.

指定主题名称的分隔符,默认为 .。

topic.cache.size

10000

用于在有界并发哈希映射中保存主题名称的容量大小。该缓存有助于确定与给定数据集合对应的主题名称。

topic.transaction

transaction

控制连接器向其发送事务元数据消息的主题名称。该主题名称遵循以下模式:

topic.prefix.topic.transaction

例如,如果主题前缀为 fulfillment,则默认主题名称为 fulfillment.transaction。

custom.metric.tags

无默认值

自定义指标标签接受键值对,用于定制 MBean 对象名称,这些键值对应追加到常规名称的末尾。每个键表示 MBean 对象名称的一个标签,其对应的值即为该标签的值。例如:k1=v1,k2=v2。

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 连接器的名称。

连接器直通配置属性

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

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

出现问题时的行为

Debezium 是一个分布式系统,用于捕获多个上游数据库中的所有变更;它绝不会遗漏或丢失事件。当系统正常运行或得到妥善管理时,Debezium 会对每条变更事件记录提供精确一次的交付保证。

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

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

配置与启动错误

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

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

在这些情况下,错误消息中包含问题的详细信息,以及可能的建议解决方法。修正配置或解决 Vitess 的问题后,重新启动连接器。

Vitess 变得不可用

当连接器正在运行时,它所连接的 Vitess 服务器(VTGate)可能由于各种原因变得不可用。如果发生这种情况,连接器会因错误而失败并停止。当服务器再次可用时,重新启动连接器。

Vitess 连接器以外部形式保存最近处理的偏移量,其形式为 Vitess VGTID。连接器重新启动并连接到服务器实例后,会与服务器通信,以便从该特定偏移量继续进行流式传输。

无效列名错误

此错误极少发生。如果你收到的消息为 Illegal prefix '@' for column: x, from schema: y, table: z,而你的表中并没有该列,那么这是由列重命名或列类型更改引起的 Vitess vstream bug。这是一个瞬时错误。你可以在短暂退避后重新启动连接器,问题应会自动解决。

Kafka Connect 进程正常停止

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

Kafka Connect 进程崩溃

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

由于在故障恢复过程中某些事件可能会被重复,消费者应始终预期可能出现重复事件。Debezium 的更改是幂等的,因此事件序列总是产生相同的状态。

在每条更改事件记录中,Debezium 连接器会插入与事件来源相关的源信息,包括 Vitess 服务器上的事件时间,以及事务更改写入到 binlog 中的位置。消费者可以跟踪这些信息,尤其是 VGTID,以判断某个事件是否为重复事件。

Kafka 变得不可用

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

连接器被停止一段时间

如果连接器被正常停止,数据库仍可继续使用。任何更改都会被记录到 Vitess 的 binlog 中。当连接器重新启动时,它会从中断的地方继续流式传输更改。也就是说,它会为连接器停止期间发生的所有数据库更改生成更改事件记录。

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

连接器在完成快照之前失败

如果快照未能完成,连接器不会自动重新尝试执行快照。如果使用之前的偏移量重启连接器,连接器将跳过快照过程并立即开始流式传输更改事件。

因此,要从该故障中恢复,请手动移除连接器偏移量,然后启动连接器。

早期 Vitess 版本的限制

Vitess 8.0.0

  • 由于 Vitess 存在一个轻微的填充问题(已在 Vitess 9.0.0 中修复),精度大于或等于 13 的小数值会在数字前产生多余的空白字符。例如,如果表定义中的列类型为 decimal(13,4),值 -1.2300 会变成 "- 1.2300",而值 1.2300 会变成 " 1.2300"。
  • 不支持 JSON 列类型。
  • VStream 8.0.0 不提供 ENUM 列允许值的附加元数据。因此,连接器不支持 ENUM 列类型。系统会输出索引编号(从 1 开始)而不是枚举值。例如,如果 ENUM 的定义为 enum('S','M','L'),则输出的值将是 "3" 而不是 "L"。

评论

登录后参与评论

正在加载评论…