源连接器

MariaDB

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

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

Debezium 连接器(MariaDB)

MariaDB 拥有二进制日志(binlog),它按照操作提交到数据库的先后顺序记录所有操作。这既包括表结构的变更,也包括表中数据的变更。MariaDB 使用 binlog 进行复制和恢复。

Debezium MariaDB 连接器读取 binlog,为行级的 INSERT、UPDATE 和 DELETE 操作生成变更事件,并将这些变更事件发布到 Kafka 主题中。客户端应用程序读取这些 Kafka 主题。

由于 MariaDB 通常会被配置为在指定时间后清除 binlog,因此 MariaDB 连接器会先对你的每个数据库执行一次初始的一致性快照。MariaDB 连接器从执行快照的位置开始读取 binlog。

有关与此连接器兼容的 MariaDB 数据库版本信息,请参阅 Debezium 发布概览。

连接器的工作原理

了解连接器所支持的 MariaDB 拓扑结构,有助于规划你的应用程序。要优化配置并运行 Debezium MariaDB 连接器,理解连接器如何跟踪表结构、公开结构变更、执行快照以及确定 Kafka 主题名称是非常有帮助的。

支持的 MariaDB 拓扑结构

Debezium MariaDB 连接器支持以下 MariaDB 拓扑结构:

独立服务器

当使用单个 MariaDB 服务器时,该服务器必须启用 binlog,Debezium MariaDB 连接器才能监控该服务器。
这通常是可以接受的,因为二进制日志也可用作增量备份。
在这种情况下,MariaDB 连接器始终连接并跟踪这个独立的 MariaDB 服务器实例。

主服务器与副本

Debezium MariaDB 连接器可以跟踪其中一台主服务器,或其中一台副本(前提是该副本已启用其 binlog),但连接器只能检测该服务器可见的集群中的变更。一般来说,除了多主拓扑之外,这不会造成问题。

连接器记录它在服务器 binlog 中的位置,而集群中每台服务器的 binlog 位置各不相同。因此,连接器只能跟踪一个 MariaDB 服务器实例。如果该服务器发生故障,必须先重启或恢复该服务器,连接器才能继续运行。

高可用集群

MariaDB 存在多种高可用解决方案,它们能显著降低应对问题和故障的难度,并几乎可以立即从中恢复。由于高可用 MariaDB 集群使用 GTID,副本能够跟踪任何主服务器上发生的所有变更。

多主

Galera 集群复制 使用一个或多个 MariaDB 副本节点,每个节点都从多个主服务器进行复制。集群复制提供了一种强大的方式,可以聚合多个 MariaDB 集群的复制。

Debezium MariaDB 连接器可以将这些多主 MariaDB 副本用作数据源,并且可以故障切换到不同的多主 MariaDB 副本,只要新副本已经追赶上旧副本即可。也就是说,新副本必须包含在第一个副本上看到的所有事务。即使连接器仅使用数据库和/或表的子集,这也同样有效,因为连接器可以配置为在尝试重新连接到新的多主 MariaDB 副本并在 binlog 中找到正确位置时,包含或排除特定的 GTID 源。

托管服务

Debezium MariaDB 连接器可以使用托管数据库选项,例如 Amazon RDS 和 Amazon Aurora。

由于这些托管选项不允许使用全局读锁,因此连接器在创建一致性快照时会使用表级锁。

架构历史主题

当数据库客户端查询数据库时,客户端使用的是数据库当前的架构。然而,数据库架构随时可能发生变化,这意味着连接器必须能够识别在记录每次插入、更新或删除操作时的架构是什么。此外,连接器未必能将当前架构应用于每个事件。如果某个事件相对较早,它可能是在应用当前架构之前记录的。

为确保正确处理架构变更之后发生的事件,MariaDB 在事务日志中不仅包含影响数据的行级变更,还包含应用于数据库的 DDL 语句。当连接器在 binlog 中遇到这些 DDL 语句时,会对其进行解析,并更新内存中每个表的架构表示。连接器使用此架构表示来识别每次插入、更新或删除操作时表的结构,并生成相应的变更事件。在一个独立的数据库架构历史 Kafka 主题中,连接器会记录所有 DDL 语句以及每条 DDL 语句在 binlog 中出现的位置。

当连接器在崩溃或正常停止后重新启动时,它会从特定位置(即特定时间点)开始读取 binlog。连接器通过读取数据库架构历史 Kafka 主题,并解析到连接器开始位置的 binlog 为止的所有 DDL 语句,来重建该时间点存在的表结构。

此数据库架构历史主题仅供连接器内部使用。此外,连接器还可以选择向面向消费者应用程序的另一个主题发出架构变更事件。

当 MariaDB 连接器捕获应用于某个表的变更,而该表正被 gh-ost 或 pt-online-schema-change 等架构变更工具处理时,迁移过程中会创建一些辅助表。你必须配置连接器,使其能够捕获这些辅助表中发生的变更。如果消费者不需要连接器为辅助表生成的记录,可以配置单条消息转换(SMT),将这些记录从连接器发出的消息中移除。

其他资源

架构变更主题

你可以配置 Debezium MariaDB 连接器,使其生成描述数据库中表所应用的架构变更的架构变更事件。连接器会将架构变更事件写入名为 <topicPrefix> 的 Kafka 主题,其中 topicPrefix 是 topic.prefix 连接器配置属性中指定的命名空间。连接器发送到架构变更主题的消息包含负载(payload),并且还可选地包含变更事件消息的架构。

架构变更事件的架构包含以下元素:

name

架构变更事件消息的名称。

type

变更事件消息的类型。

version

架构的版本。该版本是一个整数,每次架构变更时都会递增。

fields

变更事件消息中包含的字段。

示例:MariaDB 连接器架构变更主题的架构

下面的示例展示了一个以 JSON 格式表示的典型架构。

{
  "schema": {
    "type": "struct",
    "fields": [
      {
        "type": "string",
        "optional": false,
        "field": "databaseName"
      }
    ],
    "optional": false,
    "name": "io.debezium.connector.mariadb.SchemaChangeKey",
    "version": 1
  },
  "payload": {
    "databaseName": "inventory"
  }
}

架构变更事件消息的负载包含以下元素:

ddl

提供导致架构变更的 SQL CREATE、ALTER 或 DROP 语句。

databaseName

DDL 语句所应用到的数据库的名称。databaseName 的值用作消息键。

pos

这些语句在 binlog 中出现的位置。

tableChanges

架构变更后整个表架构的结构化表示。tableChanges 字段包含一个数组,其中列出了表中每个列的条目。由于结构化表示以 JSON 或 Avro 格式呈现数据,消费者无需先通过 DDL 解析器处理即可轻松读取消息。

对于处于捕获模式的表,连接器不仅会将架构变更历史存储在架构变更主题中,还会存储在内部数据库架构历史主题中。内部数据库架构历史主题仅供连接器使用,不供消费应用程序直接使用。请确保需要架构变更通知的应用程序仅从架构变更主题获取该信息。

切勿对数据库架构历史主题进行分区。为了让数据库架构历史主题正常工作,它必须保持连接器向其发出的事件记录的一致全局顺序。

为确保主题不会被拆分到多个分区中,请使用以下方法之一为该主题设置分区数量:

  • 如果手动创建数据库架构历史主题,请将分区数量指定为 1。
  • 如果使用 Apache Kafka broker 自动创建数据库架构历史主题,请将 Kafka num.partitions 配置选项的值设置为 1。
连接器向其架构变更主题发出的消息格式处于孵化阶段,可能会在不另行通知的情况下发生变更。

示例:发出到 MariaDB 连接器架构变更主题的消息

以下示例展示了一个典型的 JSON 格式架构变更消息。该消息包含表架构的逻辑表示。

{
  "schema": { },
  "payload": {
      "source": {
        "version": "3.6.3.Final",
        "connector": "mariadb",
        "name": "mariadb",
        "ts_ms": 1651535750218,
        "ts_us": 1651535750218000,
        "ts_ns": 1651535750218000000,
        "snapshot": "false",
        "db": "inventory",
        "sequence": null,
        "table": "customers",
        "server_id": 223344,
        "gtid": null,
        "file": "mariadb-bin.000003",
        "pos": 570,
        "row": 0,
        "thread": null,
        "query": null
      },
      "databaseName": "inventory",
      "schemaName": null,
      "ddl": "ALTER TABLE customers ADD middle_name varchar(255) AFTER first_name",
      "tableChanges": [
        {
          "type": "ALTER",
          "id": "\"inventory\".\"customers\"",
          "table": {
            "defaultCharsetName": "utf8mb4",
            "primaryKeyColumnNames": [
              "id"
            ],
            "columns": [
              {
                "name": "id",
                "jdbcType": 4,
                "nativeType": null,
                "typeName": "INT",
                "typeExpression": "INT",
                "charsetName": null,
                "length": null,
                "scale": null,
                "position": 1,
                "optional": false,
                "autoIncremented": true,
                "generated": true
              },
              {
                "name": "first_name",
                "jdbcType": 12,
                "nativeType": null,
                "typeName": "VARCHAR",
                "typeExpression": "VARCHAR",
                "charsetName": "utf8mb4",
                "length": 255,
                "scale": null,
                "position": 2,
                "optional": false,
                "autoIncremented": false,
                "generated": false
              },
              {
                "name": "middle_name",
                "jdbcType": 12,
                "nativeType": null,
                "typeName": "VARCHAR",
                "typeExpression": "VARCHAR",
                "charsetName": "utf8mb4",
                "length": 255,
                "scale": null,
                "position": 3,
                "optional": true,
                "autoIncremented": false,
                "generated": false
              },
              {
                "name": "last_name",
                "jdbcType": 12,
                "nativeType": null,
                "typeName": "VARCHAR",
                "typeExpression": "VARCHAR",
                "charsetName": "utf8mb4",
                "length": 255,
                "scale": null,
                "position": 4,
                "optional": false,
                "autoIncremented": false,
                "generated": false
              },
              {
                "name": "email",
                "jdbcType": 12,
                "nativeType": null,
                "typeName": "VARCHAR",
                "typeExpression": "VARCHAR",
                "charsetName": "utf8mb4",
                "length": 255,
                "scale": null,
                "position": 5,
                "optional": false,
                "autoIncremented": false,
                "generated": false
            }
          ],
          "attributes": [
            {
              "customAttribute": "attributeValue"
            }
          ]
        }
      }
    ]
  }
}

以下列表说明了上例中的部分字段:

source

payload.source 字段的结构反映了连接器写入到各表专属主题中的标准数据变更事件的结构。该字段可用于关联不同主题上的事件。

ts_ms、ts_us、ts_ns

分别以毫秒、微秒和纳秒为单位显示时间戳,表示连接器处理该事件的时间。该时间基于运行 Kafka Connect 任务的 JVM 的系统时钟。

databaseName、schemaName

这两个字段提供了包含该变更的数据库和模式的名称。databaseName 字段的值用作事件记录的消息键。

ddl

提供产生该模式变更的 DDL。ddl 字段可以包含多条 DDL 语句。每条语句都适用于 databaseName 字段中指定的数据库。多条 DDL 语句按照其应用到数据库的顺序出现。

客户端可以提交适用于多个数据库的多条 DDL 语句。如果 MariaDB 以原子方式应用这些语句,连接器会按顺序处理这些 DDL 语句,按数据库分组,并为每个分组创建一个模式变更事件。如果 MariaDB 逐条应用这些语句,连接器则为每条语句创建一个单独的模式变更事件。

tableChanges

一个数组,包含一项或多项由 DDL 命令生成的模式变更。

type

描述变更的种类。该字段可包含以下值之一:

CREATE

表已创建。

ALTER

表已被修改。

DROP

表已被删除。

id

指定被创建、修改或删除的表的完整标识符。在重命名表的情况下,该标识符由 <旧表名>,<新表名> 拼接而成。

table

表示变更应用后的表元数据。

primaryKeyColumnNames

构成该表主键的列的列表。

columns

已变更表中每一列的元数据。

attributes

每个表变更的自定义属性元数据。

有关模式变更事件的更多信息,请参阅模式历史主题。

快照

Debezium MariaDB 连接器首次启动时,会对您的数据库执行初始一致性快照。该快照使连接器能够为数据库的当前状态建立基准。

Debezium 在执行快照时可以使用不同的模式。快照模式由 snapshot.mode 配置属性决定。该属性的默认值为 initial。您可以通过修改 snapshot.mode 属性的值来定制连接器创建快照的方式。

连接器在执行快照时会完成一系列任务。具体步骤因快照模式以及数据库当前生效的表锁定策略而异。Debezium MariaDB 连接器在执行使用全局读锁或表级锁的初始快照时,会完成不同的步骤。

使用全局读锁的初始快照

你可以通过修改 snapshot.mode 属性的值来自定义连接器创建快照的方式。如果你配置了不同的快照模式,连接器将按照该工作流的修改版本来完成快照。有关在不允许使用全局读锁的环境中执行快照过程的信息,请参阅表级锁的快照工作流。

Debezium MariaDB 连接器使用全局读锁执行初始快照的默认工作流

下表列出了 Debezium 使用全局读锁创建快照时所遵循的工作流步骤。

步骤 操作
1 建立与数据库的连接。
2 确定要捕获的表。默认情况下,连接器会捕获所有非系统表的数据。快照完成后,连接器会继续为指定的表流式传输数据。如果你希望连接器仅捕获特定表的数据,可以通过设置 table.include.list 或 table.exclude.list 等属性,指示连接器仅捕获部分表或表元素的数据。
3 对要捕获的表获取全局读锁,以阻止其他数据库客户端的写入操作。 快照本身并不能阻止其他客户端执行可能干扰连接器读取 binlog 位置和表结构的 DDL 语句。连接器在读取 binlog 位置期间会保持全局读锁,并在后续步骤中释放该锁。
4 使用可重复读隔离级别启动事务,以确保事务内的所有后续读取都基于一致性快照进行。
使用这些隔离语义可能会减慢快照的进度。如果快照耗时过长,请考虑改用其他隔离配置,或者跳过初始快照,改为运行增量快照。

5

读取当前的 binlog 位置。

6

获取数据库中所有表的结构,或获取所有指定要捕获的表的结构。连接器会将结构信息持久化到其内部数据库模式历史主题中,包括所有必需的 DROP…​ 和 CREATE…​ DDL 语句。
模式历史提供了变更事件发生时生效的结构信息。

默认情况下,连接器会获取数据库中每张表的结构,包括未配置为捕获的表。如果表未配置为捕获,初始快照只会获取它们的结构,而不会捕获任何表数据。

有关为什么快照会为未包含在初始快照中的表持久化结构信息的详细信息,请参阅了解初始快照为何会捕获所有表的结构。

7

释放第 3 步中获取的全局读锁。其他数据库客户端现在可以向数据库写入数据。

8

在第 5 步中连接器读取到的 binlog 位置处,连接器开始扫描指定要捕获的表。在扫描过程中,连接器完成以下任务:

  1. 确认该表是在快照开始之前创建的。如果该表是在快照开始之后创建的,则连接器会跳过该表。快照完成并且连接器转入流式传输后,它会为在快照开始之后创建的任何表发出变更事件。
  2. 为从表中捕获的每一行生成一个 read 事件。所有 read 事件都包含相同的 binlog 位置,即第 5 步中获取的位置。
  3. 将每个 read 事件发送到源表对应的 Kafka 主题。
  4. 如果适用,释放数据表锁。

9

提交事务。

10

在连接器偏移量中记录快照已成功完成。

由此生成的初始快照捕获了被捕获表中每一行的当前状态。以此基准状态为基础,连接器会在后续变更发生时将其捕获下来。

快照进程开始后,如果该进程因连接器故障、重新平衡或其他原因而中断,则会在连接器重启后重新启动该进程。

连接器完成初始快照后,会从第 5 步中读取到的位置继续流式传输,从而不会遗漏任何更新。

如果连接器因任何原因再次停止,重启后它会从上次中断的位置继续流式传输更改。

连接器重启后,如果日志已被清理,连接器在日志中的位置可能已不可用。此时连接器会失败,并返回一个指示需要执行新快照的错误。要配置连接器在此情况下自动发起快照,请将 snapshot.mode 属性的值设置为 when_needed。有关 Debezium MariaDB 连接器故障排查的更多提示,请参阅出现问题时的行为。

使用表级锁的初始快照

在某些数据库环境中,管理员不允许使用全局读锁。如果 Debezium MariaDB 连接器检测到不允许使用全局读锁,那么它在执行快照时会使用表级锁。要使连接器执行使用表级锁的快照,Debezium 连接器用于连接 MariaDB 的数据库帐户必须拥有 LOCK TABLES 权限。

Debezium MariaDB 连接器使用表级锁执行初始快照的默认工作流程

下表展示了 Debezium 在使用表级读锁创建快照时所遵循的工作流程步骤。有关在不允许使用全局读锁的环境中执行快照流程的信息,请参阅使用全局读锁的快照工作流程。

步骤 操作
1 建立与数据库的连接。
2 确定要捕获的表。默认情况下,连接器捕获所有非系统表。要让连接器捕获表或表元素的子集,你可以设置多个 include 和 exclude 属性来过滤数据,例如 table.include.list 或 table.exclude.list。
3 获取表级锁。
4 使用可重复读语义启动事务,以确保事务内的所有后续读取都是针对一致性快照进行的。
5 读取当前的 binlog 位置。
6 读取连接器配置为捕获更改的数据库和表的结构。连接器会将结构信息持久化到其内部数据库结构历史主题中,包括所有必需的 DROP…​ 和 CREATE…​ DDL 语句。结构历史提供了更改事件发生时生效的结构信息。

默认情况下,连接器会捕获数据库中每张表的架构,包括未配置为捕获的表。如果表未配置为捕获,初始快照只会捕获其结构,而不会捕获任何表数据。

有关为什么快照会保留未包含在初始快照中的表的架构信息,请参阅了解为什么初始快照会捕获所有表的架构。

7

在第 5 步中连接器读取到的 binlog 位置上,连接器开始扫描指定为捕获的表。在扫描过程中,连接器完成以下任务:

  1. 确认该表是在快照开始之前创建的。如果该表是在快照开始之后创建的,连接器会跳过该表。快照完成后,连接器转入流式传输阶段,会为快照开始之后创建的任何表发出变更事件。
  2. 为从表中捕获的每一行数据生成一个 read 事件。所有 read 事件都包含相同的 binlog 位置,即在第 5 步中获取的位置。
  3. 将每个 read 事件发送到源表对应的 Kafka 主题中。
  4. 释放数据表锁(如果适用)。

8

提交事务。

9

释放表级锁。其他数据库客户端现在可以写入之前被锁定的任何表。

10

在连接器偏移量中记录快照已成功完成。

表 1. snapshot.mode 连接器配置属性的设置

设置 说明

always

连接器每次启动时都会执行快照。快照包含被捕获表的结构和数据。指定此值后,每次连接器启动时都会用被捕获表数据的完整表示来填充主题。快照完成后,连接器开始流式传输后续数据库变更的事件记录。

initial

连接器按照创建初始快照的默认工作流中所述执行数据库快照。快照完成后,连接器开始流式传输后续数据库变更的事件记录。

initial_only

连接器执行数据库快照。快照完成后,连接器停止运行,不会流式传输后续数据库变更的事件记录。

schema_only

已弃用,请参阅 no_data。

no_data

连接器捕获所有相关表的结构,执行创建初始快照的默认工作流中描述的所有步骤,但不会创建 READ 事件来表示连接器启动时的数据集(第 7.2 步)。

schema_only_recovery

已弃用,请参阅 recovery。

recovery

将此选项设置为恢复丢失或损坏的数据库架构历史主题。重启后,连接器会执行一次快照,根据源表重建该主题。你也可以设置该属性,定期清理意外增长的数据库架构历史主题。

警告:如果在连接器上次关闭后数据库中已提交了架构变更,请勿使用此模式执行快照。

when_needed

连接器启动后,仅在检测到以下情况之一时才执行快照:

  • 无法检测到任何主题偏移量。
  • 之前记录的偏移量所指定的日志位置在服务器上不可用。

configuration_based

将快照模式设置为 configuration_based,通过带有 snapshot.mode.configuration.based 前缀的一组连接器属性来控制快照行为。

custom

custom 快照模式允许你注入自己实现的 io.debezium.spi.snapshot.Snapshotter 接口。将 snapshot.mode.custom.name 配置属性设置为你的实现中 name() 方法提供的名称。该名称位于 Kafka Connect 集群的类路径上。如果使用 DebeziumEngine,则该名称包含在连接器 JAR 文件中。有关更多信息,请参阅自定义快照器 SPI。

有关更多信息,请参阅连接器配置属性表中的 snapshot.mode。

理解初始快照为何会捕获所有表的架构历史

连接器执行的初始快照会捕获两类信息:

表数据

关于连接器 table.include.list 属性中指定的表的 INSERT、UPDATE 和 DELETE 操作的信息。

架构数据

描述应用于表的结构变更的 DDL 语句。架构数据会持久化到内部架构历史主题中,如果配置了连接器的架构变更主题,也会持久化到该主题中。

运行初始快照后,你可能会注意到快照会捕获未指定为捕获对象的表的架构信息。默认情况下,初始快照旨在捕获数据库中每张表的架构信息,而不仅仅是被指定为捕获对象的表。连接器要求表的架构存在于架构历史主题中,才能捕获该表。通过让初始快照捕获不属于原始捕获集合的表的架构数据,Debezium 使连接器能够在日后需要时随时从这些表中捕获事件数据。如果初始快照未捕获某张表的架构,你必须先将该架构添加到历史主题中,连接器才能从该表中捕获数据。

在某些情况下,你可能希望限制初始快照中的架构捕获。当你希望缩短完成快照所需的时间时,这会很有用。或者,当 Debezium 通过一个可访问多个逻辑数据库的用户账户连接到数据库实例,但你希望连接器仅捕获特定逻辑数据库中表的变更时,也会用到该功能。

其他信息

从初始快照未捕获的表中捕获数据(无架构变更)

在某些情况下,你可能希望连接器从初始快照未捕获架构的表中捕获数据。根据连接器的配置,初始快照可能只为数据库中的特定表捕获表架构。如果历史主题中不存在该表架构,连接器将无法捕获该表,并会报告架构缺失错误。

你仍然可以从该表捕获数据,但必须执行额外的步骤来添加表架构。

前提条件

操作步骤

  1. 停止连接器。

  2. 删除由 schema.history.internal.kafka.topic property 属性指定的内部数据库架构历史主题。

  3. 对连接器配置应用以下更改:

  4. 将 snapshot.mode 设置为 recovery。

  5. 将 schema.history.internal.store.only.captured.tables.ddl 的值设置为 false。

  6. 将你希望连接器捕获的表添加到 table.include.list 中。这可以保证连接器在将来能够重建所有表的架构历史。

  7. 重启连接器。快照恢复过程会根据表的当前结构重新构建架构历史。

  8. (可选)快照完成后,发起一次增量快照,以捕获新添加表中的现有数据,以及连接器离线期间其他表发生的变更。

  9. (可选)将 snapshot.mode 重置回 no_data,以防止连接器在未来重启后再次发起恢复。

捕获初始快照未捕获的表中的数据(架构变更)

如果对表应用了架构变更,则在架构变更之前提交的记录与在变更之后提交的记录具有不同的结构。当 Debezium 从表中捕获数据时,它会读取架构历史,以确保为每个事件应用正确的架构。如果架构历史主题中不存在该架构,连接器将无法捕获该表,并会产生错误。

如果你要从初始快照未捕获的表中捕获数据,并且该表的架构已被修改,则必须将该架构添加到历史主题中(如果尚不存在)。你可以通过运行新的架构快照,或对该表运行初始快照来添加该架构。

前提条件

  • 你要从一个架构未被连接器在初始快照期间捕获的表中捕获数据。
  • 该表发生了架构变更,导致待捕获的记录不具备统一的结构。

过程

初始快照已捕获所有表的架构(store.only.captured.tables.ddl 设置为 false)

  1. 编辑 table.include.list 属性,指定你要捕获的表。
  2. 重启连接器。
  3. 如果你想从新添加的表中捕获现有数据,请发起一次增量快照。

初始快照未捕获所有表的架构(store.only.captured.tables.ddl 设置为 true)

如果初始快照未保存你要捕获的表的架构,请执行以下过程之一:

过程 1:架构快照,然后进行增量快照

在此过程中,连接器首先执行 schema 快照。然后你可以启动增量快照,使连接器能够同步数据。

  1. 停止连接器。

  2. 删除由 schema.history.internal.kafka.topic property 属性指定的内部数据库 schema 历史主题。

  3. 清除已配置的 Kafka Connect offset.storage.topic 中的偏移量。有关如何删除偏移量的更多信息,请参阅 Debezium 社区常见问题。

    删除偏移量的操作只能由具备操纵 Kafka Connect 内部数据经验的高级用户执行。该操作具有潜在的破坏性,应仅在别无选择时才作为最后手段使用。
  4. 按以下步骤为连接器配置中的属性设置值:

    1. 将 snapshot.mode 属性的值设置为 no_data。
    2. 编辑 table.include.list,添加你想要捕获的表。
  5. 重启连接器。

  6. 等待 Debezium 捕获新增和已有表的 schema。连接器停止后各表中发生的数据变更不会被捕获。

  7. 为确保不丢失数据,请启动增量快照。

过程 2:初始快照,随后可选的增量快照

在此过程中,连接器对数据库执行完整的初始快照。与任何初始快照一样,如果数据库中包含大量大表,运行初始快照可能是一项耗时的操作。快照完成后,你可以选择触发增量快照,以捕获连接器离线期间发生的任何变更。

  1. 停止连接器。
  2. 删除由 schema.history.internal.kafka.topic property 属性指定的内部数据库 schema 历史主题。
  3. 清除已配置的 Kafka Connect offset.storage.topic 中的偏移量。有关如何删除偏移量的更多信息,请参阅 Debezium 社区常见问题。
删除偏移量仅应由具备操作 Kafka Connect 内部数据经验的高级用户执行。此操作具有潜在破坏性,只能作为最后手段。
  1. 编辑 table.include.list,添加你要捕获的表。

  2. 按照以下步骤为连接器配置中的属性设置值:

    1. 将 snapshot.mode 属性的值设置为 initial。
    2. (可选)将 schema.history.internal.store.only.captured.tables.ddl 设置为 false。
  3. 重启连接器。连接器会执行完整数据库快照。快照完成后,连接器将转入流式传输阶段。

  4. (可选)若要捕获连接器停机期间发生的所有数据变更,请启动增量快照。

基于分块的并行快照

基于分块的并行快照通过将较小的工作单元分配到多个线程,加快初始快照的速度,从而改善工作负载均衡并提升快照的容错能力。

基于分块的并行快照是一项孵化中的功能。

当 snapshot.max.threads 设置为大于 1 的值时,连接器会根据主键范围将每张表划分为分块,并将这些分块分配到可用线程中。每个线程并发地对其分配的分块执行快照,连接器在遍历分块的过程中为捕获到的每一行数据发出 READ 事件。

snapshot.max.threads.multiplier 属性控制连接器相对于线程数量为每张表创建的分块数量。默认情况下,连接器为每个线程创建一个分块。设置更高的倍数会创建更多、更小的分块,从而使线程的负载更加均衡。例如,在 4 个线程且倍数为 2 的情况下,连接器会创建 8 个分块,而不是 4 个。

没有主键的表以及使用快照查询覆盖的表,会被作为一个分块处理,并回退到单线程快照方式。

下表总结了两种并行快照方式之间的差异。

特性 分块式(默认) 传统方式(每线程一表)

每个线程的工作单元

每个线程处理表中的一批行。分块由主键值的范围定义。

每个线程处理整张表。

负载分配

线程在完成当前工作后认领可用的分块,从而帮助在各线程之间均衡工作负载。

在快照期间,每个线程始终被分配给同一张表。处理完较小表的线程会保持空闲,而其他线程继续处理较大的表。

空闲连接风险

低。

线程在快照完成之前始终保持活动状态。

高。

处理完一张表的线程在等待其他线程完成期间可能保持空闲。在强制执行连接超时的环境中,空闲连接可能导致连接器在快照完成后无法正常关闭连接。关闭连接失败可能会引发异常,即使快照已成功捕获所有数据。

分块式快照是默认行为。要使用分块式快照,请将 legacy.snapshot.max.threads 设置为 false。

要恢复到传统的每线程一表行为,请将 legacy.snapshot.max.threads 设置为 true。

如果你使用传统行为并在快照完成后遇到空闲连接故障,可将 snapshot.max.threads 设置为 1 作为临时解决方案,然后重试快照。

传统的每线程一表行为已被弃用,并将在未来版本中移除。

internal.legacy.snapshot.max.threads 是 legacy.snapshot.max.threads 的已弃用别名,不应在新配置中使用。

默认情况下,连接器仅在首次启动后运行一次初始快照操作。在这次初始快照之后,正常情况下连接器不会重复执行快照过程。连接器后续捕获的任何变更事件数据都仅通过流式处理过程传入。

不过,在某些情况下,连接器在初始快照期间获取的数据可能变得过期、丢失或不完整。为了提供重新捕获表数据的机制,Debezium 提供了执行临时快照的选项。在 Debezium 环境中发生以下任何变化后,你可能需要执行临时快照:

  • 连接器配置被修改,以捕获不同的表集合。
  • Kafka 主题被删除且必须重新构建。
  • 由于配置错误或其他问题导致数据损坏。

要为之前已捕获过快照的表重新运行快照,可以通过发起所谓的临时快照(ad-hoc snapshot)来实现。临时快照需要使用信号表。你通过向 Debezium 信号表发送信号请求来发起临时快照。

当你对已有表发起临时快照时,连接器会将内容追加到该表已存在的主题中。如果之前存在的主题已被删除,只要启用了自动创建主题,Debezium 就能自动重新创建该主题。

临时快照信号会指定要包含在快照中的表。快照可以捕获数据库的全部内容,也可以只捕获数据库中表的子集。此外,快照还可以只捕获数据库中表内容的子集。

通过向信号表发送 execute-snapshot 消息来指定要捕获的表。将 execute-snapshot 信号的类型设置为 incremental 或 blocking,并按以下表格所示提供要包含在快照中的表名:

表 2. 临时 execute-snapshot 信号记录示例

字段 默认值
type incremental
指定要运行的快照类型。目前可以请求 incremental(增量)或 blocking(阻塞)快照。
data-collections 无
一个数组,包含与要纳入快照的表的完全限定名相匹配的正则表达式。对于 MariaDB 连接器,使用以下格式指定表的完全限定名:database.table。
additional-conditions 无
一个可选数组,用于指定连接器评估的一组附加条件,以确定要包含在快照中的记录子集。每个附加条件都是一个对象,用于指定过滤临时快照所捕获数据的条件。每个附加条件可以设置以下参数:

data-collection
过滤器所作用的表的完全限定名。你可以对每张表应用不同的过滤器。

filter
指定数据库记录中必须存在的列值,快照才会包含该记录,例如 "color='blue'"。

你为 filter 参数赋的值,与为阻塞快照设置 snapshot.select.statement.overrides 属性时在 SELECT 语句的 WHERE 子句中指定的值类型相同。
surrogate-key 无
一个可选字符串,用于指定连接器在快照过程中用作表主键的列名。

触发临时增量快照

你可以通过在信号表中添加 execute-snapshot 信号类型的条目,或向 Kafka 信号主题发送信号消息,来发起一次临时增量快照。连接器处理完该消息后,便会开始快照操作。快照进程会读取每个表的第一个和最后一个主键值,并将这些值作为每张表的起始点和结束点。根据表中的条目数量和所配置的块大小,Debezium 会将表划分为若干块,并依次逐块执行快照。

更多信息请参阅增量快照。

发起临时阻塞快照

你可以通过在信号表或信号主题中添加 execute-snapshot 信号类型的条目,来发起一次临时阻塞快照。连接器处理完该消息后,便会开始快照操作。连接器会暂时停止流式传输,然后按照与初始快照相同的过程,对指定的表发起快照。快照完成后,连接器会恢复流式传输。

更多信息请参阅阻塞快照。

增量快照

为了在快照管理方面提供更大的灵活性,Debezium 引入了一种补充性的快照机制,称为增量快照。增量快照依赖于 Debezium 的向 Debezium 连接器发送信号机制。增量快照基于 DDD-3 设计文档。

在增量快照中,Debezium 不像初始快照那样一次性捕获数据库的完整状态,而是分阶段、以一系列可配置的块来捕获每张表。你可以指定希望快照捕获的表,以及每个块的大小。块大小决定了快照在每次对数据库执行获取操作时收集的行数。增量快照的默认块大小为 1024 行。

随着增量快照的进行,Debezium 使用水位标记来跟踪其进度,并记录已捕获的每一行表数据。与标准的初始快照过程相比,这种分阶段捕获数据的方式具有以下优势:

  • 你可以将增量快照与流式数据捕获并行运行,而不必等快照完成后再进行流式捕获。在整个快照过程中,连接器会持续从变更日志中捕获近实时事件,两者互不阻塞。
  • 如果增量快照的进度被中断,你可以在不丢失任何数据的情况下恢复它。恢复后,快照将从中断的位置继续,而不是从头重新捕获该表。
  • 你可以随时按需运行增量快照,并根据需要重复该过程,以适应数据库的更新。例如,在修改连接器配置、向 table.include.list 属性中添加表之后,你可能会重新运行一次快照。

增量快照过程

运行增量快照时,Debezium 会按主键对每张表进行排序,然后根据配置的分块大小将表拆分为若干块。随后它逐块处理,捕获该块中的每一行数据。对于捕获的每一行,快照都会发出一个 READ 事件。该事件表示该块的快照开始时该行的值。

随着快照的进行,其他进程很可能仍在访问数据库,并可能修改表记录。为了反映此类更改,INSERT、UPDATE 或 DELETE 操作会照常被提交到事务日志中。同样,正在进行的 Debezium 流式处理过程会继续检测这些更改事件,并向 Kafka 发出相应的更改事件记录。

Debezium 如何解决具有相同主键的记录之间的冲突

在某些情况下,流式处理过程发出的 UPDATE 或 DELETE 事件可能会乱序到达。也就是说,流式处理过程可能在快照捕获包含该行 READ 事件的数据块之前,就发出了修改该表行的事件。当快照最终发出该行对应的 READ 事件时,其值已被取代。为确保乱序到达的增量快照事件能按正确的逻辑顺序被处理,Debezium 采用了一种用于解决冲突的缓冲机制。只有在快照事件与流式事件之间的冲突得到解决之后,Debezium 才会向 Kafka 发出事件记录。

快照窗口

为帮助解决迟到的 READ 事件与修改同一表行的流式事件之间的冲突,Debezium 采用了所谓的快照窗口。快照窗口标定了增量快照捕获指定表块数据的时间区间。在某个数据块的快照窗口打开之前,Debezium 保持常规行为,直接将事务日志中的事件向下游发送到目标 Kafka 主题。但从某个数据块的快照打开之时起,直到其关闭,Debezium 会执行去重步骤,以解决具有相同主键的事件之间的冲突。

对于每个数据收集,Debezium 会发出两种类型的事件,并将二者的记录存储在同一个目标 Kafka 主题中。它直接从表中捕获的快照记录以 READ 操作的形式发出。与此同时,随着用户继续更新数据收集中的记录,事务日志也会更新以反映每次提交,Debezium 会针对每项变更发出 UPDATE 或 DELETE 操作。

随着快照窗口打开并开始处理一个快照数据块,Debezium 会将快照记录送入内存缓冲区。在快照窗口期间,缓冲区中 READ 事件的主键会与传入流式事件的主键进行比较。如果没有匹配,流式事件记录将直接发送到 Kafka。如果 Debezium 检测到匹配,则会丢弃缓冲的 READ 事件,并将流式记录写入目标主题,因为流式事件在逻辑上取代了静态的快照事件。在该数据块的快照窗口关闭后,缓冲区中只剩下不存在相关事务日志事件的 READ 事件。Debezium 会将这些剩余的 READ 事件发送到该表的 Kafka 主题。

连接器会对每个快照数据块重复此过程。

要使 Debezium 能够执行增量快照,必须授予连接器写入信号表的权限。

只有能够配置为执行只读增量快照的连接器才无需写入权限(MariaDB、MySQL 或 PostgreSQL)。

目前,你可以使用以下任一方法来启动增量快照:

触发增量快照

要启动增量快照,你可以向源数据库上的信号表发送一个临时快照信号。快照信号以 SQL INSERT 查询的形式提交。

Debezium 检测到信号表中的变更后,会读取该信号,并执行所请求的快照操作。

你提交的查询会指定要包含在快照中的表,并可选择性地指定快照操作的类型。Debezium 目前支持 incremental(增量)和 blocking(阻塞)两种快照类型。

要指定包含在快照中的表,请提供一个 data-collections 数组,其中列出这些表,或者列出用于匹配表的正则表达式数组,例如:

{"data-collections": ["public.MyFirstTable", "public.MySecondTable"]}

数据集合名称区分大小写。增量快照信号的 data-collections 数组没有默认值。如果 data-collections 数组为空,Debezium 会将空数组解释为无需执行任何操作,因此不会执行快照。

如果要包含在快照中的表名含有点号(.)、空格或其他非字母数字字符,则必须使用双引号将表名括起来。例如,要包含 db1 数据库中名为 My.Table 的表,请使用以下格式:"db1.\"My.Table\""。

前提条件

使用源信号通道触发增量快照

  1. 发送一个 SQL 查询,将临时增量快照请求添加到信号表中:

    INSERT INTO <signalTable> (id, type, data) VALUES ('<id>', '<snapshotType>', '{"data-collections": ["<fullyQualfiedTableName>","<fullyQualfiedTableName>"],"type":"<snapshotType>","additional-conditions":[{"data-collection": "<fullyQualfiedTableName>", "filter": "<additional-condition>"}]}');

例如,

INSERT INTO db1.debezium_signal (id, type, data)
values ('ad-hoc-1',
    'execute-snapshot',
    '{"data-collections": ["db1.table1", "db1.table2"],
    "type":"incremental",
    "additional-conditions":[{"data-collection": "db1.table1" ,"filter":"color=\'blue\'"}]}');

命令中 id、type 和 data 参数的值对应信号表的字段。以下列表说明了上例中的参数:

INSERT INTO database.debezium_signal

指定源数据库上信号表的完全限定名称。

values

id

id 参数包含值 ad-hoc,这是一个任意字符串,用作信号请求的 id 标识符。

type

type 参数指定要执行的操作类型,本例中为 execute-snapshot。

data

信号的 data 字段包含以下字段:

data-collections

一个数组,包含要纳入快照的表名或用于匹配表名的正则表达式。

type

信号 data 字段中可选的 type 组成部分,用于指定要运行的快照操作类型。

有效值为 incremental 和 blocking。

如果未指定该值,连接器默认执行增量快照。

additional-conditions

一个可选数组,用于指定连接器据以评估并确定纳入快照的记录子集的一组附加条件。数组中的每个附加条件都是一个包含 data-collection 和 filter 属性的对象。

additional-conditions 数组中的每个 data-collection 均以其完全限定名称指定。可以指定一个或多个数据集合。每个集合都可以与一个可选的 filter 参数配对。可以为每个数据集合配对一个唯一的过滤器。filter 属性的值由一个列标签和一个值组成。运行快照时,对于数组中列出的每个集合,只会捕获包含指定过滤条件的行。

在该示例中,对于 table1 数据集合,该信号所触发的快照只会捕获字段 color 的值为 'blue' 的行。

有关 additional-conditions 参数的更多信息,请参阅使用 additional-conditions 运行临时增量快照。

使用物理行标识符作为代理键

某些数据库提供物理行标识符,这是一种表示行在磁盘上物理位置的伪列。这些标识符可以显著提升增量快照分块操作的性能。

物理行标识符在以下场景中尤为有用:

具有复合主键的表

当表没有单列代理键,而是使用多个列作为主键时,分块查询会变得复杂且低效。

索引利用效率低下

数据库查询优化器往往无法有效利用增量快照分块查询所产生的析取条件来使用索引,从而导致全表扫描。

通过使用物理行标识符作为替代键,Debezium 能够生成更简单的基于范围的查询,充分利用行标识符固有的有序性,从而显著提升性能。

Debezium 支持以下物理行标识符作为替代键:

数据库标识符说明
OracleROWIDOracle 表中某一行的物理地址。使用 ROWID 可以显著提升增量快照的性能,尤其适用于具有复合主键或索引效率较低的表。

下面的示例展示了如何使用 Oracle 的 ROWID 作为替代键来触发增量快照的 SQL 查询:

INSERT INTO db1.myschema.debezium_signal (id, type, data)
VALUES ('ad-hoc-1',
    'execute-snapshot',
    '{"data-collections": ["db1.myschema.mytable"],
      "type": "incremental",
      "surrogate-key": "ROWID"}');

物理行标识符在某些情况下可能会发生变化,从而影响快照的一致性。Oracle 的 ROWID 可能在表重组操作期间发生变化,例如 ALTER TABLE MOVE、分区维护或 SHRINK SPACE 操作。

为确保数据一致性,在使用物理行标识符的增量快照执行期间,请不要执行表维护操作,例如表移动、分区管理或收缩操作。

使用 additional-conditions 运行临时增量快照

如果希望快照仅包含表中内容的某个子集,可以通过向快照信号追加 additional-conditions 参数来修改信号请求。

典型快照的 SQL 查询形式如下:

SELECT * FROM <tableName> ....

通过添加 additional-conditions 参数,您可以向 SQL 查询中追加一个 WHERE 条件,如下例所示:

SELECT * FROM <data-collection> WHERE <filter> ....

以下示例展示了一条 SQL 查询,用于向信号表发送一个附带额外条件的即席增量快照请求:

INSERT INTO <signalTable> (id, type, data) VALUES ('<id>', '<snapshotType>', '{"data-collections": ["<fullyQualfiedTableName>","<fullyQualfiedTableName>"],"type":"<snapshotType>","additional-conditions":[{"data-collection": "<fullyQualfiedTableName>", "filter": "<additional-condition>"}]}');

例如,假设你有一个 products 表,其中包含以下列:

  • id(主键)
  • color
  • quantity

如果你想让 products 表的增量快照仅包含 color=blue 的数据项,可以使用以下 SQL 语句来触发快照:

INSERT INTO db1.debezium_signal (id, type, data) VALUES('ad-hoc-1', 'execute-snapshot', '{"data-collections": ["db1.products"],"type":"incremental", "additional-conditions":[{"data-collection": "db1.products", "filter": "color=blue"}]}');

additional-conditions 参数还支持传递基于多个列的条件。例如,使用前面示例中的 products 表,你可以提交一个查询,触发增量快照,且仅包含 color=blue 且 quantity>10 的那些商品的数据:

INSERT INTO db1.debezium_signal (id, type, data) VALUES('ad-hoc-1', 'execute-snapshot', '{"data-collections": ["db1.products"],"type":"incremental", "additional-conditions":[{"data-collection": "db1.products", "filter": "color=blue AND quantity>10"}]}');

以下示例展示了由连接器捕获的增量快照事件的 JSON 内容。

示例 1. 增量快照事件消息

{
    "before":null,
    "after": {
        "pk":"1",
        "value":"New data"
    },
    "source": {
        ...
        "snapshot":"incremental"
    },
    "op":"r",
    "ts_ms":"1620393591654",
    "ts_us":"1620393591654547",
    "ts_ns":"1620393591654547920",
    "transaction":null
}

以下列表描述了上一个增量快照事件消息示例中的部分字段:

snapshot

指定快照的类型。

op

指定操作的类型。对于快照事件,op 字段的值为 r,因为快照是一种 READ 操作。

使用 Kafka 信号通道触发增量快照

要使用 Kafka 信号通道触发一次临时增量快照,请向已配置的 Kafka 信号主题发送一条 execute-snapshot 消息。

Kafka 消息的键必须与连接器配置选项 topic.prefix 的值相匹配。

消息的值是一个包含 type 和 data 字段的 JSON 对象。

信号类型为 execute-snapshot,data 字段必须包含以下字段:

表 3. 执行快照的数据字段

字段 默认值
type incremental,要执行的快照类型。目前 Debezium 支持 incremental 和 blocking 两种类型。详见下一节。
data-collections 不适用,一个由逗号分隔的正则表达式组成的数组,用于匹配要包含在快照中的表的完全限定名称。
请使用与 signal.data.collection 配置选项相同的格式指定名称。数据集合名称区分大小写。
additional-conditions 不适用,一个可选的附加条件数组,用于指定连接器评估的条件,从而指定要包含在快照中的记录子集。
每个附加条件都是一个对象,用于指定临时快照所捕获数据的过滤条件。每个附加条件可以设置以下参数:
data-collection:过滤器所应用的表的完全限定名称。可以为每个表应用不同的过滤器。
filter:指定数据库记录中必须存在的列值,快照才会包含该记录,例如 "color='blue'"。
为 filter 参数赋的值类型与设置阻塞快照的 snapshot.select.statement.overrides 属性时在 SELECT 语句的 WHERE 子句中指定的值类型相同。

示例 2. 一条 execute-snapshot Kafka 消息

Key = `test_connector`

Value = `{"type":"execute-snapshot","data": {"data-collections": ["{collection-container}.table1", "{collection-container}.table2"], "type": "INCREMENTAL"}}`

带有 additional-conditions 的临时增量快照

Debezium 使用 additional-conditions 字段来选取表中内容的子集。

通常,Debezium 执行快照时会运行如下 SQL 查询:

SELECT * FROM <tableName> …​.

当快照请求中包含 additional-conditions 属性时,该属性的 data-collection 和 filter 参数会被追加到 SQL 查询中,例如:

SELECT * FROM <data-collection> WHERE <filter> …​.

例如,假设有一个 products 表,包含 id(主键)、color 和 brand 列,如果你希望快照只包含 color='blue' 的内容,那么在请求快照时,可以添加 additional-conditions 属性来过滤内容:

Key = `test_connector`

Value = `{"type":"execute-snapshot","data": {"data-collections": ["db1.products"], "type": "INCREMENTAL", "additional-conditions": [{"data-collection": "db1.products" ,"filter":"color='blue'"}]}}`

你还可以使用 additional-conditions 属性传入基于多个列的条件。例如,沿用前面示例中的 products 表,如果你希望快照仅包含 products 表中 color='blue' 且 brand='MyBrand' 的内容,可以发送如下请求:

Key = `test_connector`

Value = `{"type":"execute-snapshot","data": {"data-collections": ["db1.products"], "type": "INCREMENTAL", "additional-conditions": [{"data-collection": "db1.products" ,"filter":"color='blue' AND brand='MyBrand'"}]}}`

停止增量快照

在某些情况下,可能需要停止增量快照。例如,你可能发现快照的配置不正确,或者希望确保数据库的其他操作有足够的可用资源。你可以通过向源数据库上的信号表发送信号,来停止已经在运行的快照。

通过发送 SQL INSERT 查询,向信号表提交停止快照信号。停止快照信号将快照操作的 type 指定为 incremental,并且可以选择性地指定希望从当前运行的快照中排除的表。Debezium 检测到信号表中的变化后,会读取该信号,并在增量快照操作正在进行时将其停止。

其他资源

前提条件

使用源信号通道停止增量快照

  1. 向信号表发送用于停止临时增量快照的 SQL 查询:

    INSERT INTO <signalTable> (id, type, data) values ('<id>', 'stop-snapshot', '{"data-collections": ["<fullyQualfiedTableName>","<fullyQualfiedTableName>"],"type":"incremental"}');

例如,

INSERT INTO db1.debezium_signal (id, type, data)
values ('ad-hoc-1',
    'stop-snapshot',
    '{"data-collections": ["db1.table1", "db1.table2"],
    "type":"incremental"}');

信号命令中 id、type 和 data 参数的值对应于信号表的字段。

以下列表描述了前面信号示例中的各个字段:

database.debezium_signal

指定源数据库上信号表的完全限定名称。

ad-hoc-1

信号 id 参数的值。这个任意字符串提供了一个标签,有助于区分各个信号请求,并将日志消息与信号表中的条目关联起来。Debezium 本身不使用该字符串。

stop-snapshot

type 参数,用于标识该信号要触发的操作。

data-collections

信号 data 字段的一个可选组成部分,用于指定要从快照中移除的表名数组或与表名匹配的正则表达式数组。该数组列出的正则表达式按照 database.table 格式匹配表的完全限定名称。

如果在 data 字段中省略该组件,信号将停止当前正在进行的整个增量快照。

incremental

信号 data 字段的一个必填组成部分,用于指定要停止的快照操作类型。目前唯一有效的选项是 incremental。如果未指定 type 值,信号将无法停止增量快照。

使用 Kafka 信号通道停止增量快照

要使用 Kafka 信号通道停止正在进行的增量快照,请向已配置的 Kafka 信号主题发送一条 stop-snapshot 消息。

Kafka 消息的键必须与连接器配置选项 topic.prefix 的值相匹配。

消息的值是一个包含 type 和 data 字段的 JSON 对象。

信号类型为 stop-snapshot,data 字段必须包含以下字段:

字段默认值值说明
typeincremental要执行的快照类型。目前 Debezium 仅支持 incremental 类型。详情请参见下一节。
data-collections不适用一个可选的数组,其中包含以逗号分隔的正则表达式,用于匹配要从快照中移除的表的完全限定名;也可以是一个表名数组或用于匹配表名的正则表达式。请使用 database.table 格式指定表名。

表 4. 执行快照数据字段

以下示例展示了一条典型的 stop-snapshot Kafka 消息:

Key = `test_connector`

Value = `{"type":"stop-snapshot","data": {"data-collections": ["db1.table1", "db1.table2"], "type": "INCREMENTAL"}}`

只读增量快照

Debezium MariaDB 连接器支持在与数据库建立只读连接的情况下运行增量快照。为了以只读方式运行增量快照,连接器会使用已执行的全局事务 ID(GTID)集合作为高水位和低水位标记。通过将二进制日志(binlog)事件或服务器心跳的 GTID 与低水位和高水位标记进行比较,来更新数据块窗口的状态。

要切换到只读实现方式,请将 read.only 属性的值设置为 true。

前提条件

  • 启用 MariaDB GTID。

  • 如果连接器从多线程副本(即 replica_parallel_workers 的值大于 0 的副本)读取数据,你必须设置以下选项之一:

    • replica_preserve_commit_order=ON
    • slave_preserve_commit_order=ON

当 MariaDB 连接为只读时,你可以使用任何可用的信号通道,而不必使用 source 通道。

自定义快照器 SPI

若要超出标准快照模式所提供的功能来自定义快照行为,你可以实现一个或多个 Debezium 快照器 SPI 接口。这些接口可以控制快照是否运行、如何查询数据以及是否锁定表。

io.debezium.snapshot.spi.Snapshotter

控制连接器是否执行快照。

io.debezium.snapshot.spi.SnapshotQuery

控制快照期间如何查询数据。

io.debezium.snapshot.spi.SnapshotLock

控制连接器在执行快照时是否锁定表。

io.debezium.snapshot.spi.Snapshotter 接口。所有内置的快照模式都实现了此接口。

/**
 * {@link Snapshotter} is used to determine the following details about the snapshot process:
 * <p>
 * - Whether a snapshot occurs. <br>
 * - Whether streaming continues during the snapshot. <br>
 * - Whether the snapshot includes schema (if supported). <br>
 * - Whether to snapshot data or schema following an error.
 * <p>
 * Although Debezium provides many default snapshot modes,
 * to provide more advanced functionality, such as partial snapshots,
 * you can customize implementation of the interface.
 * For more information, see the documentation.
 *
 *
 *
 */
@Incubating
public interface Snapshotter extends Configurable {

    /**
     * @return the name of the snapshotter.
     *
     *
     */
    String name();

    /**
     * @param offsetExists is {@code true} when the connector has an offset context (i.e. restarted)
     * @param snapshotInProgress is {@code true} when the connector is started, but a snapshot is already in progress
     *
     * @return {@code true} if the snapshotter should take a data snapshot
     */
    boolean shouldSnapshotData(boolean offsetExists, boolean snapshotInProgress);

    /**
     * @param offsetExists is {@code true} when the connector has an offset context (i.e. restarted)
     * @param snapshotInProgress is {@code true} when the connector is started, but a snapshot is already in progress
     *
     * @return {@code true} if the snapshotter should take a schema snapshot
     */
    boolean shouldSnapshotSchema(boolean offsetExists, boolean snapshotInProgress);

    /**
     * @return {@code true} if the snapshotter should stream after taking a snapshot
     */
    boolean shouldStream();

    /**
     * @return {@code true} whether the schema can be recovered if database schema history is corrupted.
     */
    boolean shouldSnapshotOnSchemaError();

    /**
     * @return {@code true} whether the snapshot should be re-executed when there is a gap in data stream.
     */
    boolean shouldSnapshotOnDataError();

    /**
     *
     * @return {@code true} if streaming should resume from the start of the snapshot
     * transaction, or {@code false} for when a connector resumes and takes a snapshot,
     * streaming should resume from where streaming previously left off.
     */
    default boolean shouldStreamEventsStartingFromSnapshot() {
        return true;
    }

    /**
     * Lifecycle hook called after the snapshot phase is successful.
     */
    default void snapshotCompleted() {
        // no operation
    }

    /**
     * Lifecycle hook called after the snapshot phase is aborted.
     */
    default void snapshotAborted() {
        // no operation
    }
}

io.debezium.snapshot.spi.SnapshotQuery 接口。所有内置的快照查询模式都实现了该接口。

/**
 * {@link SnapshotQuery} is used to determine the query used during a data snapshot
 *
 *
 */
public interface SnapshotQuery extends Configurable, Service {

    /**
     * @return the name of the snapshot lock.
     *
     *
     */
    String name();

    /**
     * Generate a valid query string for the specified table, or an empty {@link Optional}
     * to skip snapshotting this table (but that table will still be streamed from)
     *
     * @param tableId the table to generate a query for
     * @param snapshotSelectColumns the columns to be used in the snapshot select based on the column
     *                              include/exclude filters
     * @return a valid query string, or none to skip snapshotting this table
     */
    Optional<String> snapshotQuery(String tableId, List<String> snapshotSelectColumns);

}

io.debezium.snapshot.spi.SnapshotLock 接口。所有内置的快照锁定模式都实现了该接口。

/**
 * {@link SnapshotLock} is used to determine the table lock mode used during schema snapshot
 *
 *
 */
public interface SnapshotLock extends Configurable, Service {

    /**
     * @return the name of the snapshot lock.
     *
     *
     */
    String name();

    /**
     * Returns a SQL statement for locking the given table during snapshotting, if required by the specific snapshotter
     * implementation.
     */
    Optional<String> tableLockingStatement(Duration lockTimeout, String tableId);

}

阻塞快照

阻塞快照可让你在连接器运行期间按需捕获表的完整、一致的快照,同时会临时暂停流式传输,直到快照完成。阻塞快照依赖 Debezium 的向 Debezium 连接器发送信号机制。

阻塞快照的行为与初始快照完全相同,只是你可以在运行时触发它。

在以下情况下,你可能希望运行阻塞快照,而不是使用标准的初始快照流程:

  • 你新增了一张表,并希望在连接器运行期间完成该表的快照。
  • 你新增了一张大表,并希望快照能在比增量快照更短的时间内完成。

阻塞快照流程

运行阻塞快照时,Debezium 会停止流式传输,然后按照初始快照期间所使用的相同流程,对指定的表发起快照。快照完成后,会恢复流式传输。

配置快照

你可以在信号的 data 组件中设置以下属性:

  • data-collections:用于指定必须进行快照的表。

  • data-collections:指定你希望快照包含的表。该属性接受一个以逗号分隔的正则表达式列表,用于匹配完全限定的表名。该属性的行为与 table.include.list 属性的行为类似,后者用于指定在阻塞快照中要捕获的表。

  • additional-conditions:你可以为不同的表指定不同的过滤条件。

    • data-collection 属性是要应用过滤条件的表的完全限定名称,根据数据库的不同,可能区分大小写,也可能不区分大小写。
    • filter 属性的取值应与 snapshot.select.statement.overrides 中使用的值相同,即应按大小写匹配的表的完全限定名称。

例如:

  {"type": "blocking", "data-collections": ["schema1.table1", "schema1.table2"], "additional-conditions": [{"data-collection": "schema1.table1", "filter": "SELECT * FROM [schema1].[table1] WHERE column1 = 0 ORDER BY column2 DESC"}, {"data-collection": "schema1.table2", "filter": "SELECT * FROM [schema1].[table2] WHERE column2 > 0"}]}

可能存在的重复

在发送触发快照的信号与流处理停止并开始快照之间,可能存在一段延迟。由于这段延迟,快照完成后,连接器可能会发出一些与快照所捕获的记录重复的事件记录。

通过信号设置 binlog 位置

在某些情况下,精确控制连接器从 binlog 中的哪个位置开始捕获数据库变更是很有用的。你可以在运行时发送一个 Debezium set-binlog-position 信号,以动态指定 MariaDB 连接器从数据库 binlog(二进制日志)的哪个位置开始读取。

根据你在 set-binlog-position 信号中指定的位置,连接器要么会排除日志中的事件,要么会发送重复的事件。

如果你指定的 binlog 位置晚于当前位置,连接器会跳过当前位置与新位置之间的历史变更事件。连接器不会捕获或流式传输被跳过范围内的变更。

另一方面,如果该信号指示连接器从早于当前位置的点恢复读取,连接器会重新播放它之前已发送过的日志条目,从而向目标主题发出重复事件。

你可能需要在以下情况中重置连接器的 binlog 位置:

灾难恢复

在系统故障后,指示连接器从日志中最后一个已知的有效位置恢复流式传输。从指定位置重新开始有助于确保数据不会丢失或重复,从而维护整个管道的数据一致性。

数据损坏

你在数据库或目标 Kafka 主题中检测到损坏或不正确的数据,希望跳过这些损坏的数据,并从已知的有效位置恢复。

数据库刷新后避免重新播放旧数据

在从其他来源(例如备份)刷新数据库之后,连接器可能会尝试从旧的 binlog 位置重新播放所有历史事件,这可能导致主题中出现重复数据,并造成大量资源消耗。通过从特定的近期位置恢复流式传输,你可以跳过不需要的历史事件,从而节省成本和时间。

部分复制

当将捕获的数据引导至新的 Kafka 主题时,你希望连接器从特定点开始处理日志,以排除完整的历史事件,只发送特定子集的数据库变更。通过这种方式,可以防止连接器处理过时的旧数据,使连接器仅捕获相关的近期变更。

测试

使用特定子集的数据测试连接器,排除快照捕获的完整历史数据。

set-binlog-position 信号的格式

Debezium 的 set-binlog-position 信号消息包含标准的 id、type 和 data 属性。消息的具体格式取决于你所使用的信号通道。有关 Debezium 信号格式的更多信息,请参见向 Debezium 连接器发送信号。

要在 set-binlog-position 信号中指定 binlog 位置,需要在信号的数据部分指定 binlog 文件和位置,或者指定 GTID 集合,具体取决于是否启用了 GTID(全局事务标识符)模式。同一条信号中不能同时使用这两种格式。

以下示例说明了在信号中指定 binlog 位置的两种格式:

binlog 文件和位置格式

当 GTID 模式被禁用时,在信号中指定 binlog 文件名和位置,如以下示例所示:

{"binlog_filename": "mysql-bin.000003", "binlog_position": 1234}

binlog_filename 必须符合 binlog 命名规则(例如 MySQL 的 mysql-bin.000001 或 MariaDB 的 mariadb-bin.000001)。

binlog_position 必须是非负整数。

GTID 集合格式

启用 GTID 模式时,需在信号中指定 GTID 集合。GTID 集合的值格式在 MySQL 和 MariaDB 之间有所不同:

  • MySQL 使用基于 UUID 的 GTID,例如:

    {"gtid_set": "3E11FA47-71CA-11E1-9E33-C80AA9429562:1-100"}
  • MariaDB 使用 domain-server-sequence 三元组,例如:

    {"gtid_set": "0-1-100"}

gtid_set 属性的值必须是所用数据库服务器的有效 GTID 集合字符串。

发送信号以重新定位流

在运行时发送 Debezium 的 set-binlog-position 信号,以指定 MariaDB 连接器的 binlog 读取位置。

前提条件

  • 连接器已配置为使用可用的 Debezium 信号通道 之一。
  • 源数据库上存在 信号数据集合。
  • 已设置连接器配置属性 heartbeat.interval.ms,以确保偏移量变更能够持久化。

操作步骤

  1. 确定数据库历史记录中的目标位置(binlog 文件与位置,或 GTID 集合)。

  2. 发送 set-binlog-position 信号,其中包含标准的 id、type 和 data 属性,并采用所选信号通道要求的格式。

    在信号的数据部分中,指定第 1 步中确定的目标位置。

    例如,对于源信号通道,根据是否启用了 GTID,将以下某个 SQL 命令添加到信号的 data 组件中:

    未启用 GTID 时使用的 SQL(使用 binlog 文件与位置)

    INSERT INTO debezium_signal (id, type, data) VALUES (
      'set-position-001',
      'set-binlog-position',
      '{"binlog_filename": "mariadb-bin.000003", "binlog_position": 1234}'
    );

启用 GTID 时使用的 SQL(指定 domain-server-sequence GTID 集合)

INSERT INTO debezium_signal (id, type, data) VALUES (
  'set-gtid-001',
  'set-binlog-position',
  '{"gtid_set": "0-1-100"}'
);
  1. 重启连接器,从新位置开始流式传输。

主题名称

默认情况下,MariaDB 连接器会将表中发生的所有 INSERT、UPDATE 和 DELETE 操作的变更事件写入该表对应的单个 Apache Kafka 主题。

连接器按以下约定为变更事件主题命名:

topicPrefix.databaseName.tableName

假设 fulfillment 是主题前缀,inventory 是数据库名称,且该数据库包含名为 orders、customers 和 products 的表。Debezium MariaDB 连接器将事件发送到三个 Kafka 主题,数据库中的每个表对应一个主题:

fulfillment.inventory.orders
fulfillment.inventory.customers
fulfillment.inventory.products

以下列表提供了默认名称各组成部分的定义:

topicPrefix

topic.prefix 连接器配置属性所指定的主题前缀。

schemaName

发生该操作的 schema 名称。

tableName

发生该操作的表名。

连接器采用类似的命名约定,为其内部数据库 schema 历史主题、schema 变更主题和事务元数据主题命名。

如果默认主题名称不能满足你的需求,可以配置自定义主题名称。要配置自定义主题名称,请在逻辑主题路由 SMT 中指定正则表达式。有关如何使用逻辑主题路由 SMT 自定义主题命名的更多信息,请参阅主题路由。

事务元数据

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

Debezium 接收事务元数据的限制

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

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

status

BEGIN 或 END。

id

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

ts_ms

事务边界事件(BEGIN 或 END 事件)在数据源处的时间。如果数据源未向 Debezium 提供事件时间,则该字段表示 Debezium 处理该事件的时间。

event_count(针对 END 事件)

该事务发出的事件总数。

data_collections(针对 END 事件)

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

示例

{
  "status": "BEGIN",
  "id": "0e4d5dcd-a33b-11ea-80f1-02010a22a99e:10",
  "ts_ms": 1486500577125,
  "event_count": null,
  "data_collections": null
}

{
  "status": "END",
  "id": "0e4d5dcd-a33b-11ea-80f1-02010a22a99e:10",
  "ts_ms": 1486500577691,
  "event_count": 2,
  "data_collections": [
    {
      "data_collection": "s1.a",
      "event_count": 1
    },
    {
      "data_collection": "s2.a",
      "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": "1580390884335",
  "ts_us": "1580390884335472",
  "ts_ns": "1580390884335472987",
  "transaction": {
    "id": "0e4d5dcd-a33b-11ea-80f1-02010a22a99e:10",
    "total_order": "1",
    "data_collection_order": "1"
  }
}

数据变更事件

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

Debezium 和 Kafka Connect 的设计围绕着持续的事件消息流。然而,这些事件的结构可能会随时间变化,这可能难以让消费者处理。为了解决这个问题,每个事件包含其内容的模式(schema),或者如果您正在使用模式注册表,则包含一个模式 ID,消费者可以使用该 ID 从注册表获取模式。这使得每个事件都是自包含的。

下面的骨架 JSON 显示了变更事件的四个基本组件。Debezium 在变更消息中表示这些组件的确切方式取决于您在应用程序中如何配置 Kafka Connect 转换器。只有当您配置转换器生成 schema 字段时,变更事件才会包含 schema 字段。同样,只有当您配置转换器生成事件键和事件负载时,变更事件才会包含事件键和事件负载。如果您使用 JSON 转换器并配置其生成所有四个基本变更事件组件,则变更事件的结构如下:

{
 "schema": {
   ...
  },
 "payload": {
   ...
 },
 "schema": {
   ...
 },
 "payload": {
   ...
 },
}

以下列表详细说明了前一个变更事件结构中所示的各个组件:

schema

事件消息中的第一个 schema 字段是事件键的一部分。事件键的 schema 是一个 Kafka Connect schema,它定义了事件键的 payload 部分中所包含的结构和数据类型。事件键结构反映了被数据库操作所变更的表的主键结构。如果源表没有主键,则事件键结构映射的是源表唯一键的结构。

可以通过设置 message.key.columns 连接器配置属性来覆盖表的主键。在这种情况下,第一个 schema 字段描述的是由该属性标识的键的结构。

payload

第一个 payload 字段是事件键的一部分。其结构与前面的事件键 schema 的结构相对应,其中包含被变更行的键。

schema

第二个 schema 字段是事件值的一部分。事件值的这个 schema 表示定义了事件值 payload 部分中所包含的结构和数据类型的 Kafka Connect schema。这个 schema 的值反映了被数据库操作所变更的行的结构。通常,事件值 schema 中包含嵌套的 schema。

payload

第二个 payload 字段是事件值的一部分。其结构与前面的事件值 schema 的结构相对应,其中包含被变更行的实际数据。

默认情况下,连接器会将变更事件记录流式传输到名称与其事件来源表相同的主题中。有关更多信息,请参阅主题名称。

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

如果逻辑服务器名称、数据库名或表名中包含无效字符,并且区分不同名称的字符恰好都是无效的、因此被替换为下划线,就可能导致意外的冲突。

变更事件键

变更事件的键包含已变更表的键的 schema 以及被变更行的实际键。schema 及其对应的 payload 都包含一个字段,对应连接器创建事件时,已变更表的 PRIMARY KEY(或唯一约束)中的每一列。

请看下面的 customers 表,其后是该表的变更事件键示例。

CREATE TABLE customers (
  id INTEGER NOT NULL AUTO_INCREMENT PRIMARY KEY,
  first_name VARCHAR(255) NOT NULL,
  last_name VARCHAR(255) NOT NULL,
  email VARCHAR(255) NOT NULL UNIQUE KEY
) AUTO_INCREMENT=1001;

每个捕获 customers 表变更的变更事件都具有相同的事件键 schema。因此,Debezium 为该表发出的每个变更事件,其事件消息都具有相同的键结构。只要 customers 表的 schema 保持不变,键结构就保持稳定。上述表定义会产生如下 JSON 示例中的键结构:

{
 "schema": {
    "type": "struct",
    "name": "mariadb-server-1.inventory.customers.Key",
    "optional": false,
    "fields": [
      {
        "field": "id",
        "type": "int32",
        "optional": false
      }
    ]
  },
 "payload": {
    "id": 1001
  }
}

以下列表描述了 Debezium 从 customers 表发出的变更事件消息中,事件键的一些元素:

schema

事件键的 schema 元素指定了 Kafka Connect 模式,该模式定义了事件键 payload 部分中存在的结构和数据类型。

name

指定定义键负载(payload)结构的模式名称。该模式描述了发生变更的表的主键结构。事件键的模式名称格式如下:connector-name.database-name.table-name.Key。

在前面的事件消息中,模式名称由以下元素组成:

mariadb-server-1

指定生成此事件的连接器名称。

inventory

指定包含发生变更的表的数据库。

customers

指定发生变更的表。

fields

指定 payload 中的必填字段。fields 元素包含以下属性:

field

字段的名称,例如 id。

type

字段的数据类型,例如 int32。

optional

一个布尔值,用于指定事件键的负载字段是否必须包含值。当表没有主键时,键的负载字段中的值是可选的。

payload

指定产生此变更事件的行的键。在前面的示例中,payload 包含一个字段 id,其值为 1001。

变更事件值

变更事件值包含操作前后完整行的状态,以及源元数据和操作类型,其结构为 Kafka Connect 的 Envelope。

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

考虑用于展示变更事件键示例的同一张示例表:

CREATE TABLE customers (
  id INTEGER NOT NULL AUTO_INCREMENT PRIMARY KEY,
  first_name VARCHAR(255) NOT NULL,
  last_name VARCHAR(255) NOT NULL,
  email VARCHAR(255) NOT NULL UNIQUE KEY
) AUTO_INCREMENT=1001;

针对此表的变更事件,其值部分的说明如下:

create 事件

create 事件包含 INSERT 操作的完整行数据,该新行的状态被记录在事件值负载的 after 字段中。

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

{
  "schema": {
    "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": "mariadb-server-1.inventory.customers.Value",
        "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": "mariadb-server-1.inventory.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": "boolean",
            "optional": true,
            "default": false,
            "field": "snapshot"
          },
          {
            "type": "string",
            "optional": false,
            "field": "db"
          },
          {
            "type": "string",
            "optional": true,
            "field": "table"
          },
          {
            "type": "int64",
            "optional": false,
            "field": "server_id"
          },
          {
            "type": "string",
            "optional": true,
            "field": "gtid"
          },
          {
            "type": "string",
            "optional": false,
            "field": "file"
          },
          {
            "type": "int64",
            "optional": false,
            "field": "pos"
          },
          {
            "type": "int32",
            "optional": false,
            "field": "row"
          },
          {
            "type": "int64",
            "optional": true,
            "field": "thread"
          },
          {
            "type": "string",
            "optional": true,
            "field": "query"
          }
        ],
        "optional": false,
        "name": "io.debezium.connector.mariadb.Source",
        "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": "mariadb-server-1.inventory.customers.Envelope"
  },
  "payload": {
    "op": "c",
    "ts_ms": 1465491411815,
    "ts_us": 1465491411815437,
    "ts_ns": 1465491411815437158,
    "before": null,
    "after": {
      "id": 1004,
      "first_name": "Anne",
      "last_name": "Kretchmar",
      "email": "[email protected]"
    },
    "source": {
      "version": "3.6.3.Final",
      "connector": "mariadb",
      "name": "mariadb-server-1",
      "ts_ms": 0,
      "ts_us": 0,
      "ts_ns": 0,
      "snapshot": false,
      "db": "inventory",
      "table": "customers",
      "server_id": 0,
      "gtid": null,
      "file": "mariadb-bin.000003",
      "pos": 154,
      "row": 0,
      "thread": 7,
      "query": "INSERT INTO customers (first_name, last_name, email) VALUES ('Anne', 'Kretchmar', '[email protected]')"
    }
  }
}

以下列表描述了前述 create 事件消息中值部分的选定字段:

schema

指定事件值的模式(schema)。该模式描述了值中负载(payload)的结构。只要表结构保持不变,Debezium 为某张表发出的每个变更事件都使用相同的值模式。

name

schema 部分可以包含多个 name 字段。每个 name 字段指定事件值负载中某个字段的模式。

mariadb-server-1.inventory.customers.Value 是负载中 before 和 after 字段的模式。该模式是 customers 表所特有的。

before 和 after 字段的模式名称采用 logicalName.tableName.Value 的形式。这种格式可确保模式名称在数据库内唯一。在使用 Avro 转换器 的环境中,模式名称唯一可确保每个逻辑源中每张表的 Avro 模式拥有各自独立的演进历史和版本历史。

"name": "io.debezium.connector.mariadb.Source"

io.debezium.connector.mariadb.Source 是负载中 source 字段的模式。该模式是 MariaDB 连接器所特有的,连接器将其用于它生成的所有事件。

"name": "mariadb-server-1.inventory.customers.Envelope"

指定负载整体结构的模式名称。该模式名称由以下几部分组成:

mariadb-server-1

指定生成此事件的连接器名称。

inventory

指定包含被修改表的数据库名称。

customers

指定被修改的表名。

payload

指定被修改行的实际数据。

由于事件的 JSON 表示形式同时包含消息的模式部分和负载部分,其体积通常比它所描述的行更大。为了减小连接器流式传输到 Kafka 主题的消息体积,你可以使用 Avro 转换器。

op

指定导致连接器生成此事件的操作类型。在此示例中,c 表示执行了 create 操作,从而产生了一条新行。该字段可以包含以下值之一:

c创建一行。
u更新一行。
d删除一行。
r读取一行(仅适用于快照)。

ts_ms、ts_us、ts_ns

以毫秒、微秒和纳秒为单位显示时间戳,指示连接器处理该事件的时间。该时间基于运行 Kafka Connect 任务的 JVM 中的系统时钟。

before

一个可选字段,用于指定事件发生前行的状态。由于前面示例中的 op 字段为 c(create,创建),该变更事件描述的是新数据,因此 before 字段的值为 null。

after

一个可选字段,用于指定事件发生后行的状态。在此示例中,after 字段包含新行的 id、first_name、last_name 和 email 列的值。

source

一个必填字段,用于描述事件的源元数据。该字段包含的信息可用于将此事件与其他事件进行比较,包括事件的来源、事件发生的顺序,以及这些事件是否属于同一个事务。源元数据提供以下信息:

version

Debezium 版本。

name

连接器名称。

connector

连接器类型。

ts_ms、ts_us、ts_ns

以毫秒、微秒和纳秒为单位显示的时间戳,表示变更在数据库中发生的时间。

snapshot

指定该事件是否由快照操作产生。

db、table

包含新行的数据库和表的名称。

server_id

MariaDB 服务器 ID(如果可用)。

file

记录该事件的二进制日志文件名。

pos

二进制日志位置。

row

事件中的行。

thread

创建该事件的 MariaDB 线程 ID(仅非快照事件)。

query

在源表上执行该操作的 SQL 命令。

如果在 MariaDB 数据库配置中启用了 binlog_annotate_row_events 选项,并且在连接器配置中启用了 include.query 属性,那么 source 字段还会提供一个 query 字段,其中包含导致该变更事件的原始 SQL 语句。

update 事件

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

{
  "schema": { ... },
  "payload": {
    "before": {
      "id": 1004,
      "first_name": "Anne",
      "last_name": "Kretchmar",
      "email": "[email protected]"
    },
    "after": {
      "id": 1004,
      "first_name": "Anne Marie",
      "last_name": "Kretchmar",
      "email": "[email protected]"
    },
    "source": {
      "version": "3.6.3.Final",
      "name": "mariadb-server-1",
      "connector": "mariadb",
      "ts_ms": 1465581029100,
      "ts_us": 1465581029100000,
      "ts_ns": 1465581029100000000,
      "snapshot": false,
      "db": "inventory",
      "table": "customers",
      "server_id": 223344,
      "gtid": null,
      "file": "mariadb-bin.000003",
      "pos": 484,
      "row": 0,
      "thread": 7,
      "query": "UPDATE customers SET first_name='Anne Marie' WHERE id=1004"
    },
    "op": "u",
    "ts_ms": 1465581029523,
    "ts_us": 1465581029523758,
    "ts_ns": 1465581029523758914
  }
}

以下列表描述了前述 update 事件消息中值部分的部分字段:

before

显示事件发生前行的状态。在 update 事件的值中,before 字段为每个表列都包含一个子字段,其值是该列在数据库提交之前的值。在本示例中,first_name 的值为 Anne。

after

显示事件发生后行的状态。你可以比较 before 和 after 结构,以确定对该行所做的更新。在示例中,first_name 的值现在是 Anne Marie。

source

表示事件的来源元数据。source 字段的结构与 create 事件中的结构相同,但其中某些值不同,例如,示例 update 事件来自 binlog 中的不同位置。update 事件值中的 source 字段提供以下元数据:

version

Debezium 版本。

name

连接器名称。

connector

连接器的类型。

ts_ms, ts_us, ts_ns

以毫秒、微秒和纳秒显示时间戳,表示更改在数据库中发生的时间。

snapshot

指定该事件是否由快照操作产生。

db、table

包含新行的数据库和表的名称。

server_id

MariaDB 服务器 ID(如果可用)。

file

记录该事件的二进制日志(binary log)名称。

pos

二进制日志位置。

row

事件内的行号。

thread

创建该事件的 MariaDB 线程 ID(仅限非快照事件)。

query

在源表上执行该操作的 SQL 命令。

如果 MariaDB 数据库配置中启用了 binlog_annotate_row_events 选项,并且你在连接器配置中启用了 include.query 属性,则 source 字段还会提供一个 query 字段,其中包含导致该变更事件的原始 SQL 语句。

op

指定操作的类型。在 update 事件的值中,op 字段的值为 u,表示该行因更新而发生变化。

ts_ms、ts_us、ts_ns

以毫秒、微秒和纳秒显示时间戳,表示连接器处理该事件的时间。该时间基于运行 Kafka Connect 任务的 JVM 中的系统时钟。

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

更新行的主键或唯一键的列会更改该行键的值。键发生变化后,Debezium 会生成以下三个事件:

  • 一个 DELETE 事件。
  • 一个墓碑事件,指定该行的旧键。
  • 一个 CREATE 事件,提供该行的新主键。

欲了解更多信息,请参阅主键更新。

主键更新

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

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

delete 事件

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

{
  "schema": { ... },
  "payload": {
    "before": {
      "id": 1004,
      "first_name": "Anne Marie",
      "last_name": "Kretchmar",
      "email": "[email protected]"
    },
    "after": null,
    "source": {
      "version": "3.6.3.Final",
      "connector": "mariadb",
      "name": "mariadb-server-1",
      "ts_ms": 1465581902300,
      "ts_us": 1465581902300000,
      "ts_ns": 1465581902300000000,
      "snapshot": false,
      "db": "inventory",
      "table": "customers",
      "server_id": 223344,
      "gtid": null,
      "file": "mariadb-bin.000003",
      "pos": 805,
      "row": 0,
      "thread": 7,
      "query": "DELETE FROM customers WHERE id=1004"
    },
    "op": "d",
    "ts_ms": 1465581902461,
    "ts_us": 1465581902461842,
    "ts_ns": 1465581902461842579
  }
}

以下列表描述了前面 delete 事件消息值部分中的部分字段:

before

指定事件发生前行的状态。在 delete 事件值中,before 字段包含该行在随数据库提交被删除之前的值。

after

指定事件发生后行的状态。在 delete 事件值中,after 字段为 null,表示该行已不存在。

source

指定事件的源元数据。在 delete 事件值中,source 字段的结构与同一表的 create 和 update 事件相同,许多 source 字段的值也相同。在 delete 事件值中,ts_ms 和 pos 字段的值以及其他值可能已发生变化。delete 事件值中的 source 字段提供以下元数据:

version

Debezium 版本。

name

连接器名称。

connector

连接器类型。

ts_ms、ts_us、ts_ns

以毫秒、微秒和纳秒为单位显示时间戳,指示更改在数据库中发生的时间。

snapshot

指定该事件是否由快照操作产生。

db、table

包含新行的数据库和表的名称。

server_id

MariaDB 服务器 ID(如果可用)。

file

记录该事件的二进制日志文件名。

pos

二进制日志位置。

row

事件中的行。

thread

创建该事件的 MariaDB 线程 ID(仅限非快照事件)。

query

在源表上执行该操作的 SQL 命令。

如果在 MariaDB 数据库配置中启用了 binlog_annotate_row_events 选项,并且在连接器配置中启用了 include.query 属性,那么 source 字段还会提供一个 query 字段,其中包含导致该变更事件的原始 SQL 语句。

op

指定操作类型。在本示例中,值 d 表示执行了 DELETE 操作,导致某行被移除。

ts_ms、ts_us、ts_ns

以毫秒、微秒和纳秒为单位显示时间戳,指示连接器处理该事件的时间。该时间基于运行 Kafka Connect 任务的 JVM 中的系统时钟。

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

delete 变更事件记录为消费者提供了处理该行移除所需的信息。之所以包含旧值,是因为某些消费者可能需要这些值才能正确处理移除操作。

MariaDB 连接器的事件设计为与 Kafka 日志压缩协同工作。日志压缩允许在保证每个键至少保留最新一条消息的前提下,删除较早的一些消息。这样 Kafka 就可以在回收存储空间的同时,确保主题中包含一份完整的数据集,可用于重新加载基于键的状态。

墓碑事件

当某一行被删除时,delete 事件的值仍能与日志压缩配合工作,因为 Kafka 可以删除所有具有相同键的较早消息。但是,要让 Kafka 删除所有具有相同键的消息,消息的值必须为 null。为此,在 Debezium MariaDB 连接器发出 delete 事件之后,连接器会发出一个特殊的墓碑事件,该事件具有相同的键,但值为 null。

truncate 事件

truncate 变更事件表示某个表已被截断。truncate 事件的消息键为 null。消息值类似于以下示例:

{
    "schema": { ... },
    "payload": {
        "source": {
            "version": "3.6.3.Final",
            "name": "mariadb-server-1",
            "connector": "mariadb",
            "name": "mariadb-server-1",
            "ts_ms": 1465581029100,
            "ts_us": 1465581029100000,
            "ts_ns": 1465581029100000000,
            "snapshot": false,
            "db": "inventory",
            "table": "customers",
            "server_id": 223344,
            "gtid": null,
            "file": "mariadb-bin.000003",
            "pos": 484,
            "row": 0,
            "thread": 7,
        },
        "op": "t",
        "ts_ms": 1465581029523,
        "ts_us": 1465581029523468,
        "ts_ns": 1465581029523468471
    }
}

以下列表说明了前面 truncate 事件消息中值部分的选定字段:

source

描述事件源元数据的必填字段。在 truncate 事件值中,source 字段的结构与同一表的 create、update 和 delete 事件相同。source 字段提供以下元数据:

version

Debezium 版本。

name

连接器名称。

connector

连接器类型。

ts_ms, ts_us, ts_ns

以毫秒、微秒和纳秒显示的时间戳,表示变更在数据库中发生的时间。

snapshot

指定该事件是否由快照操作产生。

db, table

包含新行的数据库和表的名称。

server_id

MariaDB 服务器 ID(如果可用)。

file

记录该事件的二进制日志文件名。

pos

二进制日志位置。

row

事件中的行号。

thread

创建该事件的 MariaDB 线程 ID(仅非快照事件)。

op

指定操作类型。在此示例中,值 t 表示执行了 TRUNCATE 操作,导致行被移除。op:描述操作类型的必填字符串。op 字段的值为 t,表示该表已被截断。

ts_ms, ts_us, ts_ns

以毫秒、微秒和纳秒显示的时间戳,表示连接器处理该事件的时间。该时间基于运行 Kafka Connect 任务的 JVM 中的系统时钟。

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

如果单个 TRUNCATE 语句作用于多个表,则连接器会为每个被截断的表发出一条 truncate 变更事件记录。

truncate 事件表示对整张表所做的更改,且它没有消息键。因此,对于具有多个分区的主题,与某张表相关的更改事件(create、update 等)或 truncate 事件没有顺序保证。例如,如果消费者从多个分区读取某张表的事件,它可能会在从一个分区收到删除该表所有数据的 truncate 事件之后,才从另一个分区收到该表的 update 事件。只有使用单个分区的主题才有顺序保证。

如果你不希望连接器捕获 truncate 事件,可以使用 skipped.operations 选项将其过滤掉。

数据类型映射

Debezium MariaDB 连接器用与该行所在表结构相同的事件来表示行的更改。事件中为每个列值包含一个字段。该列的 MariaDB 数据类型决定了 Debezium 在事件中如何表示该值。

存储字符串的列在 MariaDB 中通过字符集和排序规则来定义。MariaDB 连接器在读取 binlog 事件中列值的二进制表示时,会使用该列的字符集。

连接器可以将 MariaDB 数据类型映射为字面类型和语义类型。

  • 字面类型:使用 Kafka Connect 模式类型表示该值的方式。
  • 语义类型:Kafka Connect 模式如何捕获字段的含义(模式名称)。

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

基本类型

下表展示了连接器如何映射基本的 MariaDB 数据类型。

表 5. 基本类型映射说明

MariaDB 类型 字面类型 语义类型

BOOLEAN, BOOL

BOOLEAN

无

BIT(1)

BOOLEAN

无

BIT(>1)

BYTES

io.debezium.data.Bits

length 模式参数包含一个整数,表示位数。byte[] 以小端形式包含这些位,其大小足以容纳指定的位数。例如,当 n 为位数时:
numBytes = n/8 + (n%8== 0 ? 0 : 1)

TINYINT[(M)]

INT16

无

SMALLINT[(M)]

INT16

无

MEDIUMINT[(M)]

INT32

无

INT[(M)], INTEGER[(M)]

INT32

无

BIGINT[(M)]

INT64

n/a

REAL[(M,D)]

FLOAT32

n/a

FLOAT[(P)]

FLOAT32 或 FLOAT64

精度仅用于确定存储大小。精度 P 在 0 到 23 之间时,生成 4 字节的单精度 FLOAT32 列;精度 P 在 24 到 53 之间时,生成 8 字节的双精度 FLOAT64 列。

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 或 STRING

n/a

根据 binary.handling.mode 连接器配置属性的设置,可以是原始字节(默认)、base64 编码的字符串、base64 URL 安全编码的字符串,或十六进制编码的字符串。

VARBINARY(M)]

BYTES 或 STRING

n/a

根据 binary.handling.mode 连接器配置属性的设置,可以是原始字节(默认)、base64 编码的字符串、base64 URL 安全编码的字符串,或十六进制编码的字符串。

TINYBLOB

BYTES 或 STRING

n/a

根据 binary.handling.mode 连接器配置属性的设置,可以是原始字节(默认)、base64 编码的字符串、base64 URL 安全编码的字符串,或十六进制编码的字符串。

TINYTEXT

STRING

n/a

BLOB[(M)]

BYTES 或 STRING

n/a

根据 binary.handling.mode 连接器配置属性的设置,可以是原始字节(默认)、base64 编码的字符串、base64 URL 安全编码的字符串,或十六进制编码的字符串。

仅支持大小不超过 2GB 的值。建议使用凭证支取(claim check)模式将大型列值外部化。

TEXT[(M)]

STRING

n/a

仅支持大小不超过 2GB 的值。建议使用凭证支取(claim check)模式将大型列值外部化。

MEDIUMBLOB

BYTES 或 STRING

n/a

根据 binary.handling.mode 连接器配置属性的设置,可以是原始字节(默认)、base64 编码的字符串、base64 URL 安全编码的字符串,或十六进制编码的字符串。

MEDIUMTEXT

STRING

n/a

LONGBLOB

BYTES 或 STRING

n/a

根据 binary.handling.mode 连接器配置属性的设置,可以是原始字节(默认)、base64 编码的字符串、base64 URL 安全编码的字符串,或十六进制编码的字符串。

仅支持大小不超过 2GB 的值。建议使用凭证支取(claim check)模式将大型列值外部化。

LONGTEXT

STRING

n/a

仅支持大小不超过 2GB 的值。建议使用凭证支取(claim check)模式将大型列值外部化。

JSON

STRING

io.debezium.data.Json

包含 JSON 文档、数组或标量的字符串表示。

ENUM

STRING

io.debezium.data.Enum

allowed 架构参数包含以逗号分隔的允许值列表。

SET

STRING

io.debezium.data.EnumSet

allowed 架构参数包含以逗号分隔的允许值列表。

YEAR[(4)]

INT32

io.debezium.time.Year

TIMESTAMP[(fsp)]

STRING

io.debezium.time.ZonedTimestamp

采用 ISO 8601 格式,精度为微秒。MariaDB 允许小数秒精度 fsp 的取值范围为 0-6。

时间类型

除 TIMESTAMP 数据类型外,MariaDB 的时间类型取决于 time.precision.mode 连接器配置属性的值。对于默认值指定为 CURRENT_TIMESTAMP 或 NOW 的 TIMESTAMP 列,Kafka Connect 架构中使用 1970-01-01 00:00:00 作为默认值。

MariaDB 允许 DATE、DATETIME 和 TIMESTAMP 列出现零值,因为有时零值比 NULL 值更可取。当列定义允许 NULL 值时,MariaDB 连接器将零值表示为 NULL 值;当列不允许 NULL 值时,则表示为纪元日(epoch day)。

不带时区的时间值

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

  • 值为 2018-06-20 06:37:03 的 DATETIME 会变为 1529476623000。
  • 值为 2018-06-20 06:37:03 的 TIMESTAMP 会变为 2018-06-20T13:37:03Z。

此类列会根据服务器(或会话)的当前时区,转换为 UTC 中等效的 io.debezium.time.ZonedTimestamp。默认情况下,时区将从服务器查询获取。

如果此查询失败,则必须通过数据库的 timezone MariaDB 配置选项显式指定。例如,如果数据库的时区(无论是全局设置的,还是通过 timezone 选项为连接器配置的)为 "America/Los_Angeles",则 TIMESTAMP 值 "2018-06-20 06:37:03" 将由值为 "2018-06-20T13:37:03Z" 的 ZonedTimestamp 表示。

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

有关时间值相关属性的更多详情,请参阅 MariaDB 连接器配置属性 的文档。

time.precision.mode=adaptive_time_microseconds(默认)

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

表 6. time.precision.mode=adaptive_time_microseconds 时的映射关系

MariaDB 类型 字面类型 语义类型

DATE

INT32

io.debezium.time.Date
表示自纪元以来的天数。

TIME[(fsp)]

INT64

io.debezium.time.MicroTime
以微秒表示时间值,不包含时区信息。MariaDB 允许小数秒精度 fsp 的取值范围为 0-6。

DATETIME, DATETIME(0), DATETIME(1), DATETIME(2), DATETIME(3)

INT64

io.debezium.time.Timestamp
表示自纪元以来的毫秒数,不包含时区信息。

DATETIME(4), DATETIME(5), DATETIME(6)

INT64

io.debezium.time.MicroTimestamp
表示自纪元以来的微秒数,不包含时区信息。

time.precision.mode=connect

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

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

MariaDB 类型 字面类型 语义类型

DATE

INT32

org.apache.kafka.connect.data.Date
表示自纪元以来的天数。

TIME[(fsp)]

INT64

org.apache.kafka.connect.data.Time
表示自午夜起以微秒计的时间值,不包含时区信息。

DATETIME[(fsp)]

INT64

org.apache.kafka.connect.data.Timestamp
表示自纪元以来的毫秒数,不包含时区信息。

time.precision.mode=isostring

MariaDB 连接器将日期、时间和日期时间值表示为 UTC 时区下 ISO-8601 格式的字符串。由于这些值以字符串形式表示,此模式可以保留 connect 模式可能丢失的小数秒精度。

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

MariaDB 类型 字面类型 语义类型

DATE

STRING

io.debezium.time.IsoDate

按照 ISO 8601 标准,以 UTC 格式表示日期值,例如 2017-09-15Z。

TIME[(fsp)]

STRING

io.debezium.time.IsoTime

按照 ISO 8601 标准,以 UTC 格式表示时间值,例如 04:05:11.789Z。

DATETIME[(fsp)]

STRING

io.debezium.time.IsoTimestamp

以 UTC 格式(遵循 ISO 8601 标准)表示日期时间值,例如 2019-07-09T02:28:57.123456Z。

time.precision.mode 连接器属性还接受 microseconds 和 nanoseconds 这两个值,尽管上面的映射表中目前并未展示这两种模式。MariaDB 的 DATETIME 和 TIME 列精度最高只到微秒(DATETIME(0 - 6)、TIME(0 - 6)),因此 nanoseconds 相比 microseconds 并不会提供更高的精度。如果将 time.precision.mode=nanoseconds,请注意,自 Unix 纪元以来的纳秒数存储在有符号的 INT64 中,可表示的范围被限制在大约 1677-09-21T00:12:43Z 到 2262-04-11T23:47:16Z 之间。超出该范围的值——例如 9999-12-31 这样的时间尽头哨兵值——无法表示,会导致值转换器抛出 IllegalArgumentException。该异常会按照标准的 event.processing.failure.handling.mode 配置进行处理,因此运维人员可以选择 fail(默认值)、warn 或 skip。请使用 time.precision.mode=microseconds(或 isostring)以避免这一限制。

小数类型

Debezium 连接器根据 decimal.handling.mode 连接器配置属性的设置来处理小数。

decimal.handling.mode=precise

表 9. decimal.handling.mode=precise 时的映射

MariaDB 类型 字面量类型 语义类型

NUMERIC[(M[,D])]

BYTES

org.apache.kafka.connect.data.Decimal
scale 架构参数包含一个整数,表示小数点移动了多少位。

DECIMAL[(M[,D])]

BYTES

org.apache.kafka.connect.data.Decimal
scale 架构参数包含一个整数,表示小数点移动了多少位。

decimal.handling.mode=double

MariaDB 类型字面量类型语义类型
NUMERIC[(M[,D])]FLOAT64不适用
DECIMAL[(M[,D])]FLOAT64不适用

表 10. decimal.handling.mode=double 时的映射关系

decimal.handling.mode=string

MariaDB 类型字面量类型语义类型
NUMERIC[(M[,D])]STRING不适用
DECIMAL[(M[,D])]STRING不适用

表 11. decimal.handling.mode=string 时的映射关系

布尔值

MariaDB 以特定的方式在内部处理 BOOLEAN 值。BOOLEAN 列在内部映射为 TINYINT(1) 数据类型。在流处理期间创建表时,由于 Debezium 能接收到原始的 DDL,因此会使用正确的 BOOLEAN 映射。而在执行快照时,Debezium 会执行 SHOW CREATE TABLE 来获取表定义,对于 BOOLEAN 列和 TINYINT(1) 列,返回的都是 TINYINT(1)。此时 Debezium 无法获知原始的类型映射,因此会将其映射为 TINYINT(1)。

为了便于你将源列转换为布尔数据类型,Debezium 提供了 TinyIntOneToBooleanConverter 自定义转换器,你可以通过以下任一方式使用它:

  • 将所有 TINYINT(1) 或 TINYINT(1) UNSIGNED 列映射为 BOOLEAN 类型。

  • 使用以逗号分隔的正则表达式列表来指定要转换的部分列。
    要使用这种转换方式,你必须设置 converters 配置属性,并使用 selector 参数,如以下示例所示:

    converters=boolean
    boolean.type=io.debezium.connector.binlog.converters.TinyIntOneToBooleanConverter
    boolean.selector=db1.table1.*, db1.table2.column1
  • 注意:在某些情况下,快照执行 SHOW CREATE TABLE 时,数据库可能不会显示 tinyint unsigned 的长度,这意味着该转换器无法生效。新的 length.checker 选项可以解决此问题,其默认值为 true。请禁用 length.checker,并使用 selected 属性指定需要转换的列,而不是根据类型转换所有列,如以下示例所示:

    converters=boolean
    boolean.type=io.debezium.connector.binlog.converters.TinyIntOneToBooleanConverter
    boolean.length.checker=false
    boolean.selector=db1.table1.*, db1.table2.column1

空间类型

目前,Debezium MariaDB 连接器支持以下空间数据类型。

表 12. 空间类型映射说明

MariaDB 类型字面量类型语义类型
GEOMETRY
LINESTRING
POLYGON
MULTIPOINT
MULTILINESTRING
MULTIPOLYGON
GEOMETRYCOLLECTION
STRUCTio.debezium.data.geometry.Geometry,包含具有以下两个字段的结构:

- srid (INT32):空间参考系统 ID,用于定义结构中存储的几何对象类型
- wkb (BYTES):以 Well-Known-Binary(wkb)格式编码的几何对象二进制表示。更多详情请参阅开放地理空间联盟。

向量类型

向量类型

目前,Debezium MariaDB 连接器支持以下向量数据类型。

MariaDB 类型字面量类型语义类型
VECTORARRAY (FLOAT32)io.debezium.data.FloatVector

表 13. 向量类型映射说明

自定义转换器

默认情况下,Debezium MariaDB 连接器为 MariaDB 数据类型提供了多个 CustomConverter 实现。这些自定义转换器根据连接器配置,为特定数据类型提供替代映射。要向连接器添加 CustomConverter,请按照自定义转换器文档中的说明进行操作。

TINYINT(1) 转换为布尔值

默认情况下,在连接器快照期间,Debezium MariaDB 连接器从 JDBC 驱动程序获取列类型,该驱动程序会将 TINYINT(1) 类型分配给 BOOLEAN 列。Debezium 随后使用这些 JDBC 列类型来定义快照事件的架构。当连接器从快照阶段过渡到流式传输阶段后,默认映射所产生的更改事件架构可能导致 BOOLEAN 列的映射不一致。为帮助确保 MariaDB 统一发出 BOOLEAN 列,你可以应用自定义的 TinyIntOneToBooleanConverter,如以下配置示例所示。

示例:TinyIntOneToBooleanConverter 配置

converters=tinyint-one-to-boolean
tinyint-one-to-boolean.type=io.debezium.connector.binlog.converters.TinyIntOneToBooleanConverter
tinyint-one-to-boolean.selector=.*.MY_TABLE.DATA
tinyint-one-to-boolean.length.checker=false

在前面的示例中,selector 和 length.checker 属性是可选的。默认情况下,转换器会检查 TINYINT 数据类型是否符合长度为 1 的要求。如果将 length.checker 设置为 false,转换器就不会显式确认 TINYINT 数据类型是否符合长度为 1 的要求。selector 根据提供的正则表达式指定要转换的表或列。如果省略 selector 属性,转换器会将所有 TINYINT 列映射为逻辑 BOOL 字段类型。如果你未配置 selector 选项,并且希望将 TINYINT 列映射为 TINYINT(1),则应省略 length.checker 属性,或将它的值设置为 true。

JDBC sink 数据类型

如果你将 Debezium JDBC sink 连接器与 Debezium MariaDB 源连接器集成,MariaDB 连接器在快照阶段和流式处理阶段会以不同的方式发出某些列属性。为了让 JDBC sink 连接器能够一致地消费来自快照阶段和流式处理阶段的变更,你必须将 JdbcSinkDataTypesConverter 转换器作为 MariaDB 源连接器配置的一部分,如下例所示:

示例:JdbcSinkDataTypesConverter 配置

converters=jdbc-sink
jdbc-sink.type=io.debezium.connector.binlog.converters.JdbcSinkDataTypesConverter
jdbc-sink.selector.boolean=.*.MY_TABLE.BOOL_COL
jdbc-sink.selector.real=.*.MY_TABLE.REAL_COL
jdbc-sink.selector.string=.*.MY_TABLE.STRING_COL
jdbc-sink.treat.real.as.double=true

在前面的示例中,selector.* 与 treat.real.as.double 配置属性是可选的。

selector.* 属性用于指定以逗号分隔的正则表达式列表,这些表达式指明转换器所应用的表和列。默认情况下,转换器会将以下规则应用于所有表中的布尔型、实数型和基于字符串的列数据类型:

  • BOOLEAN 数据类型始终以 INT16 逻辑类型输出,其中 1 表示 true,0 表示 false
  • REAL 数据类型始终以 FLOAT64 逻辑类型输出。
  • 基于字符串的列始终包含 __debezium.source.column.character_set 架构参数,该参数中存放列的字符集。

对于每种数据类型,你都可以配置一个选择器规则,以覆盖默认的作用范围,使选择器仅应用于特定的表和列。例如,要设置布尔转换器的作用范围,可在连接器配置中添加如下规则(如前面的示例所示):converters.jdbc-sink.selector.boolean=.*.MY_TABLE.BOOL_COL

零日期回退

默认情况下,MySQL/MariaDB 允许插入 0000-00-00(DATE)和 0000-00-00 00:00:00(DATETIME / TIMESTAMP)值,这些是哨兵式"零日期"值,并不表示任何真实的时间点。Debezium 无法在其标准时间逻辑类型(io.debezium.time.Date、io.debezium.time.Timestamp、io.debezium.time.ZonedTimestamp)中表示这些哨兵值,因此默认情况下会将其转换为 null(当列可为空时),或者在列为 NOT NULL 时转换为纪元时间(1970-01-01[T00:00:00Z])。

ZeroDateFallbackConverter 自定义转换器允许你按 JDBC 类型或按列指定连接器为零日期行输出的值。你也可以选择为 NOT NULL 列输出 null —— 转换器会将该列的 Kafka Connect 架构提升为可选,并清空其 defaultValue 槽位,这样在 Avro 序列化时,null 行就不会被列的 DDL 默认值回填。

示例:ZeroDateFallbackConverter 配置

converters=zero-date
zero-date.type=io.debezium.connector.binlog.converters.ZeroDateFallbackConverter
zero-date.target.types=DATE,DATETIME,TIMESTAMP
zero-date.fallback.date=NULL
zero-date.fallback.datetime=2000-01-01 00:00:00
zero-date.fallback.timestamp=1970-01-01T00:00:00Z
zero-date.column.appdb.users.deleted_at.fallback=NULL
zero-date.column.appdb.orders.shipped_at.fallback=2099-12-31 23:59:59
zero-date.selector=.*
属性默认值描述
target.typesDATE,DATETIME,TIMESTAMP以逗号分隔的 DATE、DATETIME、TIMESTAMP 子集,用于控制转换器适用于哪些 JDBC 时间类型。
fallback.dateNULL为零值日期的 DATE 行输出的值。设为 NULL(或不设置)时,会将列模式提升为可选并输出 null。否则输出一个 ISO-8601 日期字面量(yyyy-MM-dd)。
fallback.datetimeNULL为零值日期的 DATETIME 行输出的值。设为 NULL(或不设置)时,会将列模式提升为可选并输出 null。否则输出一个 UTC 日期时间字面量(yyyy-MM-dd HH:mm:ss)。
fallback.timestampNULL为零值日期的 TIMESTAMP 行输出的值。设为 NULL(或不设置)时,会将列模式提升为可选并输出 null。否则输出一个 ISO-8601 带偏移量的日期时间字面量(例如 1970-01-01T00:00:00Z)。
column.<db>.<table>.<column>.fallback(未设置)针对特定列的回退值,覆盖类型级别的回退设置。允许的取值与类型级别选项相同(NULL 或与类型相匹配的字面量)。查找键使用 RelationalColumn#dataCollection() 返回的完全限定列名。
selector.*以逗号分隔的正则表达式列表,与 <db>.<table>.<column> 进行匹配。只有匹配的列才会交由转换器处理。默认通配符会将转换器应用于所有 JDBC 类型包含在 target.types 中的列。

生效的回退值的解析顺序为:列级覆盖 > 类型级回退 > NULL。

该转换器依赖 binlog 反序列化器 / JDBC 驱动在 SPI 输入处将零日期值以 null 形式呈现,从而区分零日期行与真实的 1970-01-01 行。因此,在行级输出上不会产生误报:真正的 1970-01-01 00:00:00 UTC 行会像未使用该转换器时一样,继续输出 0(epoch 天数 / 毫秒),而零日期行则输出所配置的回退值(或 null)。

该转换器会在构建 schema 时,将用户配置的回退值写入 Kafka Connect schema 的 defaultValue 槽位,而不管该列在 DDL 中的默认值是什么。这依赖于 TableSchemaBuilder 在构建 schema 期间,为了解析该列的默认值而恰好调用一次已注册的转换器 lambda。这是一个内部约定,而非对外公布的 SPI 保证。

设置 MariaDB

在安装和运行 Debezium 连接器之前,需要完成一些 MariaDB 的设置工作。

创建用户

Debezium MariaDB 连接器需要一个 MariaDB 用户账户。Debezium MariaDB 连接器为其捕获变更的所有数据库,都要求该 MariaDB 用户拥有相应的权限。

前置条件

  • 一个 MariaDB 服务器。
  • SQL 命令的基本知识。

操作步骤

  1. 创建 MariaDB 用户:

    mariadb> CREATE USER 'user'@'localhost' IDENTIFIED BY 'password';
  2. 为该用户授予权限:

    mariadb> GRANT SELECT, RELOAD, SHOW DATABASES, REPLICATION SLAVE, REPLICATION CLIENT ON *.* TO 'user' IDENTIFIED BY 'password';

有关所需权限的说明,请参阅用户权限说明。

如果使用不允许全局读锁的托管服务(如 Amazon RDS 或 Amazon Aurora),则会使用表级锁来创建一致性快照。在这种情况下,你还需要为你创建的用户授予 LOCK TABLES 权限。更多详情请参阅快照。
3. 完成用户权限的设置:
mariadb> FLUSH PRIVILEGES;

用户权限说明

下表介绍了必须授予 Debezium 连接器用户账户的 MariaDB 权限,并说明了为何需要每项权限。

表 14. 用户权限说明

关键字 说明

SELECT

使连接器能够从数据库的表中查询行数据。仅在执行快照时使用。

RELOAD

使连接器能够使用 FLUSH 语句清除或重新加载内部缓存、刷新表或获取锁。仅在执行快照时使用。

SHOW DATABASES

使连接器能够通过执行 SHOW DATABASE 语句查看数据库名称。仅在执行快照时使用。

REPLICATION SLAVE

使连接器能够连接并读取 MariaDB 服务器的 binlog。

REPLICATION CLIENT

使连接器能够使用以下语句:

  • SHOW MASTER STATUS
  • SHOW SLAVE STATUS
  • SHOW BINARY LOGS

连接器始终需要此权限。

ON

指定这些权限所适用的数据库。

TO 'user'

指定要授予权限的用户。

IDENTIFIED BY 'password'

指定该用户的 MariaDB 密码。

启用 binlog

你必须为 MariaDB 复制启用二进制日志记录。二进制日志以一种使副本能够传播这些更改的方式记录事务更新。

前提条件

  • 一个 MariaDB 服务器。
  • 适当的 MariaDB 用户权限。

操作步骤

  1. 检查 log-bin 选项是否已启用:

    mariadb> SHOW VARIABLES LIKE '%log_bin%';
  2. 如果 binlog 为 OFF,请将下表中的属性添加到 MariaDB 服务器的配置文件中:

    server-id         = 223344 # 查询变量名为 server_id,例如 SELECT variable_value FROM information_schema.global_variables WHERE variable_name='server_id';
    log_bin                     = mariadb-bin
    binlog_format               = ROW
    binlog_row_image            = FULL
    binlog_expire_logs_seconds  = 864000
  3. 再次检查 binlog 状态,以确认您的更改:

    mariadb> SHOW VARIABLES LIKE '%log_bin%';
  1. 如果你在 Amazon RDS 上运行 MariaDB,必须为数据库实例启用自动备份,二进制日志才会生成。如果数据库实例未配置为执行自动备份,即使你应用了前述步骤中描述的设置,binlog 也处于禁用状态。

MariaDB binlog 配置属性说明

下表介绍了必须设置的 MariaDB 服务器配置属性,这些属性用于启用二进制日志并将其配置为与 Debezium 连接器配合使用。

属性说明
server-idserver-id 的值在 MariaDB 集群中的每个服务器和每个复制客户端上都必须是唯一的。
log_binlog_bin 的值是 binlog 文件序列的基础名称。
binlog_formatbinlog-format 必须设置为 ROW 或 row。
binlog_row_imagebinlog_row_image 必须设置为 FULL 或 full。
binlog_expire_logs_secondsbinlog_expire_logs_seconds 对应于已废弃的系统变量 expire_logs_days。它表示自动删除 binlog 文件的秒数。默认值为 2592000,即 30 天。请根据环境的需要设置该值。更多信息请参阅 MariaDB 清除 Debezium 使用的 binlog 文件。
log_bin_compress二进制日志是否可以被压缩。Debezium 的 MariaDB 连接器不支持压缩的二进制日志条目,因此 log_bin_compress 必须设置为 0(默认值),即不进行压缩。

表 15. MariaDB binlog 配置属性说明

MariaDB 11.4 引入了新的二进制日志格式。当您将 Debezium 与 MariaDB 11.4 或更高版本一起使用时,必须将 MariaDB 服务器变量 binlog_legacy_event_pos 设置为 1(ON),以确保连接器能够消费新格式的事件。如果该变量保持默认设置(0,即 OFF),那么在连接器重启后,Debezium 可能无法找到恢复点。

Debezium MariaDB 连接器目前不支持二进制日志压缩。为使连接器能够消费所有事件,必须将 MariaDB 服务器配置选项 log_bin_compress 设置为 0。

启用 GTID

全局事务标识符(GTID)唯一标识集群中服务器上发生的事务。虽然 Debezium MariaDB 连接器并不强制要求使用 GTID,但使用 GTID 可以简化复制过程,并使您更轻松地确认主服务器和副本服务器是否一致。

对于 MariaDB,GTID 默认已启用,无需进行额外配置。

配置会话超时

对大型数据库执行初始一致快照时,在读取表的过程中,已建立的连接可能会超时。您可以通过在 MariaDB 配置文件中配置 interactive_timeout 和 wait_timeout 来防止这种情况发生。

前提条件

  • 一个 MariaDB 服务器。
  • SQL 命令的基本知识。
  • 访问 MariaDB 配置文件的权限。

操作步骤

  1. 配置 interactive_timeout:

    mariadb> interactive_timeout=<duration-in-seconds>
  2. 配置 wait_timeout:

    mariadb> wait_timeout=<duration-in-seconds>

表 16. MariaDB 会话超时选项说明

选项 说明
interactive_timeout 服务器在关闭交互式连接之前等待该连接活动的秒数。 +1
有关更多信息,请参阅 MariaDB 文档。
wait_timeout 服务器在关闭非交互式连接之前等待该连接活动的秒数。有关更多信息,请参阅 MariaDB 文档。

启用查询日志事件

你可能希望查看每个 binlog 事件对应的原始 SQL 语句。在 MariaDB 配置中启用 binlog_annotate_row_events 选项即可实现此目的。

前提条件

  • 一个 MariaDB 服务器。
  • SQL 命令的基础知识。
  • 访问 MariaDB 配置文件的权限。

操作步骤

  • 在 MariaDB 中启用 binlog_annotate_row_events:

    mariadb> binlog_annotate_row_events=ON

binlog_annotate_row_events 被设置为某个值,用于启用/禁用在 binlog 条目中包含原始 SQL 语句的支持。

  • ON = 启用
  • OFF = 禁用

验证 binlog 行值选项

在数据库中检查 binlog_row_value_options 变量的设置。要使连接器能够消费 UPDATE 事件,必须将该变量设置为 PARTIAL_JSON 以外的值。

前置条件

  • 一个 MariaDB 服务器。
  • 基本的 SQL 命令知识。
  • 访问 MariaDB 配置文件的权限。

操作步骤

  1. 检查当前变量值

    mariadb> show global variables where variable_name = 'binlog_row_value_options';

结果

+--------------------------+-------+
| Variable_name            | Value |
+--------------------------+-------+
| binlog_row_value_options |       |
+--------------------------+-------+
  1. 如果该变量的值设置为 PARTIAL_JSON,请运行以下命令将其取消设置:

    mariadb> set @@global.binlog_row_value_options="" ;

部署

要部署 Debezium MariaDB 连接器,需要安装 Debezium MariaDB 连接器归档文件、配置连接器,并将其配置添加到 Kafka Connect 中以启动连接器。

前提条件

操作步骤

  1. 下载 Debezium MariaDB 连接器插件。
  2. 将文件解压到您的 Kafka Connect 环境中。
  3. 将包含 JAR 文件的目录添加到 Kafka Connect 的 plugin.path 中。
  4. 配置连接器,并将配置添加到您的 Kafka Connect 集群中。
  5. 重启 Kafka Connect 进程,以加载新的 JAR 文件。

如果您使用的是不可变容器,请参阅 Debezium 的容器镜像,其中提供了预装 MariaDB 连接器、可直接运行的 Apache Kafka、MariaDB 和 Kafka Connect 镜像。

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

您还可以在 Kubernetes 和 OpenShift 上运行 Debezium。

MariaDB 连接器配置示例

以下示例展示了一个连接器实例的配置,该连接器从位于 192.168.99.100、端口为 3306 的 MariaDB 数据库服务器中捕获数据,我们将其逻辑命名为 fullfillment。通常,您通过设置连接器可用的配置属性,在 JSON 文件中配置 Debezium MariaDB 连接器。

您可以选择仅为数据库中部分模式和表生成事件。此外,还可以忽略、掩码或截断包含敏感数据的列、超出指定大小的列,或下游应用程序不需要的列。

{
    "name": "inventory-connector",
    "config": {
        "connector.class": "io.debezium.connector.{context}.{connector-name}Connector",
        "database.hostname": "192.168.99.100",
        "database.port": "3306",
        "database.user": "debezium-user",
        "database.password": "debezium-user-pw",
        "database.server.id": "184054",
        "topic.prefix": "fullfillment",
        "database.include.list": "inventory",
        "schema.history.internal.kafka.bootstrap.servers": "kafka:9092",
        "schema.history.internal.kafka.topic": "schemahistory.fullfillment",
        "include.schema.changes": "true"
    }
}

以下列表描述了前面配置示例中的各个字段:

name

指定要向 Kafka Connect 服务注册的连接器名称(inventory-connector)。

connector.class

指定 Java 连接器类的名称。

database.hostname

指定 MariaDB 服务器的地址或主机名。

database.port

指定连接器用于访问 MariaDB 服务器的端口号。

database.user

指定连接器用于访问数据库的 MariaDB 用户账户名称。指定的账户必须具有足够的权限,如用户权限说明表中所述。

database.password

指定 Debezium 用于访问数据库的 MariaDB 用户账户的密码。

database.server.id

指定连接器的唯一 ID。

topic.prefix

指定 MariaDB 服务器或集群的主题前缀。

database.include.list

指定你希望连接器从指定服务器上捕获的数据库。

schema.history.internal.kafka.bootstrap.servers

指定连接器用于向数据库 schema 历史主题写入和恢复 DDL 语句的 Kafka broker。

schema.history.internal.kafka.topic

指定数据库 schema 历史主题的名称。该主题仅供内部使用,消费者不应使用。

include.schema.changes

指定连接器是否应为 DDL 变更生成事件,并将其发布到 fulfillment schema 变更主题,以供消费者使用。

有关可为 Debezium MariaDB 连接器设置的完整配置属性列表,请参阅 MariaDB 连接器配置属性。

你可以通过 POST 命令将此配置发送到正在运行的 Kafka Connect 服务。该服务会记录配置,并启动一个执行以下操作的连接器任务:

  • 连接到 MariaDB 数据库。
  • 读取处于捕获模式的表的变更数据表。
  • 将变更事件记录流式传输到 Kafka 主题。

添加连接器配置

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

前提条件

操作步骤

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

结果

连接器启动后,会对为其配置的 MariaDB 数据库执行一致性快照。随后,连接器开始针对行级操作生成数据变更事件,并将变更事件记录流式传输到 Kafka 主题。

连接器属性

Debezium MariaDB 连接器具有众多配置属性,你可以使用它们来为应用实现合适的连接器行为。许多属性都有默认值。

MariaDB 连接器配置属性的相关信息组织如下:

必需的 Debezium MariaDB 连接器配置属性

以下列表说明了部署 MariaDB 连接器所需的配置属性,除非另有默认值说明。

bigint.unsigned.handling.mode

默认值

long

描述

指定连接器在变更事件中如何表示 BIGINT UNSIGNED 列。

设置以下选项之一:

long

使用 Java long 数据类型表示 BIGINT UNSIGNED 列的值。虽然 long 类型不提供最高的精度,但在大多数消费者中都易于实现。在大多数环境中,这是首选设置。

precise

使用 java.math.BigDecimal 数据类型表示值。连接器使用 Kafka Connect 的 org.apache.kafka.connect.data.Decimal 数据类型,以编码二进制格式表示值。如果连接器通常处理大于 2^63 的值,请设置此选项。long 数据类型无法表示该量级的值。

binary.handling.mode

默认值

bytes

描述

指定连接器在变更事件中如何表示二进制列(如 blob、binary、varbinary)的值。

可设置以下选项之一:

bytes

以字节数组的形式表示二进制数据。

base64

以 base64 编码的字符串形式表示二进制数据。

base64-url-safe

以 base64-url-safe 编码的字符串形式表示二进制数据。

hex

以十六进制(base16)编码的字符串形式表示二进制数据。

column.exclude.list

默认值

空字符串

描述

一个可选的、以逗号分隔的正则表达式列表,用于匹配要从变更事件记录值中排除的列的完全限定名称。源记录中的其他列照常被捕获。列的完全限定名称格式为 databaseName.tableName.columnName。

为了匹配列名,Debezium 会将你指定的正则表达式作为锚定正则表达式来应用。也就是说,指定的表达式会与列的整个名称字符串进行匹配,而不会匹配列名中可能出现的子字符串。如果你在配置中包含此属性,请不要再设置 column.include.list 属性。

column.include.list

默认值

空字符串

描述

一个可选的、以逗号分隔的正则表达式列表,用于匹配要包含在变更事件记录值中的列的完全限定名称。其他列将从事件记录中省略。列的完全限定名称格式为 databaseName.tableName.columnName。

为了匹配列名,Debezium 会将你指定的正则表达式作为锚定正则表达式来应用。也就是说,指定的表达式会与列的整个名称字符串进行匹配,而不会匹配列名中可能出现的子字符串。
如果你在配置中包含此属性,请不要设置 column.exclude.list 属性。

column.mask.hash.hashAlgorithm.with.salt.salt

column.mask.hash.v2.hashAlgorithm.with.salt.salt

默认值

无默认值

描述

一个可选的、以逗号分隔的正则表达式列表,用于匹配字符类型列的完全限定名称。列的完全限定名称格式为 <databaseName>.<tableName>.<columnName>。

为了匹配列的名称,Debezium 会将你指定的正则表达式作为锚定正则表达式来应用。也就是说,指定的表达式会与列的整个名称字符串进行匹配;该表达式不会匹配列名称中可能存在的子字符串。在生成的更改事件记录中,所指定列的值会被替换为化名列(pseudonym)。

化名由应用指定的 hashAlgorithm(哈希算法)和 salt(盐值)后生成的哈希值组成。所使用的哈希函数可以在保持引用完整性的同时,将列值替换为化名。支持的哈希函数详见 Java 加密架构标准算法名称文档中的 MessageDigest 章节。

在下面的示例中,CzQMA0cB5K 是随机选取的盐值。

column.mask.hash.SHA-256.with.salt.CzQMA0cB5K = inventory.orders.customerName, inventory.shipment.customerName

如有需要,伪名会自动缩短到与列长度一致。连接器配置可以包含多个属性,分别指定不同的哈希算法和盐值。

根据所使用的 hashAlgorithm、所选的 salt 以及实际数据集,最终得到的数据集可能并未被完全屏蔽。

哈希策略版本 2 可确保在不同位置或系统中散列的值保持一致。

column.mask.with.length.chars

默认值

无默认值

说明

一个可选的、以逗号分隔的正则表达式列表,用于匹配基于字符的列的完全限定名称。如果你希望连接器屏蔽某一组列的值(例如这些列包含敏感数据),请设置此属性。将 length 设置为正整数,即可用属性名称中 length 指定的星号(*)字符数量替换指定列中的数据。将 length 设置为 0(零),则用空字符串替换指定列中的数据。

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

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

column.propagate.source.type

默认值

无默认值

说明

一个可选的、以逗号分隔的正则表达式列表,用于匹配你希望连接器为其发出表示列元数据的额外参数的列的完全限定名称。设置此属性后,连接器会向事件记录的模式中添加以下字段:

  • __debezium.source.column.type

  • __debezium.source.column.length

  • __debezium.source.column.scale

    这些参数分别传递列的原始类型名称和长度(针对可变宽度类型)。

    允许连接器发出这些额外数据,有助于在目标数据库中正确设置特定数值列或字符列的大小。

    列的完全限定名称遵循以下格式之一:databaseName.tableName.columnName,或 databaseName.schemaName.tableName.columnName。

    为匹配列名,Debezium 会将你指定的正则表达式作为锚定正则表达式来应用。也就是说,指定的表达式要与列的整个名称字符串匹配;它不会匹配列名中可能存在的子字符串。

column.truncate.to.length.chars

默认值

无默认值

说明

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

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

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

connect.timeout.ms

默认值

30000(30 秒)

说明

一个正整数值,用于指定连接器在连接请求超时之前,等待与 MariaDB 数据库服务器建立连接的最长时间(以毫秒为单位)。

connector.class

默认值

无默认值

说明

连接器的 Java 类名。对于 MariaDB 连接器,请始终指定 io.debezium.connector.mariadb.MariaDbConnector。

database.exclude.list

默认值

空字符串

说明

一个可选的、以逗号分隔的正则表达式列表,用于匹配你希望连接器不捕获其变更的数据库名称。连接器会捕获所有未列入 database.exclude.list 的数据库中的变更。

为了匹配数据库名称,Debezium 会将你指定的正则表达式作为锚定正则表达式来应用。也就是说,指定的表达式会与数据库的整个名称字符串进行匹配,而不会匹配数据库名称中可能存在的子字符串。

如果你在配置中包含此属性,请勿同时设置 database.include.list 属性。

database.hostname

默认值

无默认值

说明

MariaDB 数据库服务器的 IP 地址或主机名。

database.include.list

默认值

空字符串

说明

一个可选的、以逗号分隔的正则表达式列表,用于匹配连接器要从中捕获变更的数据库名称。对于名称不在 database.include.list 中的数据库,连接器不会捕获其变更。默认情况下,连接器会捕获所有数据库中的变更。

为了匹配数据库的名称,Debezium 会将你指定的正则表达式作为锚定正则表达式来应用。也就是说,指定的表达式会与数据库的整个名称字符串进行匹配;它不会匹配数据库名称中可能出现的子字符串。

如果在配置中包含此属性,请不要同时设置 database.exclude.list 属性。

database.password

默认值

无默认值

描述

连接器用于连接 MariaDB 数据库服务器的 MariaDB 用户的密码。

database.port

默认值

3306

描述

MariaDB 数据库服务器的整数端口号。

database.server.id

默认值

无默认值

描述

此数据库客户端的数字 ID。指定的 ID 必须在 MariaDB 集群中所有当前运行的数据库进程中唯一。为了能够读取 binlog,连接器会使用这个唯一的 ID 作为另一台服务器加入 MariaDB 数据库集群。

database.user

默认值

无默认值

描述

连接器用于连接 MariaDB 数据库服务器的 MariaDB 用户的名称。

decimal.handling.mode

默认值

precise

描述

指定连接器如何处理变更事件中 DECIMAL 和 NUMERIC 列的值。

可以设置以下选项之一:

precise

使用二进制形式的 java.math.BigDecimal 值来精确表示值。

double

使用 double 数据类型来表示值。此选项可能会导致精度损失,但对大多数使用者来说更易于使用。

string

将值编码为格式化的字符串。此选项易于使用,但可能导致关于真实类型语义信息的丢失。

event.deserialization.failure.handling.mode 已弃用

默认值

fail

描述

指定连接器在 binlog 事件反序列化过程中发生异常后的响应方式。

此选项已弃用。

请改用 event.processing.failure.handling.mode 属性。

此属性接受以下选项:

fail

传播该异常(异常中会指出有问题的事件及其 binlog 偏移量),并导致连接器停止。

warn

记录有问题的事件及其 binlog 偏移量,然后跳过该事件。

ignore

跳过有问题的事件,且不记录任何内容。

field.name.adjustment.mode

默认值

无默认值

描述

指定应如何调整字段名称,以兼容连接器所使用的消息转换器。

设置以下选项之一:

none

不做任何调整。

avro

将 Avro 名称中无效的字符替换为下划线字符。

avro_unicode

将下划线字符或不能用于 Avro 名称的字符替换为相应的 unicode 表示形式,例如 _uxxxx。

下划线字符(_)表示转义序列,类似于 Java 中的反斜杠

更多信息请参阅:Avro 命名。

gtid.source.excludes

默认值

无默认值

说明

以逗号分隔的正则表达式列表,用于匹配连接器在 MariaDB 服务器上查找 binlog 位置时所用 GTID 集合中的源域 ID。设置此属性后,连接器仅使用其源 UUID 不匹配任何指定 exclude 模式的 GTID 范围。

为了匹配 GTID 的值,Debezium 会将你指定的正则表达式作为锚定正则表达式来应用。也就是说,指定的表达式会与 GTID 的域标识符进行匹配。

如果设置了此属性,请勿同时设置 gtid.source.includes 属性。

gtid.source.includes

默认值

无默认值

说明

以逗号分隔的正则表达式列表,用于匹配连接器在 MariaDB 服务器上查找 binlog 位置时所用 GTID 集合中的源域 ID。设置此属性后,连接器仅使用其源 UUID 匹配某个指定 include 模式的 GTID 范围。

为了匹配 GTID 的值,Debezium 会将你指定的正则表达式作为锚定正则表达式来应用。也就是说,指定的表达式会与 GTID 的域标识符进行匹配。

如果设置了此属性,请勿同时设置 gtid.source.excludes 属性。

gtid.ignore.on.recovery

默认值

false

说明

布尔值,用于指定连接器在恢复过程中是否忽略已存储偏移量中的 GTID 集合。设置为 true 时,连接器将从已存储的 binlog 文件和位置开始,而不尝试基于 GTID 的定位。

流式传输恢复后,连接器会正常捕获 GTID,并在后续偏移量中存储刷新后的 GTID 状态。

include.query

默认值

false

说明

布尔值,用于指定连接器发出的变更事件是否包含生成该变更的 SQL 查询。

将此属性设置为 true 可能会暴露你通过其他设置明确排除或屏蔽的表或字段的信息。

要启用此属性,必须将数据库属性 binlog_annotate_row_events 设置为 ON。

设置此属性对快照过程生成的事件没有影响。快照事件不包含原始 SQL 查询。

有关如何配置数据库以在每个日志事件中返回原始 SQL 语式的更多信息,请参阅启用查询日志事件。

include.schema.changes

默认值

true

说明

一个布尔值,用于指定连接器是否将数据库架构的变更发布到与主题前缀同名的 Kafka 主题中。连接器以包含数据库名称的键和描述架构更新的 JSON 结构作为值来记录每次架构变更。这种记录架构变更的机制与连接器内部对数据库架构历史变更的记录是相互独立的。

inconsistent.schema.handling.mode

默认值

fail

说明

指定当 binlog 事件引用的表不存在于内部架构表示中时(即内部表示与数据库不一致),连接器应如何响应。

可设置为以下选项之一:

fail

连接器抛出一个异常,报告有问题的事件及其 binlog 偏移量,然后停止运行。

warn

连接器记录有问题的事件及其 binlog 偏移量,然后跳过该事件。

skip

连接器跳过有问题的事件,并且不在日志中报告。

message.key.columns

默认值

无默认值

说明

一个表达式列表,用于指定连接器为发布到指定表对应的 Kafka 主题的变更事件记录构建自定义消息键时所使用的列。默认情况下,Debezium 使用表的主键列作为其发出的记录的消息键。若要替换默认设置,或者为没有主键的表指定键,你可以基于一个或多个列配置自定义消息键。

要为表建立自定义消息键,请列出该表,后接用作消息键的列。每个列表条目采用以下格式:

<fully-qualified_tableName>:<keyColumn>,<keyColumn>

若要基于多个列名构建表键,请在各列名之间插入逗号。

每个完全限定的表名都是采用以下格式的正则表达式:

<databaseName>.<tableName>

该属性可以包含多个表的条目。在列表中使用分号分隔各个表条目。

以下示例为 inventory.customers 和 purchase.orders 表设置消息键:

inventory.customers:pk1,pk2;(.*).purchaseorders:pk3,pk4

对于 inventory.customer 表,pk1 和 pk2 列被指定为消息键。对于任何数据库中的 purchaseorders 表,pk3 和 pk4 列被用作消息键。

用于创建自定义消息键的列数没有限制。不过,最好只使用指定唯一键所需的最少列数。

name

默认值

无默认值

描述

连接器的唯一名称。如果尝试使用相同的名称注册多个连接器,注册将会失败。所有 Kafka Connect 连接器都要求设置此属性。

schema.name.adjustment.mode

默认值

无默认值

描述

指定连接器如何调整架构名称以兼容连接器所使用的消息转换器。

可设置以下选项之一:

none

不作任何调整。

avro

将 Avro 名称中无效的字符替换为下划线字符。

avro_unicode

将下划线字符或无法用于 Avro 名称的字符替换为相应的 unicode,例如 _uxxxx。

_ 是一个转义序列,类似于 Java 中的反斜杠

skip.messages.without.change

默认值

false

描述

指定当连接器未检测到所包含列发生变化时,是否仍为记录发出消息。如果列列在 column.include.list 中,或未列在 column.exclude.list 中,则该列被视为已包含。将此值设置为 true,可在所包含的列没有任何变化时阻止连接器捕获记录。

table.exclude.list

默认值

空字符串

描述

一个可选的、以逗号分隔的正则表达式列表,用于匹配你不想让连接器从中捕获更改的表的完全限定表标识符。连接器会捕获 table.exclude.list 中未包含的任何表的更改。每个标识符的形式为 databaseName.tableName。

为了匹配表的名称,Debezium 会将你指定的正则表达式作为锚定正则表达式应用。也就是说,指定的表达式会与表的整个名称字符串进行匹配;它不会匹配表名中可能出现的子字符串。

如果设置了此属性,请不要同时设置 table.include.list 属性。

table.include.list

默认值

空字符串

描述

一个可选的、以逗号分隔的正则表达式列表,用于匹配你想要捕获其更改的表的完全限定表标识符。连接器不会捕获未包含在 table.include.list 中的任何表的更改。每个标识符的形式为 databaseName.tableName。默认情况下,连接器会捕获其配置为捕获更改的所有数据库中全部非系统表的更改。

为了匹配表名,Debezium 会将你指定的正则表达式作为锚定正则表达式来应用。也就是说,指定的表达式会与表的整个名称字符串进行匹配;它不会匹配表名中可能存在的子字符串。

如果你设置了此属性,请不要同时设置 table.exclude.list 属性。

tasks.max

默认值

1

描述

要为该连接器创建的最大任务数。由于 MariaDB 连接器始终使用单个任务,更改默认值不会产生任何效果。

time.precision.mode

默认值

adaptive_time_microseconds

描述

指定连接器用于表示时间、日期和时间戳值的精度类型。

设置以下选项之一:

adaptive_time_microseconds

连接器根据数据库列的类型,使用毫秒、微秒或纳秒精度值,完全按照数据库中的原样捕获 date、datetime 和 timestamp 值,但 TIME 类型字段除外,这些字段始终以微秒为单位捕获。

adaptive

(已弃用)连接器根据列的数据类型,使用毫秒、微秒或纳秒精度值,完全按照数据库中的原样捕获时间和时间戳值。

connect

连接器始终使用 Kafka Connect 内置的 Time、Date 和 Timestamp 表示形式来表示时间和时间戳值,无论数据库列的精度如何,均使用毫秒精度。

isostring

连接器将时间、日期和 datetime 值表示为 UTC 时区下 ISO-8601 格式的字符串,使用 io.debezium.time.IsoDate、io.debezium.time.IsoTime 和 io.debezium.time.IsoTimestamp 语义类型。

tombstones.on.delete

默认值

true

描述

指定 delete 事件之后是否跟一个墓碑事件。在源记录被删除后,连接器可以发出墓碑事件(默认行为),以便在为该主题启用日志压缩时,Kafka 能够完全删除与被删除行的键相关的所有事件。

设置以下选项之一:

true

连接器通过发出一个 delete 事件及随后的墓碑事件来表示删除操作。

false

连接器仅发出 delete 事件。

topic.prefix

默认值

无默认值

说明

一个字符串,用于指定 Debezium 捕获更改的 MariaDB 数据库服务器或集群的命名空间。由于主题前缀用于命名接收该连接器发出的所有事件的 Kafka 主题,因此主题前缀在所有连接器中保持唯一非常重要。取值只能包含字母数字字符、连字符、点和下划线。

设置此属性后,请勿更改其值。如果更改了该值,连接器重新启动后将不再继续向原有主题发出事件,而是将后续事件发往基于新值命名的主题。连接器还将无法恢复其数据库架构历史主题。

Debezium MariaDB 连接器高级配置属性

以下列表介绍了 MariaDB 连接器的高级配置属性,借助这些属性,您可以针对特定环境微调连接器的行为,包括架构历史记录处理、信号通道、事件过滤和自定义转换。这些属性的默认值很少需要更改,因此您不必在连接器配置中指定它们。

binlog.buffer.size

默认值

0

说明

binlog 读取器使用的前置缓冲区大小。默认设置 0 表示禁用缓冲。

在特定条件下,MariaDB binlog 中可能包含由 ROLLBACK 语句回滚的未提交数据。典型示例包括使用保存点,或在单个事务中混合临时表和常规表的更改。

当检测到事务开始时,Debezium 会尝试向前推进 binlog 位置,查找 COMMIT 或 ROLLBACK,从而确定是否需要流式传输该事务中的更改。binlog 缓冲区的大小定义了 Debezium 在查找事务边界时能够缓冲的最大更改数量。如果事务的大小超过缓冲区,则 Debezium 必须回退并重新读取在流式传输过程中未能放入缓冲区的事件。

此功能处于孵化阶段,欢迎提供反馈。该功能预计尚未完全完善。

binlog.net.read.timeout

默认值

0

说明

在服务器超时之前,等待从 binlog 连接完成一次读取的秒数。值为 0 表示连接器使用 MySQL 服务器的默认值。

在高延迟的网络环境中,服务器的 net_read_timeout 默认值可能过低,导致服务器以 EOFException 提前关闭 binlog 流式传输连接。增大此值可防止在 binlog 流式传输过程中发生意外断开。

此值通过 SET net_read_timeout=<值> 在 binlog 流式连接上以会话级别进行设置,不会影响服务器的全局设置。

binlog.net.write.timeout

默认值

0

说明

在服务器超时之前,等待向 binlog 连接完成一次写入的秒数。值为 0 表示连接器使用 MySQL 服务器的默认值。

当服务器通过 binlog 传输大量数据时,服务器的 net_write_timeout 默认值可能过低,导致服务器以 EOFException 关闭连接。增大此值可防止在流式传输大型事务时发生意外断开。

此值通过 SET net_write_timeout=<值> 在 binlog 流式连接上以会话级别进行设置,不会影响服务器的全局设置。

connect.keep.alive

默认值

true

说明

一个布尔值,用于指定是否使用单独的线程来保持与 MariaDB 服务器或集群的连接处于活动状态。

converters

默认值

无默认值

说明

列出一个以逗号分隔的符号名称列表,表示连接器可以使用的自定义转换器实例。例如 boolean。

要让连接器使用自定义转换器,必须设置此属性。

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

<转换器符号名称>.type

例如,

boolean.type: io.debezium.connector.binlog.converters.TinyIntOneToBooleanConverter

若要进一步控制已配置转换器的行为,可以添加一个或多个配置参数,以便向转换器传递值。要将这些额外的配置参数与某个转换器关联,需要在参数名前加上该转换器的符号名。

例如,要定义一个 selector 参数来指定 boolean 转换器处理的列子集,请添加以下属性:

boolean.selector=db1.table1.*, db1.table2.column1

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 的默认屏蔽模式;指定的值不会扩展或补充默认模式。

database.initial.statements

默认值

无默认值

描述

一个以分号分隔的 SQL 语句列表,在建立到数据库的 JDBC 连接时执行(不是读取事务日志的连接)。若要在 SQL 语句中将分号指定为字符而非分隔符,请使用两个分号(;;)。

连接器可能会自行决定建立 JDBC 连接,因此此属性仅用于配置会话参数,不用于执行 DML 语句。

database.query.timeout.ms

默认值

600000(10 分钟)

描述

指定连接器等待查询完成的时间(毫秒)。

将该值设置为 0(零)可取消超时限制。

database.ssl.keystore

默认值

无默认值

描述

可选设置,用于指定密钥库文件的位置。密钥库文件可用于客户端与 MariaDB 服务器之间的双向认证。

database.ssl.keystore.password

默认值

无默认值

描述

密钥库文件的密码。仅在配置了 database.ssl.keystore 时才指定密码。

database.ssl.mode

默认值

preferred

描述

指定连接器是否使用加密连接。

可用设置如下:

disabled

指定使用非加密连接。

preferred

如果服务器支持安全连接,连接器将建立加密连接;如果服务器不支持安全连接,连接器将回退为使用非加密连接。

required

连接器建立加密连接。如果无法建立加密连接,连接器将失败。

verify_ca

连接器的行为与设置 required 选项时相同,但它还会使用配置的证书颁发机构(CA)证书来验证服务器的 TLS 证书。如果服务器的 TLS 证书与任何有效的 CA 证书都不匹配,连接器将失败。

verify_identity

连接器的行为与设置 verify_ca 选项时相同,但它还会验证服务器证书是否与远程连接的主机匹配。

database.ssl.truststore

默认值

无默认值

描述

用于验证服务器证书的信任库文件的位置。

database.ssl.truststore.password

默认值

无默认值

描述

信任库文件的密码。用于检查信任库的完整性并解锁信任库。

ddl.parser.type

默认值

default

描述

指定用于解析 MySQL DDL 语句的 ANTLR 语法。

可设置以下选项之一:

default

(默认)使用 Oracle MySQL 语法,该语法处于积极维护状态,支持 MySQL 8.0+ 的特性。此语法基于官方 MySQL 语法规范,推荐用于生产环境。

legacy

使用 Positive Technologies(PT)MySQL 语法。此选项用于与现有部署保持向后兼容。PT 语法可能不支持所有 MySQL 8.0+ 的特性。

大多数情况下,应使用 default 语法。只有在迁移过程中遇到现有 DDL 语句的兼容性问题时,才使用 legacy 语法。

enable.time.adjuster

默认值

true

说明

布尔值,指示连接器是否将 2 位年份表示转换为 4 位年份表示。如果希望完全由数据库执行转换,请将该值设置为 false。

MariaDB 用户可以插入 2 位或 4 位的年份值。2 位值会被映射到 1970 - 2069 范围内的年份。默认情况下,由连接器执行此转换。

errors.max.retries

默认值

-1

说明

指定当某个操作产生可重试错误(例如连接错误)后,连接器的响应方式。

可设置以下选项之一:

-1

无限制。无论之前失败过多少次,连接器始终会自动重新启动并重试该操作。

0

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

>0

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

event.converting.failure.handling.mode

默认值

warn

说明

指定当列的数据类型与 Debezium 内部架构指定的类型不匹配,导致无法转换表记录时,连接器的响应方式。可设置以下选项之一:

fail

抛出异常,报告由于字段的数据类型与架构类型不匹配而导致转换失败,并指出可能需要以 schema _only_recovery 模式重新启动连接器,以便成功完成转换。

warn

连接器将 null 值写入转换失败的列对应的事件字段,并向警告日志写入一条消息。

skip

连接器将 null 值写入转换失败的列对应的事件字段,并向调试日志写入一条消息。

event.processing.failure.handling.mode

默认值

fail

说明

指定连接器如何处理处理事件时发生的失败,例如遇到损坏的事件。可用设置如下:

fail

连接器抛出异常,报告有问题的事件及其位置,随后连接器停止。

warn

连接器不抛出异常,而是记录有问题的事件及其位置,然后跳过该事件。

ignore

连接器忽略有问题的事件,且不生成日志记录。

heartbeat.action.query

默认值

无默认值

描述

指定连接器在发送心跳消息时,在源数据库上执行的查询。

例如,以下查询会定期捕获源数据库中已执行的 GTID 集合的状态。

INSERT INTO gtid_history_table (select @gtid_executed)

heartbeat.interval.ms

默认值

0

描述

指定连接器向 Kafka 主题发送心跳消息的频率。默认情况下,连接器不发送心跳消息。

心跳消息可用于监控连接器是否正在接收数据库的变更事件。心跳消息有助于减少连接器重新启动时需要重新发送的变更事件数量。要发送心跳消息,请将此属性设置为一个正整数,该整数表示两条心跳消息之间的毫秒数。

incremental.snapshot.allow.schema.changes

默认值

false

描述

指定连接器是否允许在增量快照期间进行模式更改。当该值设置为 true 时,连接器会在增量快照期间检测模式更改,并重新选择当前块以避免锁定 DDL。

不支持对主键的更改。在增量快照期间更改主键可能导致结果不正确。另一个限制是,如果模式更改仅影响列的默认值,则在从 binlog 流中处理该 DDL 之前不会检测到此更改。这不会影响快照事件的值,但这些快照事件的模式可能包含过期的默认值。

incremental.snapshot.chunk.size

默认值

1024

描述

连接器在获取增量快照块时获取并读入内存的最大行数。增大数据块大小可以提高效率,因为快照执行的查询次数更少,但每次查询的规模更大。不过,较大的数据块大小也需要更多内存来缓冲快照数据。请根据你的环境调整数据块大小,以获得最佳性能。

incremental.snapshot.watermarking.strategy

默认值

insert_insert

描述

指定连接器在增量快照期间使用的水印机制,用于对可能先被增量快照捕获、随后在流式传输恢复后再次被捕获的事件进行去重。

你可以指定以下选项之一:

insert_insert(默认)

当你发送信号以启动增量快照时,对于 Debezium 在快照期间读取的每个块,它都会向信号数据集合写入一条记录,以记录打开快照窗口的信号。快照完成后,Debezium 会插入第二条记录,以记录关闭窗口的信号。

insert_delete

当你发送信号以启动增量快照时,对于 Debezium 读取的每个块,它都会向信号数据集合写入一条记录,以记录打开快照窗口的信号。快照完成后,该条记录会被删除。不会为关闭快照窗口的信号创建任何记录。设置此选项可防止信号数据集合快速增长。

max.batch.size

默认值

2048

描述

一个正整数值,用于指定此连接器在每次迭代中应处理的每批事件的最大数量。

max.queue.size

默认值

8192

描述

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

max.queue.size.in.bytes

默认值

0

描述

一个长整型值,用于指定阻塞队列的最大字节容量。默认情况下,不会为阻塞队列指定容量限制。要指定队列可消耗的字节数,请将此属性设置为正的长整型值。如果同时设置了 max.queue.size,则当队列的大小达到其中任一属性指定的限制时,向队列的写入将被阻塞。例如,如果你设置 max.queue.size=1000、max.queue.size.in.bytes=5000,那么当队列包含 1000 条记录后,或者当队列中记录的数据量达到 5000 字节后,向队列的写入就会被阻塞。

min.row.count.to.stream.results

默认值

1000

描述

在快照期间,连接器会查询其配置为捕获更改的每张表。连接器使用每个查询结果来生成读取事件,该事件包含该表中所有行的数据。此属性决定 MariaDB 连接器是将表的结果放入内存(速度快,但需要大量内存),还是流式传输结果(可能较慢,但适用于非常大的表)。此属性的设置指定了在连接器开始流式传输结果之前,表必须包含的最小行数。

要跳过所有表大小检查并在快照期间始终流式传输所有结果,请将此属性设置为 0。

notification.enabled.channels

默认值

无默认值

说明

为连接器启用的通知通道名称列表。默认情况下,提供以下通道:

  • sink
  • log
  • jmx

此外,你还可以实现自定义通知通道。

poll.interval.ms

默认值

500(0.5 秒)

说明

一个正整数值,指定连接器在开始处理一批事件之前,等待新更改事件出现的毫秒数。

provide.transaction.metadata

默认值

false

说明

决定连接器是否生成带有事务边界的事件,并用事务元数据丰富更改事件封装。如果要让连接器执行此操作,请指定 true。有关更多信息,请参阅事务元数据。

read.only

默认值

false

说明

指定连接器是否将水位线写入信号数据集合,以跟踪增量快照的进度。将值设置为 true,可使对数据库具有只读连接的连接器采用无需向信号数据集合写入的增量快照水位线策略。

signal.data.collection

默认值

无默认值

说明

用于向连接器发送信号的数据集合的完全限定名称。请使用以下格式指定集合名称:

<databaseName>.<tableName>

集合名称区分大小写。

signal.enabled.channels

默认值

无默认值

说明

为连接器启用的信号通道名称列表。默认情况下,提供以下通道:

  • source
  • kafka
  • file
  • jmx

你也可以选择实现一个自定义信号通道。

skipped.operations

默认值

t

说明

以逗号分隔的列表,用于指定连接器在流式传输期间要跳过的操作类型。

设置以下选项之一来指定要跳过的操作:

c

插入/创建操作。

u

更新操作。

d

删除操作。

t

清空(truncate)操作。

none

连接器不跳过任何操作。

snapshot.delay.ms

默认值

无默认值

说明

连接器启动后,在执行快照之前应等待的时间间隔(以毫秒为单位)。如果你在集群中启动多个连接器,此属性有助于避免快照中断,因为快照中断可能导致连接器重新平衡。

snapshot.fetch.size

默认值

未设置

说明

默认情况下,在快照期间,连接器按行批次读取表内容。设置此属性可指定单个批次中的最大行数。

snapshot.include.collection.list

默认值

table.include.list 中指定的所有表。

说明

一个可选的、以逗号分隔的正则表达式列表,用于匹配要包含在快照中的表的全限定名称(<databaseName>.<tableName>)。指定的项必须出现在连接器的 table.include.list 属性中。

此属性仅在连接器的 snapshot.mode 属性设置为 no_data 以外的值时才生效。此属性不会影响增量快照的行为。

为匹配表名,Debezium 会将你指定的正则表达式作为锚定正则表达式来应用。也就是说,指定的表达式会与表的整个名称字符串进行匹配,而不会匹配表名中可能出现的子字符串。

snapshot.lock.timeout.ms

默认值

10000

说明

正整数,用于指定执行快照时等待获取表锁的最长时间(以毫秒为单位)。如果连接器在此时间间隔内无法获取表锁,快照将失败。

有关更多信息,请参阅描述 MariaDB 连接器如何执行数据库快照 的文档。

snapshot.locking.mode

默认值

minimal

说明

指定连接器是否持有全局 MariaDB 读锁以及持有多长时间,该锁会在连接器执行快照期间阻止对数据库的任何更新。

可使用以下设置:

minimal

连接器仅在快照的初始阶段持有全局读锁,该阶段连接器读取数据库模式和其他元数据。在快照的下一阶段,连接器会在从每张表中选择所有行时释放该锁。为了以一致的方式执行 SELECT 操作,连接器使用 REPEATABLE READ 事务。虽然释放全局读锁允许其他 MariaDB 客户端更新数据库,但 REPEATABLE READ 隔离级别的使用可确保快照的一致性,因为连接器在整个事务期间持续读取相同的数据。

extended

在整个快照期间阻塞所有写操作。如果客户端提交的并发操作与 MariaDB 中的 REPEATABLE READ 隔离级别不兼容,请使用此设置。

none

阻止连接器在快照期间获取任何表锁。虽然此选项适用于所有快照模式,但仅在快照运行期间不会发生模式变更时使用它才是安全的。使用 MyISAM 引擎定义的表始终会获取表锁,因此即使设置了此选项,这类表也会被锁定。这种行为与使用 InnoDB 引擎定义的表不同,后者获取的是行级锁。

custom

连接器根据 snapshot.locking.mode.custom.name 属性指定的实现执行快照,该属性是 io.debezium.spi.snapshot.SnapshotLock 接口的自定义实现。

snapshot.locking.mode.custom.name

默认值

无默认值

描述

当 snapshot.locking.mode 设置为 custom 时,使用此设置指定由 'io.debezium.spi.snapshot.SnapshotLock' 接口定义的 name() 方法所提供的自定义实现的名称。

有关更多信息,请参阅自定义快照器 SPI。

snapshot.max.threads

默认值

1

描述

指定连接器执行初始快照时使用的线程数。该值必须是大于或等于 1 的正整数。当设置为大于 1 的值时,连接器会基于主键范围将每张表划分为多个块,并在所有可用线程之间并发处理这些块,从而执行并行快照。

基于分块的并行快照是一项孵化中特性。
没有主键的表,以及使用快照选择覆盖的表,会被作为一个块进行处理,并回退到单线程快照。

snapshot.max.threads.multiplier

默认值

1

描述

一个全局乘数,用于控制并行快照期间每张表创建的块数量。默认情况下,Debezium 为每个线程创建一个块。设置更高的值会使连接器在执行快照时为每张表创建更多、更小的块。较小的块能让线程负载更加均衡。例如,在 4 个线程、乘数为 2 的情况下,Debezium 会创建 8 个块,而不是 4 个。

要为特定表覆盖该乘数,请将 snapshot.max.threads.multiplier.<fully_qualified_table_name> 设置为所需的值。

与提高乘数相比,增加线程数通常是提高快照吞吐量更有效的方式。
基于块的并行快照是一项孵化中(incubating)的功能。

legacy.snapshot.max.threads

默认值

false

描述

默认值(false)使连接器能够通过将表划分成块,并使用独立线程并发处理每个块,从而加速初始快照。

当启用旧版行为时,完成自身表快照的线程会保持空闲,等待其他线程完成。在强制执行连接超时的环境中,空闲连接可能会导致连接器在快照完成后无法干净地关闭连接。这样一来,即使快照已成功捕获所有数据,也可能会抛出异常。如果您遇到此问题,请将 snapshot.max.threads 设置为 1,然后重试快照。
传统的「每线程一张表」行为已弃用,并将在未来版本中移除。属性 internal.legacy.snapshot.max.threads 是 legacy.snapshot.max.threads 的已弃用别名,不应在新配置中使用。

snapshot.mode

默认值

initial

说明

指定连接器启动时执行快照的判定标准。

可用设置如下:

always

连接器每次启动都会执行快照。快照包含被捕获表的结构和数据。指定此值后,每次连接器启动时,都会用被捕获表数据的完整表示来填充主题。

initial

仅当该逻辑服务器名称尚无已记录的偏移量,或连接器检测到之前的快照未能完成时,连接器才执行快照。快照完成后,连接器开始流式传输后续数据库更改的事件记录。

initial_only

仅当该逻辑服务器名称尚无已记录的偏移量时,连接器才执行快照。快照完成后,连接器停止运行。它不会转入流式传输阶段去读取 binlog 中的更改事件。

schema_only

已弃用,请参阅 no_data。

no_data

连接器执行仅捕获结构而不捕获任何表数据的快照。如果不需要主题包含数据的一致性快照,但希望捕获上次连接器重启之后发生的所有结构更改,请设置此选项。

schema_only_recovery

已弃用,请参阅 recovery。

recovery

设置此选项可恢复丢失或损坏的数据库结构历史主题。重启后,连接器会执行快照,从源表重建该主题。您也可以设置该属性,以定期清理意外膨胀的数据库结构历史主题。

如果在上次连接器关闭之后有结构更改被提交到数据库中,请勿使用此模式执行快照。

when_needed

连接器启动后,仅在检测到以下情况之一时才执行快照:

  • 无法检测到任何主题偏移量。
  • 之前记录的偏移量指定的 binlog 位置或 GTID 在服务器上不可用。

configuration_based

使用此选项,您可以通过一组以 snapshot.mode.configuration.based 为前缀的连接器属性来控制快照行为。

custom

连接器按照 snapshot.mode.custom.name 属性所指定的实现执行快照,该属性定义了 io.debezium.spi.snapshot.Snapshotter 接口的自定义实现。

snapshot.mode.configuration.based.snapshot.data

默认值

false

说明

如果 snapshot.mode 设置为 configuration_based,请设置此属性以指定连接器在执行快照时是否包含表数据。

snapshot.mode.configuration.based.snapshot.on.data.error

默认值

false

说明

如果 snapshot.mode 设置为 configuration_based,请设置此属性以指定在事务日志中不再提供数据时,连接器是否在快照中包含表数据。

snapshot.mode.configuration.based.snapshot.on.schema.error

默认值

false

说明

如果 snapshot.mode 设置为 configuration_based,请设置此属性以指定在架构历史主题不可用时,连接器是否在快照中包含表架构。

snapshot.mode.configuration.based.snapshot.schema

默认值

false

说明

如果 snapshot.mode 设置为 configuration_based,请设置此属性以指定连接器在执行快照时是否包含表架构。

snapshot.mode.configuration.based.start.stream

默认值

false

说明

如果 snapshot.mode 设置为 configuration_based,请设置此属性以指定连接器是否在快照完成后开始流式传输变更事件。

snapshot.mode.custom.name

默认值

无默认值

说明

如果 snapshot.mode 设置为 custom,请使用此设置指定自定义实现的名称,该名称由 'io.debezium.spi.snapshot.Snapshotter' 接口中定义的 name() 方法提供。连接器重启后,Debezium 会调用指定的自定义实现,以确定是否执行快照。有关更多信息,请参阅自定义快照器 SPI。

snapshot.query.mode

默认值

select_all

说明

指定连接器在执行快照时如何查询数据。

请设置以下选项之一:

select_all(默认)

连接器使用 select all 查询从被捕获的表中检索行,并可根据列的 include 和 exclude 列表配置,有选择地调整所选列。

custom

连接器根据 snapshot.query.mode.custom.name 属性指定的实现执行快照查询,该属性定义了 io.debezium.spi.snapshot.SnapshotQuery 接口的自定义实现。

与使用 snapshot.select.statement.overrides 属性相比,此设置使您能够更灵活地管理快照内容。

snapshot.query.mode.custom.name

默认值

无默认值

描述

当 snapshot.query.mode 设置为 custom 时,使用此设置指定由 'io.debezium.spi.snapshot.SnapshotQuery' 接口定义的 name() 方法所提供的自定义实现的名称。有关更多信息,请参阅自定义快照器 SPI。

snapshot.select.statement.overrides

默认值

无默认值

描述

指定连接器使用自定义 SELECT 语句来确定哪些行应包含在快照中的表。

此属性仅影响快照。它不适用于连接器在流式传输阶段从日志中读取的事件。

此属性由两个协同工作的部分组成:

主属性

一个以逗号分隔的完全限定表名列表,格式为 <databaseName>.<tableName>。该列表标识了您想要指定自定义快照查询的表。

例如:

"snapshot.select.statement.overrides": "inventory.products,customers.orders"

如果完全限定的模式名或表名包含特殊字符,例如空格、方括号([ 或 ])或句点字符(.),请将该字符串用双引号括起来,以防止连接器将这些特殊字符解释为分隔符。如果完全限定的表名不包含空格或特殊字符,则表名周围的双引号是可选的。

次要属性

对于主属性中列出的每个表,你都必须定义对应的 snapshot.select.statement.overrides.<databaseName>.<tableName> 属性,用于指定快照期间要运行的自定义 SELECT 语句。

SELECT 语句决定表中哪些行会被包含在快照中。

例如,要为 customers.orders 表指定 SELECT 语句,请添加以下属性:

snapshot.select.statement.overrides.customers.orders

如果某个表在主属性中列出,但缺少对应的次要属性,连接器会记录一条警告日志,并对该表使用默认的快照行为。

配置示例

以下示例如何配置 snapshot.select.statement.overrides 属性,以便对 customers.orders 表执行快照,并且仅包含未被软删除的记录;也就是说,软删除字段 delete_flag 的值为 0。

"snapshot.select.statement.overrides": "customer.orders",
"snapshot.select.statement.overrides.customer.orders": "SELECT * FROM customers.orders WHERE delete_flag = 0 ORDER BY id DESC"

snapshot.tables.order.by.row.count

默认值

disabled

说明

指定连接器在执行初始快照时处理表的顺序。

设置以下选项之一:

descending

连接器按表中行数从多到少的顺序执行快照。

ascending

连接器按表中行数从少到多的顺序执行快照。

disabled

连接器在执行初始快照时忽略行数。

source.struct.version

默认值

v2

说明

Debezium 事件中 source 块的架构版本。Debezium 0.10 对 source 块的结构进行了一些破坏性变更,以统一所有连接器暴露的结构。

将此选项设置为 v1 可生成早期版本使用的结构。但是,不建议使用此设置,并计划在未来版本的 Debezium 中移除该设置。

statistics.metrics.enabled

默认值

true

说明

指定连接器是否为流式传输指标收集高级统计数据,例如分位数。

设置为 true 时,连接器收集以下统计数据:

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

统计数据使用概率数据结构(DDSketch)计算,该结构可提供相对精度为 1% 的近似分位数值。

+ 设置为 false 时,连接器不收集分位数,分位数 JMX 指标将返回 null 值。连接器继续收集最小值、最大值和平均值。禁用分位数收集可略微减少内存开销。

+ 有关更多信息,请参阅流式传输指标。

streaming.delay.ms

默认值

0

说明

指定连接器在完成快照后延迟启动流式传输过程的时间(以毫秒为单位)。设置延迟间隔有助于防止在快照完成之后、流式传输过程开始之前立即发生故障时,连接器重新启动快照。设置的延迟值应高于为 Kafka Connect 工作程序设置的 offset.flush.interval.ms 属性的值。

table.ignore.builtin

默认值

true

说明

一个布尔值,指定是否应忽略内置系统表。无论表包含列表和排除列表如何设置,该设置都会生效。默认情况下,对系统表中值的更改不会被捕获,Debezium 也不会为系统表的更改生成事件。

topic.cache.size

默认值

10000

说明

指定可在有界并发哈希映射中存储的 Topic 名称数量。连接器使用该缓存来帮助确定与某个数据集合相对应的 Topic 名称。

topic.delimiter

默认值

.(句点)

说明

指定连接器在 Topic 名称的各个组成部分之间插入的分隔符。

topic.heartbeat.prefix

默认值

__debezium-heartbeat

说明

指定连接器向其发送心跳消息的 Topic 名称。该 Topic 名称遵循以下格式:

topic.heartbeat.prefix.topic.prefix

例如,当该属性设置为默认值,且 Topic 前缀为 fulfillment 时,Topic 名称为 __debezium-heartbeat.fulfillment。

如果设置了 topic.heartbeat.name,则会忽略此属性。

topic.heartbeat.name

默认值

空

说明

指定连接器向其发送心跳消息的 Topic 的显式完整名称,该名称会覆盖由 topic.heartbeat.prefix 和 topic.prefix 推导出的基于前缀的命名方式。

设置后,无论连接器的 topic.prefix 是什么,所有心跳消息都会被路由到这个确切的 Topic 名称。当运行多个连接器并将心跳事件统一汇总到一个共享 Topic 中时,此属性非常有用,可以避免创建大量单分区的心跳 Topic。

例如,将此属性设置为 debezium-heartbeat 后,所有心跳消息都会被路由到名为 debezium-heartbeat 的 Topic。

如果此属性为空或未设置,连接器将回退到默认行为:topic.heartbeat.prefix.topic.prefix。

topic.naming.strategy

默认值

io.debezium.schema.DefaultTopicNamingStrategy

说明

连接器所使用的 TopicNamingStrategy 类的名称。指定的策略决定了连接器如何为存储数据更改事件记录、架构更改、事务、心跳等内容的 Topic 命名。

topic.transaction

默认值

transaction

说明

指定连接器向其发送事务元数据消息的 Topic 名称。该 Topic 名称遵循以下模式:

topic.prefix.topic.transaction

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

use.nongraceful.disconnect

默认值

false

说明

一个布尔值,用于指定二进制日志客户端的保活线程是否将 SO_LINGER 套接字选项设置为 0,以立即关闭僵死的 TCP 连接。如果连接器在 SSLSocketImpl.close 中遇到死锁,请将该值设置为 true。有关更多信息,请参阅 mysql-binlog-connector-java GitHub 仓库中的 Issue 133。

guardrail.collections.max

默认值

0

说明

指定连接器可以捕获的最大表数。此限制不仅适用于正在捕获数据变更的表,也适用于架构历史中记录其架构的任何表。超过此限制将触发 guardrail.collections.limit.action 指定的操作。将此属性设置为 0 可防止连接器触发防护操作。

guardrail.collections.limit.action

默认值

warn

说明

指定当连接器捕获的表数量超过 guardrail.collections.max 属性中指定的数量时要触发的操作。将该属性设置为以下值之一:

fail

连接器失败并报告异常。

warn

连接器记录一条警告。

Debezium 连接器数据库架构历史配置属性

Debezium 提供了一组 schema.history.internal.* 属性,用于控制连接器与架构历史主题的交互方式。

下表描述了用于配置 Debezium 连接器的 schema.history.internal 属性。

表 17. 连接器数据库架构历史配置属性

属性 默认值 说明

schema.history.internal.kafka.topic

无默认值

连接器存储数据库架构历史的 Kafka 主题的完整名称。

schema.history.internal.kafka.bootstrap.servers

无默认值

连接器用于与 Kafka 集群建立初始连接的主机/端口对列表。此连接用于检索连接器先前存储的数据库架构历史,以及写入从源数据库读取的每条 DDL 语句。每个主机/端口对应指向 Kafka Connect 进程所使用的同一个 Kafka 集群。

schema.history.internal.kafka.recovery.poll.interval.ms

100

一个整数值,用于指定连接器在启动/恢复期间轮询已持久化数据时应等待的最大毫秒数。默认值为 100 毫秒。

schema.history.internal.kafka.query.timeout.ms

3000

一个整数值,用于指定连接器在使用 Kafka 管理客户端获取集群信息时应等待的最大毫秒数。

schema.history.internal.kafka.create.timeout.ms

30000

一个整数值,用于指定连接器在使用 Kafka 管理客户端创建 Kafka 历史主题时应等待的最大毫秒数。

schema.history.internal.kafka.recovery.attempts

100

连接器在恢复失败并抛出错误之前,尝试读取已持久化历史数据的最大次数。在未接收到数据后等待的最长时间为 recovery.attempts × recovery.poll.interval.ms。

schema.history.internal.skip.unparseable.ddl

false

一个布尔值,用于指定连接器是忽略格式错误或无法识别的数据库语句,还是停止处理以便由人工修复该问题。安全的默认值为 false。只有在谨慎的情况下才应使用跳过选项,因为在处理 binlog 时,该选项可能导致数据丢失或损坏。

schema.history.internal.store.only.captured.tables.ddl

false

一个布尔值,用于指定连接器是记录某个模式或数据库中所有表的架构结构,还是只记录被指定为捕获目标的表的架构结构。请指定以下值之一:

false(默认值)

在进行数据库快照时,连接器会记录数据库中所有非系统表的架构数据,包括未被指定为捕获目标的表。最好保留默认设置。如果日后决定捕获最初未被指定为捕获目标的表的变更,连接器可以轻松开始从这些表中捕获数据,因为它们的架构结构已经存储在架构历史主题中。Debezium 需要表的架构历史,以便识别变更事件发生时该表所呈现的结构。

true

在数据库快照期间,连接器只为 Debezium 捕获变更事件的表记录表结构。如果你更改了该默认值,而稍后又配置连接器从该数据库中的其他表捕获数据,连接器将缺少从这些表捕获变更事件所需的结构信息。

schema.history.internal.store.only.captured.databases.ddl

false

一个布尔值,用于指定连接器是否记录该数据库实例中所有逻辑数据库的结构信息。可指定以下值之一:

true

连接器仅为 Debezium 捕获变更事件所在的逻辑数据库和 schema 中的表记录结构信息。

false

连接器记录所有逻辑数据库的结构信息。

schema.history.internal.memory.optimization

off

控制 Debezium 如何使用驻留器(interner)在内存中对相同的结构对象(表、列、属性)进行去重。可指定以下值之一:

off

不执行任何去重(默认值)。

on

每个连接器使用其各自独立的驻留池。可在不影响其他连接器的情况下,降低单个连接器的堆内存使用量。

shared

所有配置为 shared 的连接器共享同一个全局驻留池。当有多个连接器跟踪结构相似的表时,可最大限度地提高去重效果。

MariaDB 连接器透传配置属性

你可以在连接器配置中设置透传属性,以自定义 Apache Kafka 生产者和消费者的行为。有关 Kafka 生产者和消费者全部配置属性的信息,请参阅 Kafka 文档。

用于配置生产者和消费者客户端与 schema 历史主题交互方式的透传属性

Debezium 依赖 Apache Kafka 生产者将结构变更写入数据库结构历史主题。同样,它在连接器启动时依赖 Kafka 消费者从数据库结构历史主题中读取数据。你可以通过为一组以 schema.history.internal.producer.* 和 schema.history.internal.consumer.* 为前缀的透传配置属性赋值,来定义 Kafka 生产者和消费者客户端的配置。这些透传的生产者和消费者数据库结构历史属性可控制一系列行为,例如这些客户端如何与 Kafka 代理建立安全连接,如下例所示:

schema.history.internal.producer.security.protocol=SSL
schema.history.internal.producer.ssl.keystore.location=/var/private/ssl/kafka.server.keystore.jks
schema.history.internal.producer.ssl.keystore.password=test1234
schema.history.internal.producer.ssl.truststore.location=/var/private/ssl/kafka.server.truststore.jks
schema.history.internal.producer.ssl.truststore.password=test1234
schema.history.internal.producer.ssl.key.password=test1234

schema.history.internal.consumer.security.protocol=SSL
schema.history.internal.consumer.ssl.keystore.location=/var/private/ssl/kafka.server.keystore.jks
schema.history.internal.consumer.ssl.keystore.password=test1234
schema.history.internal.consumer.ssl.truststore.location=/var/private/ssl/kafka.server.truststore.jks
schema.history.internal.consumer.ssl.truststore.password=test1234
schema.history.internal.consumer.ssl.key.password=test1234

Debezium 在将属性传递给 Kafka 客户端之前,会去除属性名称中的前缀。

有关 Kafka 生产者配置属性 和 Kafka 消费者配置属性 的更多信息,请参阅 Apache Kafka 文档。

用于配置 MariaDB 连接器与 Kafka 信号主题交互方式的透传属性

Debezium 提供了一组 signal.* 属性,用于控制连接器与 Kafka 信号主题的交互方式。

下表描述了 Kafka signal 属性。

表 18. Kafka 信号配置属性

属性 默认值 说明

signal.kafka.topic

<topic.prefix>-signal

连接器用于监控临时信号的 Kafka 主题名称。

如果自动创建主题已禁用,你必须手动创建所需的信号主题。为保证信号的顺序,必须使用信号主题。信号主题必须只有一个分区。

signal.kafka.groupId

kafka-signal

Kafka 消费者所使用的组 ID 名称。

signal.kafka.bootstrap.servers

无默认值

连接器用于建立与 Kafka 集群初始连接的主机和端口对列表。每个端口对引用 Debezium Kafka Connect 进程所使用的 Kafka 集群。

signal.kafka.poll.timeout.ms

100

一个整数值,指定连接器在轮询信号时等待的最长毫秒数。

用于配置信号通道 Kafka 消费者客户端的透传属性

Debezium 连接器提供了对信号 Kafka 消费者的透传配置功能。透传信号属性以 signal.consumer.* 为前缀。例如,连接器会将 signal.consumer.security.protocol=SSL 这样的属性传递给 Kafka 消费者。

Debezium 在将属性传递给 Kafka 信号消费者之前,会去除这些属性的前缀。

用于配置 MariaDB 连接器接收器通知通道的透传属性

下表描述了可用于配置 Debezium sink notification 通道的属性。

属性默认值描述
notification.sink.topic.name无默认值接收来自 Debezium 通知的主题名称。当你将 notification.enabled.channels 属性配置为包含 sink 作为启用的通知通道之一时,此属性为必填项。

表 19. Sink 通知配置属性

Debezium 连接器数据库驱动透传配置属性

Debezium 连接器支持数据库驱动的透传配置。透传数据库属性以 driver.* 为前缀。例如,连接器会将 driver.foobar=false 这类属性传递给 JDBC URL。

Debezium 在将属性传递给数据库驱动之前,会去除属性中的前缀。

监控

除了 Kafka 和 Kafka Connect 提供的内置 JMX 指标支持之外,Debezium MariaDB 连接器还提供三种类型的指标。

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

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

默认情况下,MariaDB 连接器对流式指标使用以下 MBean 名称:

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

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

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

快照指标

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

下表列出了可用于监控 Debezium 快照操作的 JMX 指标,包括行计数、表进度、持续时间和队列容量。除非快照操作正在进行中,或者自上次连接器启动以来已发生过快照,否则不会暴露快照指标。

AttributesTypeDescription
LastEventstring连接器读取到的最后一个快照事件。
MilliSecondsSinceLastEventlong自连接器读取并处理最近一个事件以来经过的毫秒数。
NumberOfErroneousEventslong记录连接器在快照操作期间识别为错误的变更事件数量。每当在初始快照、增量快照或临时快照期间遇到无法处理的事件时,连接器就会递增此指标。事件处理失败的原因可能是格式不正确、与模式不兼容,或在转换过程中发生故障。该指标值在连接器任务的整个生命周期内保持。如果快照被中断且连接器任务重新启动,则该指标计数会重置为 0。
TotalNumberOfEventsSeenlong自上次启动或重置以来,此连接器看到的事件总数。
NumberOfEventsFilteredlong被连接器上配置的包含/排除列表过滤规则过滤掉的事件数量。
CapturedTablesstring[]连接器捕获的表列表。
QueueTotalCapacityint用于在快照器和主 Kafka Connect 循环之间传递事件的队列长度。
QueueRemainingCapacityint用于在快照器和主 Kafka Connect 循环之间传递事件的队列的剩余容量。
TotalTableCountint纳入快照的表的总数。
RemainingTableCountint快照尚未复制的表的数量。
SnapshotRunningboolean快照是否已启动。
SnapshotPausedboolean快照是否已暂停。
SnapshotAbortedboolean快照是否已中止。
SnapshotCompletedboolean快照是否已完成。
SnapshotSkippedboolean快照是否被跳过。
SnapshotDurationInSecondslong快照迄今为止花费的总秒数,即使尚未完成。其中也包含快照处于暂停状态的时间。
SnapshotPausedDurationInSecondslong快照处于暂停状态的总秒数。如果快照被多次暂停,暂停时间会累加。
RowsScannedMap<String, Long>一个映射,包含快照中每张表已扫描的行数。处理过程中会逐步将表添加到该映射中。每扫描 10,000 行以及每完成一张表时更新一次。
TableChunkCountsMap<String, Long>一个映射,在使用基于分块的多线程快照时,包含快照中每张表的分块数量。
TableChunksCompletedCountsMap<String, Long>一个映射,在使用基于分块的多线程快照时,包含快照中每张表已完成的分块数量。
MaxQueueSizeInByteslong队列的最大缓冲字节数。当 max.queue.size.in.bytes 被设置为正的 long 值时,此指标可用。
CurrentQueueSizeInByteslong队列中记录的当前字节数。

下表列出了连接器执行增量快照时可用的其他 JMX 指标,包括可用于跟踪快照进度的块和表边界标识符。

属性类型描述
ChunkIdstring当前快照块的标识符。
ChunkFromstring定义当前块的主键集的下界。
ChunkTostring定义当前块的主键集的上界。
TableFromstring当前正在快照的表的主键集的下界。
TableTostring当前正在快照的表的主键集的上界。

Debezium MariaDB 连接器还提供了 HoldingGlobalLock 自定义快照指标。该指标的值为布尔值,表示连接器当前是否持有全局写锁或表写锁。

流处理指标

Debezium MariaDB 连接器提供了三种类型的指标,这些指标是对 Kafka 和 Kafka Connect 内置的 JMX 指标支持的补充。

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

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

下表列出了可用于监控 Debezium 流式处理操作的 JMX 指标,包括按类型统计的事件数量、相对源端的延迟、队列容量以及连接状态。

AttributesTypeDescription
LastEventstring连接器读取的最后一个流式事件。
MilliSecondsSinceLastEventlong自连接器读取并处理最近一个事件以来经过的毫秒数。
NumberOfErroneousEventslong记录连接器在流式传输过程中识别为错误的变更事件数量。在流式会话的生命周期内,每当连接器遇到无法处理的事件时,该指标就会递增。事件处理失败的原因可能包括格式错误、与模式不兼容,或在转换过程中发生失败。该指标值在连接器任务的生命周期内持续保留。连接器重启后,该指标计数将重置为 0。
TotalNumberOfEventsSeenlong自上次连接器启动或指标重置以来,源数据库报告的数据变更事件总数。表示 Debezium 需要处理的数据变更工作负载。
TotalNumberOfCreateEventsSeenlong自连接器上次启动或指标重置以来,连接器处理的创建事件总数。
TotalNumberOfUpdateEventsSeenlong自连接器上次启动或指标重置以来,连接器处理的更新事件总数。
TotalNumberOfDeleteEventsSeenlong自连接器上次启动或指标重置以来,连接器处理的删除事件总数。
NumberOfEventsFilteredlong被连接器上配置的包含/排除列表过滤规则过滤掉的事件数量。
NumberOfUnchangedEventsSkippedlong自上次连接器启动或指标重置以来,由于没有被监控列发生变化而跳过的更新事件数量。如果 skip.messages.without.change 为 false,则默认值为 -1。
CapturedTablesstring[]连接器捕获的表列表。
QueueTotalCapacityint用于在流式处理器和主 Kafka Connect 循环之间传递事件的队列长度。
QueueRemainingCapacityint用于在流式处理器和主 Kafka Connect 循环之间传递事件的队列剩余可用容量。
Connectedboolean表示连接器当前是否已连接到数据库服务器的标志。
MilliSecondsBehindSourcelong最后一个变更事件的时间戳与连接器处理该事件之间相差的毫秒数。该值会包含数据库服务器和连接器所在机器之间时钟的任何差异。
MilliSecondsBehindSourceMinValuelong连接器运行期间观测到的相对源的最小延迟毫秒数。
MilliSecondsBehindSourceMaxValuelong连接器运行期间观测到的相对源的最大延迟毫秒数。
MilliSecondsBehindSourceAverageValuedouble连接器运行期间所有观测值计算出的相对源的平均延迟毫秒数。
MilliSecondsBehindSourceP50double相对源延迟毫秒数的第 50 百分位(中位数)。与平均值相比,该指标受离群值影响较小,因此能更稳健地衡量典型延迟。当 statistics.metrics.enabled 设置为 true(默认值)时可用。
MilliSecondsBehindSourceP95double相对源延迟毫秒数的第 95 百分位。该指标表明 95% 的延迟测量值低于此值,有助于识别尾部延迟并设定 SLA 阈值。当 statistics.metrics.enabled 设置为 true(默认值)时可用。
MilliSecondsBehindSourceP99double相对源延迟毫秒数的第 99 百分位。该指标表明 99% 的延迟测量值低于此值,有助于了解最坏情况下的性能表现。当 statistics.metrics.enabled 设置为 true(默认值)时可用。
NumberOfCommittedTransactionslong已提交的已处理事务数量。

| SourceEventPosition | Map<String, String> | 最后接收到的事件的坐标

Debezium MariaDB 连接器还提供以下额外的流式传输指标:

属性类型描述
BinlogFilenamestring连接器最近读取的 binlog 文件名。
BinlogPositionlong连接器最近读取的 binlog 位置(以字节计)。
IsGtidModeEnabledboolean标志位,表示连接器当前是否正在跟踪 MariaDB 服务器的 GTID。
GtidSetstring连接器在读取 binlog 时处理的最近一个 GTID 集合的字符串表示。
NumberOfSkippedEventslongMariaDB 连接器跳过的事件数量。通常,事件被跳过是因为 MariaDB 的 binlog 中存在格式错误或无法解析的事件。
NumberOfDisconnectslongMariaDB 连接器断开连接的次数。
NumberOfRolledBackTransactionslong已处理但被回滚且未被流式传输的事务数量。
NumberOfNotWellFormedTransactionslong不符合 BEGIN + COMMIT/ROLLBACK 预期协议的事务数量。在正常情况下,该值应为 0。
NumberOfLargeTransactionslong未能放入预读缓冲区的事务数量。为获得最佳性能,该值应明显小于 NumberOfCommittedTransactions 和 NumberOfRolledBackTransactions。

表 20. 其他 MariaDB 流式传输指标说明

Schema 历史指标

MBean 为 debezium.mariadb:type=connector-metrics,context=schema-history,server=<topic.prefix>。

下表列出了可用于监控连接器 schema 历史处理过程的 JMX 指标,包括恢复状态、已应用的 schema 变更数量,以及最近变更的时间戳。

属性类型描述
Statusstring数据库架构历史记录的状态,取值为 STOPPED、RECOVERING(正在从存储中恢复历史记录)或 RUNNING。
RecoveryStartTimelong恢复开始的时间,以纪元秒数表示。
ChangesRecoveredlong在恢复阶段读取到的变更数量。
ChangesAppliedlong在恢复和运行期间已应用的架构变更总数。
MilliSecondsSinceLast​RecoveredChangelong自上次从历史存储中恢复变更以来经过的毫秒数。
MilliSecondsSinceLast​AppliedChangelong自上次应用变更以来经过的毫秒数。
LastRecoveredChangestring从历史存储中恢复的最近一次变更的字符串表示。
LastAppliedChangestring最近一次已应用变更的字符串表示。

出错时的行为

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

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

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

配置和启动错误

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

  • 连接器的配置无效。
  • 连接器无法使用指定的连接参数成功连接到 MariaDB 服务器。
  • 连接器尝试从 binlog 中的某个位置重新启动,而 MariaDB 已不再保留该位置的历史记录。

在这些情况下,错误消息中会包含问题的详细信息,可能还会给出建议的解决办法。修正配置或解决 MariaDB 的问题后,请重新启动连接器。

MariaDB 变得不可用

如果你的 MariaDB 服务器变得不可用,Debezium MariaDB 连接器将因错误而失败,连接器随之停止。当服务器再次可用时,请重新启动连接器。

不过,如果你连接的是高可用 MariaDB 集群,则可以立即重新启动连接器。它会连接到集群中的另一台 MariaDB 服务器,找到该服务器 binlog 中表示最后一个事务的位置,并从该特定位置开始读取新服务器的 binlog。

Kafka Connect 正常停止

当 Kafka Connect 正常停止时,会有一个短暂的延迟,期间 Debezium MariaDB 连接器任务会被停止并在新的 Kafka Connect 进程上重新启动。

Kafka Connect 进程崩溃

如果 Kafka Connect 崩溃,进程会停止,且任何 Debezium MariaDB 连接器任务都会终止,而不会记录其最近处理的偏移量。在分布式模式下,Kafka Connect 会在其他进程上重新启动连接器任务。但是,MariaDB 连接器会从先前进程记录的最后一个偏移量处恢复。因此,替代任务可能会重新生成崩溃之前已处理过的某些事件,从而产生重复事件。

每条更改事件消息都包含源特定的信息,你可以用它们来识别重复事件,例如:

  • 事件来源
  • MariaDB 服务器的事件时间
  • binlog 文件名和位置 +1
  • GTID。

Kafka 变得不可用

Kafka Connect 框架通过使用 Kafka 生产者 API 将 Debezium 变更事件记录到 Kafka 中。如果 Kafka broker 变得不可用,Debezium MariaDB 连接器会暂停,直到连接重新建立,然后连接器会从中断处继续工作。

MariaDB 清除 binlog 文件

如果 Debezium MariaDB 连接器停止的时间过长,MariaDB 服务器会清除较旧的 binlog 文件,连接器的最后位置可能会丢失。当连接器重新启动时,MariaDB 服务器已不再保留起始点,连接器将再次执行初始快照。如果快照被禁用,连接器会因错误而失败。

有关 MariaDB 连接器如何执行初始快照的详细信息,请参阅快照。

评论

登录后参与评论

正在加载评论…