MySQL
Debezium MySQL 连接器
MySQL 拥有一个二进制日志(binlog),它按照操作提交到数据库的先后顺序记录下所有操作。这既包括表结构的变更,也包括表中数据的变更。MySQL 使用 binlog 进行复制和恢复。
Debezium MySQL 连接器读取 binlog,针对行级的 INSERT、UPDATE 和 DELETE 操作生成变更事件,并将这些变更事件发送到 Kafka 主题中。客户端应用读取这些 Kafka 主题。
由于 MySQL 通常配置为在经过指定时间段后清除 binlog,因此 MySQL 连接器会先对你的每个数据库执行一次初始的一致性快照。MySQL 连接器从生成快照的位置开始读取 binlog。
有关与此连接器兼容的 MySQL 数据库版本信息,请参阅 Debezium 发行版概览。
连接器的工作原理
了解连接器所支持的 MySQL 拓扑结构,有助于你规划应用。要以最佳方式配置和运行 Debezium MySQL 连接器,需要理解连接器如何跟踪表结构、暴露 schema 变更、执行快照以及确定 Kafka 主题名称。
支持的 MySQL 拓扑
Debezium MySQL 连接器支持以下 MySQL 拓扑:
单机模式
当使用单个 MySQL 服务器时,该服务器必须启用 binlog,Debezium MySQL 连接器才能监控该服务器。
这通常是可以接受的,因为二进制日志也可以用作增量备份。
在这种情况下,MySQL 连接器始终连接并跟随这个单机 MySQL 服务器实例。
主库与副本
Debezium MySQL 连接器可以跟随其中一个主服务器,或者跟随其中一个副本(前提是该副本已启用 binlog),但连接器只能检测到该服务器可见的集群中的变更。通常这不成问题,除了多主拓扑之外。
连接器记录它在服务器 binlog 中的位置,而集群中每个服务器的 binlog 位置各不相同。因此,连接器必须只跟随一个 MySQL 服务器实例。如果该服务器发生故障,必须先重启或恢复该服务器,连接器才能继续运行。
高可用集群
你可以使用各种提供冗余的方案来保障 MySQL 的可用性。部署冗余的副本主机,可以显著提升容错能力,并几乎立即从问题和故障中恢复。由于大多数高可用 MySQL 集群使用 GTID,副本能够追踪在任意主服务器上发生的所有变更。
多主模式
网络数据库(NDB)集群复制 使用一个或多个 MySQL 副本节点,每个副本节点从多个主服务器复制。集群复制提供了一种强大的方式来聚合多个 MySQL 集群的复制。这种拓扑结构需要使用 GTID。
Debezium MySQL 连接器可以将这些多主 MySQL 副本用作源,并且可以故障转移到不同的多主 MySQL 副本,前提是新副本已追赶上旧副本,即新副本拥有在第一个副本上见过的所有事务。即使连接器只使用部分数据库和/或表,这种方式同样有效,因为连接器可以配置为在尝试重新连接到新的多主 MySQL 副本并在 binlog 中找到正确位置时,包含或排除特定的 GTID 来源。
托管式
Debezium MySQL 连接器可以使用 Amazon RDS 和 Amazon Aurora 等托管式数据库选项。
由于这些托管式选项不允许使用全局读锁,因此连接器在创建一致快照时会使用表级锁。
Schema 历史主题
当数据库客户端查询数据库时,客户端使用的是数据库当前的 schema。但是,数据库 schema 可能随时发生变化,这意味着连接器必须能够确定在记录每次插入、更新或删除操作时的 schema 是什么。此外,连接器不一定能将当前 schema 应用于每一个事件。如果某个事件相对较旧,它可能是在当前 schema 应用之前记录的。
为确保正确处理 schema 变更之后发生的事件,MySQL 在事务日志中不仅包含影响数据的行级变更,还包含应用于数据库的 DDL 语句。当连接器在 binlog 中遇到这些 DDL 语句时,会对其进行解析,并更新每张表 schema 的内存表示。连接器使用该 schema 表示来确定每次插入、更新或删除操作发生时表的结构,并生成相应的变更事件。在一个独立的数据库 schema 历史 Kafka 主题中,连接器会记录所有 DDL 语句,以及每条 DDL 语句在 binlog 中出现的位置。
当连接器在崩溃或正常停止之后重新启动时,它会从一个特定的位置,即一个特定的时间点开始读取 binlog。连接器通过读取数据库 schema 历史 Kafka 主题,并解析到连接器开始读取位置为止的所有 DDL 语句,来重建该时间点存在的表结构。
该数据库 schema 历史主题仅供连接器内部使用。你还可以选择让连接器将 schema 变更事件发送到另一个面向消费者应用程序的主题。
当 MySQL 连接器捕获某个表的变更,而该表正在使用 gh-ost 或 pt-online-schema-change 等 schema 变更工具时,迁移过程中会创建一些辅助表。你必须配置连接器,使其捕获这些辅助表中发生的变更。如果消费者不需要连接器为辅助表生成的记录,可以配置单消息转换(SMT),从连接器发出的消息中移除这些记录。
其他资源
- 接收 Debezium 事件记录的主题默认名称。
Schema 变更主题
你可以配置 Debezium MySQL 连接器生成 schema 变更事件,用以描述应用于数据库表的 schema 变更。连接器将 schema 变更事件写入名为 <topicPrefix> 的 Kafka 主题,其中 topicPrefix 是 topic.prefix 连接器配置属性中指定的命名空间。连接器发送到 schema 变更主题的消息包含一个 payload,并且可以选择性地包含变更事件消息的 schema。
schema 变更事件的 schema 包含以下元素:
name
schema 变更事件消息的名称。
type
变更事件消息的类型。
version
schema 的版本。版本是一个整数,每当 schema 发生更改时就会递增。
fields
变更事件消息中包含的字段。
示例:MySQL 连接器 schema 变更主题的 schema
以下示例展示了一个典型的 JSON 格式 schema。
{
"schema": {
"type": "struct",
"fields": [
{
"type": "string",
"optional": false,
"field": "databaseName"
}
],
"optional": false,
"name": "io.debezium.connector.mysql.SchemaChangeKey",
"version": 1
},
"payload": {
"databaseName": "inventory"
}
}架构变更事件消息的负载包含以下元素:
ddl
提供导致架构变更的 SQL CREATE、ALTER 或 DROP 语句。
databaseName
DDL 语句所应用到的数据库名称。databaseName 的值用作消息键。
pos
这些语句在 binlog 中出现的位置。
tableChanges
架构变更后整张表架构的结构化表示。tableChanges 字段包含一个数组,其中列出了该表的每个列的条目。由于结构化表示以 JSON 或 Avro 格式呈现数据,消费者无需先通过 DDL 解析器处理,即可轻松读取消息。
| 对于处于捕获模式的表,连接器不仅会将架构变更历史存储到架构变更主题中,还会存储到内部数据库架构历史主题中。内部数据库架构历史主题仅供连接器使用,不打算供消费应用程序直接使用。请确保需要接收架构变更通知的应用程序仅从架构变更主题中消费这些信息。 |
|---|
切勿对数据库架构历史主题进行分区。为了让数据库架构历史主题正常工作,它必须保持连接器向其发出的事件记录的一致全局顺序。
为确保该主题不会被划分到多个分区中,可通过以下任一方法设置该主题的分区数:
- 如果你手动创建数据库架构历史主题,请将分区数指定为
1。 - 如果你使用 Apache Kafka 代理自动创建数据库架构历史主题,请在创建主题时,将 Kafka
num.partitions配置选项的值设置为1。
| 连接器向其架构变更主题发出的消息格式尚处于孵化阶段,可能会在不另行通知的情况下发生变化。 |
|---|
示例:向 MySQL 连接器架构变更主题发出的消息
以下示例展示了 JSON 格式的典型架构变更消息。该消息包含表架构的逻辑表示。
{
"schema": { },
"payload": {
"source": {
"version": "3.6.3.Final",
"connector": "mysql",
"name": "mysql",
"ts_ms": 1651535750218,
"ts_us": 1651535750218000,
"ts_ns": 1651535750218000000,
"snapshot": "false",
"db": "inventory",
"sequence": null,
"table": "customers",
"server_id": 223344,
"gtid": null,
"file": "mysql-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 语句。如果 MySQL 原子地应用这些语句,连接器会按顺序处理这些 DDL 语句,按数据库对它们进行分组,并为每个分组创建一个模式变更事件。如果 MySQL 逐条应用它们,连接器会为每条语句分别创建一个模式变更事件。
tableChanges
一个包含一个或多个项的数组,这些项包含 DDL 命令生成的模式变更。
type
描述变更的种类。该字段可以包含以下值之一:
CREATE
已创建表。
ALTER
已修改表。
DROP
已删除表。
id
指定被创建、修改或删除的表的完整标识符。在重命名表的情况下,该标识符是 <旧表名>,<新表名> 的连接。
table
表示应用变更之后的表元数据。
primaryKeyColumnNames
构成表主键的列的列表。
columns
变更表中每一列的元数据。
attributes
每个表变更的自定义属性元数据。
有关模式变更事件的更多信息,请参阅 模式历史主题。
快照
当 Debezium MySQL 连接器首次启动时,它会执行数据库的初始一致性快照。该快照使连接器能够为数据库的当前状态建立基线。
Debezium 在执行快照时可以使用不同的模式。快照模式由 snapshot.mode 配置属性决定。该属性的默认值为 initial。您可以通过更改 snapshot.mode 属性的值来自定义连接器创建快照的方式。
连接器在执行快照时会完成一系列任务。具体步骤因快照模式以及数据库当前生效的表锁定策略而异。Debezium MySQL 连接器在执行使用全局读锁或表级锁的初始快照时,完成的步骤各不相同。
使用全局读锁的初始快照
你可以通过修改 snapshot.mode 属性的值来定制连接器创建快照的方式。如果你配置了不同的快照模式,连接器将使用此工作流的修改版本来完成快照。有关在不允许使用全局读锁的环境中进行快照处理的信息,请参阅表级锁的快照工作流。
Debezium MySQL 连接器使用全局读锁执行初始快照的默认工作流
下表显示了 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 位置处,连接器开始扫描指定要捕获的表。在扫描过程中,连接器完成以下任务:
- 确认该表是在快照开始之前创建的。如果该表是在快照开始之后创建的,连接器会跳过该表。快照完成并转入流式传输阶段后,连接器会为快照开始之后创建的任何表发出变更事件。
- 为从表中捕获的每一行生成一个
read事件。所有read事件都包含相同的 binlog 位置,即第 5 步中获取的位置。 - 将每个
read事件发送到源表对应的 Kafka 主题。 - 如适用,释放数据表锁。
9
提交事务。
10
在连接器偏移量中记录快照已成功完成。
由此产生的初始快照捕获了被捕获表中每一行的当前状态。以该基线状态为基础,连接器会捕获此后发生的变更。
快照过程开始后,如果由于连接器故障、重新平衡或其他原因导致该过程中断,那么连接器重启后该过程也会重新开始。
连接器完成初始快照后,会从第 5 步中读取的位置继续流式传输,以确保不会遗漏任何更新。
如果连接器因任何原因再次停止,重启后它会从之前中断的位置继续流式传输变更。
连接器重启后,如果日志已被清理,连接器在日志中的位置可能已不可用。此时连接器会失败,并返回一个指示需要进行新快照的错误。要让连接器在此情况下自动发起快照,请将 snapshot.mode 属性的值设置为 when_needed。有关 Debezium MySQL 连接器故障排查的更多提示,请参阅出现问题时的行为。
使用表级锁的初始快照
在某些数据库环境中,管理员不允许使用全局读锁。如果 Debezium MySQL 连接器检测到不允许使用全局读锁,连接器在执行快照时会使用表级锁。要使连接器执行使用表级锁的快照,Debezium 连接器用于连接 MySQL 的数据库账户必须具有 LOCK TABLES 权限。
Debezium MySQL 连接器使用表级锁执行初始快照的默认工作流
下表展示了 Debezium 使用表级读锁创建快照时遵循的工作流步骤。有关在不允许使用全局读锁的环境中进行快照的信息,请参阅全局读锁的快照工作流。
| 步骤 | 操作 |
|---|---|
| 1 | 建立与数据库的连接。 |
| 2 | 确定要捕获的表。默认情况下,连接器捕获所有非系统表。要让连接器捕获部分表或表元素,可以设置若干 include 和 exclude 属性来过滤数据,例如 table.include.list 或 table.exclude.list。 |
| 3 | 获取表级锁。 |
| 4 | 启动具有可重复读语义的事务,以确保事务内的所有后续读取都是针对一致性快照进行的。 |
| 5 | 读取当前的 binlog 位置。 |
| 6 | 读取连接器已配置为捕获变更的数据库和表的架构。连接器会将架构信息持久化到其内部数据库架构历史主题中,包括所有必需的 DROP… 和 CREATE… DDL 语句。 |
架构历史提供了变更事件发生时生效的结构信息。
默认情况下,连接器会捕获数据库中每张表的结构,包括那些未配置为捕获的表。如果表未配置为捕获,初始快照只会捕获其结构,不会捕获任何表数据。
有关为什么快照会持久化未包含在初始快照中的表的结构信息的详细说明,请参阅了解为何初始快照会捕获所有表的结构。
7
在第 5 步中连接器读取到的 binlog 位置处,连接器开始扫描被指定为捕获的表。在扫描过程中,连接器完成以下任务:
- 确认该表是在快照开始之前创建的。如果该表是在快照开始之后创建的,连接器会跳过该表。快照完成后,当连接器转换为流式处理时,它会为快照开始之后创建的任何表发出变更事件。
- 为从表中捕获的每一行生成一个
read事件。所有read事件都包含相同的 binlog 位置,即在第 5 步中获取的位置。 - 将每个
read事件发送到源表对应的 Kafka 主题。 - 释放数据表锁(如果适用)。
8
提交事务。
9
释放表级锁。其他数据库客户端现在可以向之前被锁定的任何表写入数据。
10
在连接器偏移量中记录快照已成功完成。
表 1. snapshot.mode 连接器配置属性的设置
设置说明
always
连接器每次启动时都会执行快照。快照包括被捕获表的结构和数据。指定此值后,连接器每次启动时都会用被捕获表数据的完整表示来填充主题。快照完成后,连接器开始流式传输后续数据库变更的事件记录。
initial
连接器按照创建初始快照的默认工作流中所述执行数据库快照。快照完成后,连接器开始流式传输后续数据库变更的事件记录。
initial_only
连接器执行数据库快照。快照完成后,连接器停止运行,不会流式传输后续数据库变更的事件记录。
schema_only
已弃用,请参阅 no_data。
no_data
连接器捕获所有相关表的结构,执行创建初始快照的默认工作流中描述的所有步骤,但不会创建 READ 事件来表示连接器启动时的数据集(第 7.2 步)。
schema_only_recovery
已弃用,请参阅 recovery。
recovery
将此选项设置为恢复丢失或损坏的数据库 schema history 主题。重启后,连接器会运行一次快照,根据源表重建该主题。你也可以设置该属性,定期清理意外增长的数据库 schema history 主题。
警告:如果在连接器上次关闭之后数据库中已提交了 schema 变更,请勿使用此模式执行快照。
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。
理解为何初始快照会捕获所有表的 schema 历史
连接器运行的初始快照会捕获两类信息:
表数据
关于连接器 table.include.list 属性中所列表的 INSERT、UPDATE 和 DELETE 操作的信息。
Schema 数据
描述应用于表的结构性变更的 DDL 语句。Schema 数据会持久化到内部 schema history 主题中;如果配置了 schema change 主题,也会持久化到连接器的 schema change 主题中。
运行初始快照后,你可能会注意到快照捕获了未指定为捕获目标的表的 schema 信息。默认情况下,初始快照被设计为捕获数据库中存在的每一张表的 schema 信息,而不仅仅是被指定为捕获目标的表。连接器要求在捕获某张表之前,该表的 schema 必须已存在于 schema history 主题中。通过让初始快照捕获不属于原始捕获集合的表的 schema 数据,Debezium 使连接器能够在日后需要时轻松捕获这些表的事件数据。如果初始快照未捕获某张表的 schema,则必须先将该 schema 添加到历史主题中,连接器才能从该表捕获数据。
在某些情况下,你可能希望在初始快照中限制架构捕获。当你希望减少完成快照所需的时间时,或者当 Debezium 通过一个可访问多个逻辑数据库的用户帐户连接到数据库实例,但你只希望连接器从特定逻辑数据库中的表捕获更改时,这非常有用。
更多信息
- 从初始快照未捕获的表中捕获数据(无架构变更)
- 从初始快照未捕获的表中捕获数据(架构变更)
- 设置
schema.history.internal.store.only.captured.tables.ddl属性以指定从哪些表捕获架构信息。 - 设置
schema.history.internal.store.only.captured.databases.ddl属性以指定从哪些逻辑数据库捕获架构变更。
从初始快照未捕获的表中捕获数据(无架构变更)
在某些情况下,你可能希望连接器从初始快照未捕获其架构的表中捕获数据。根据连接器的配置,初始快照可能仅捕获数据库中特定表的架构。如果历史主题中不存在该表的架构,连接器将无法捕获该表,并报告缺少架构的错误。
你仍然可能从该表中捕获数据,但必须执行额外的步骤来添加表架构。
前提条件
- 你希望从一张架构在初始快照期间未被连接器捕获的表中捕获数据。
- 在事务日志中,该表的所有条目均使用相同的架构。有关从经历结构变更的新表中捕获数据的信息,请参阅从初始快照未捕获的表中捕获数据(架构变更)。
操作步骤
停止连接器。
删除由
schema.history.internal.kafka.topic属性指定的内部数据库架构历史主题。对连接器配置应用以下更改:
将
snapshot.mode设置为recovery。将
schema.history.internal.store.only.captured.tables.ddl的值设置为false。将你要连接器捕获的表添加到
table.include.list中。这样可以保证在将来,连接器能够为所有表重建架构历史。重启连接器。快照恢复过程会根据表的当前结构重新构建架构历史。
(可选)快照完成后,启动增量快照,以捕获新添加表中的现有数据,以及在连接器离线期间其他表发生的变更。
(可选)将
snapshot.mode重置回no_data,以防止连接器在将来重启后再次发起恢复。
从初始快照未捕获的表中捕获数据(架构变更)
如果对表应用了架构变更,则在架构变更之前提交的记录与在变更之后提交的记录结构不同。当 Debezium 从表中捕获数据时,它会读取架构历史,以确保为每个事件应用正确的架构。如果架构主题中不存在该架构,连接器将无法捕获该表,并会导致错误。
如果你要从初始快照未捕获的表中捕获数据,并且该表的架构已被修改,则必须将该架构添加到历史主题中(如果尚未存在)。你可以通过运行新的架构快照或为该表运行初始快照来添加架构。
前提条件
- 你要从架构未在初始快照中被连接器捕获的表中捕获数据。
- 该表已应用架构变更,导致待捕获的记录不具有统一的结构。
操作步骤
初始快照已捕获所有表的架构(store.only.captured.tables.ddl 设置为 false)
- 编辑
table.include.list属性,指定你要捕获的表。 - 重启连接器。
- 如果你要从新添加的表中捕获现有数据,请启动增量快照。
初始快照未捕获所有表的架构(store.only.captured.tables.ddl 设置为 true)
如果初始快照未保存你要捕获的表的架构,请执行以下任一操作步骤:
操作步骤 1:架构快照,随后执行增量快照
在此操作步骤中,连接器首先执行架构快照。然后你可以启动增量快照,使连接器能够同步数据。
停止连接器。
删除由
schema.history.internal.kafka.topic属性指定的内部数据库模式历史主题。清除已配置的 Kafka Connect
offset.storage.topic中的偏移量。有关如何删除偏移量的更多信息,请参阅 Debezium 社区常见问题。删除偏移量的操作应仅由具有操作 Kafka Connect 内部数据经验的高级用户执行。此操作具有潜在的破坏性,应仅作为最后手段使用。 按以下步骤为连接器配置中的属性设置值:
- 将
snapshot.mode属性的值设置为no_data。 - 编辑
table.include.list,添加你想要捕获的表。
- 将
重启连接器。
等待 Debezium 捕获新表和现有表的模式。连接器停止之后任何表中发生的数据变更都不会被捕获。
为确保不丢失数据,请启动增量快照。
过程 2:初始快照,随后可选进行增量快照
在此过程中,连接器对数据库执行完整的初始快照。与任何初始快照一样,如果数据库中包含许多大表,运行初始快照可能是一项耗时的操作。快照完成后,你可以选择触发增量快照,以捕获连接器离线期间发生的任何变更。
- 停止连接器。
- 删除由
schema.history.internal.kafka.topic属性指定的内部数据库模式历史主题。 - 清除已配置的 Kafka Connect
offset.storage.topic中的偏移量。有关如何删除偏移量的更多信息,请参阅 Debezium 社区常见问题。
| 移除偏移量应仅由具备操作 Kafka Connect 内部数据经验的高级用户执行。此操作具有潜在破坏性,只应作为最后的手段。 | |
|---|---|
4. 编辑 table.include.list,添加你想要捕获的表。 |
|
| 5. 按照以下步骤为连接器配置中的属性设置值: |
- 将
snapshot.mode属性的值设置为initial。 - (可选)将
schema.history.internal.store.only.captured.tables.ddl设置为false。 - 重启连接器。连接器会执行完整数据库快照。快照完成后,连接器将转入流式传输阶段。
- (可选)若要捕获连接器离线期间发生的任何数据变更,请启动增量快照。
基于分块的并行快照
基于分块的并行快照通过将较小的工作单元分配到多个线程来加速初始快照,从而改善工作负载均衡并提升快照的容错能力。
| 基于分块的并行快照是一项孵化中特性。 |
|---|
当 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 |
不适用 | 一个数组,包含与要纳入快照的表的完全限定名相匹配的正则表达式。对于 MySQL 连接器,请使用以下格式指定表的完全限定名: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\""。 |
|---|
先决条件
-
- 源数据库上存在信号数据集合。
signal.data.collection属性中已指定该信号数据集合。
使用源信号通道触发增量快照
发送 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 支持以下物理行标识符作为代理键:
| 数据库 | 标识符 | 说明 |
|---|---|---|
| Oracle | ROWID | Oracle 表中某一行的物理地址。使用 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(主键)colorquantity
如果你想让 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 消息的 key 必须与 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"}}`带有附加条件的即席增量快照
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 表,如果希望快照仅包含 color='blue' 且 brand='MyBrand' 的 products 表内容,可以发送以下请求:
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 检测到信号表中的更改后,会读取该信号,并在增量快照操作正在进行时将其停止。
其他资源
前提条件
-
- 源数据库上存在信号数据集合。
signal.data.collection属性中指定了信号数据集合。
使用源信号通道停止增量快照
向信号表发送 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 消息的 key 必须与 topic.prefix 连接器配置选项的值匹配。
消息的值是一个包含 type 和 data 字段的 JSON 对象。
信号类型为 stop-snapshot,data 字段必须包含以下字段:
| 字段 | 默认值 | 值 |
|---|---|---|
type | incremental | 要执行的快照类型。目前 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 MySQL 连接器支持以只读方式连接数据库来运行增量快照。在只读访问下运行增量快照时,连接器使用已执行的全局事务 ID(GTID)集合作为高低水位标记。通过将二进制日志(binlog)事件或服务器心跳的 GTID 与低水位标记和高水位标记进行比较,来更新块窗口的状态。
要切换到只读实现,请将 read.only 属性的值设置为 true。
前提条件
如果连接器从多线程副本读取(即
replica_parallel_workers的值大于0的副本),必须设置以下选项之一:replica_preserve_commit_order=ONslave_preserve_commit_order=ON
当 MySQL 连接为只读时,你可以使用任何可用的信号通道,而无需使用 source 通道。
快照事件的操作类型
MySQL 连接器将快照事件作为 READ 操作发出("op" : "r")。如果你希望连接器将快照事件作为 CREATE(c)事件发出,请配置 Debezium 的 ReadToInsertEvent 单消息转换器(SMT)来修改事件类型。
以下示例展示了如何配置该 SMT:
示例:使用 ReadToInsertEvent SMT 更改快照事件的类型
transforms=snapshotasinsert,...
transforms.snapshotasinsert.type=io.debezium.connector.mysql.transforms.ReadToInsertEvent自定义快照器 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 信号,动态指定 MySQL 连接器从数据库 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 信号,以指定 MySQL 连接器的 binlog 读取位置。
前置条件
- 连接器已配置为使用其中一种可用的 Debezium 信号通道。
- 源数据库中存在信号数据集合。
- 已设置连接器配置属性
heartbeat.interval.ms,以确保偏移量更改被持久化。
操作步骤
确定数据库历史记录中的目标位置(binlog 文件与位置,或 GTID 集合)。
发送一个
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": "mysql-bin.000003", "binlog_position": 1234}' );
启用 GTID 时使用的 SQL(指定基于 UUID 的 GTID 集合)
INSERT INTO debezium_signal (id, type, data) VALUES (
'set-gtid-001',
'set-binlog-position',
'{"gtid_set": "3E11FA47-71CA-11E1-9E33-C80AA9429562:1-100"}'
);- 重启连接器,从新位置开始流式传输。
主题名称
默认情况下,MySQL 连接器会将表中发生的所有 INSERT、UPDATE 和 DELETE 操作的变更事件写入该表对应的单个 Apache Kafka 主题中。
连接器使用以下约定来命名变更事件主题:
topicPrefix.databaseName.tableName
假设主题前缀为 fulfillment,数据库名称为 inventory,该数据库包含名为 orders、customers 和 products 的表。Debezium MySQL 连接器将事件发送到三个 Kafka 主题,每个表对应一个主题:
fulfillment.inventory.orders
fulfillment.inventory.customers
fulfillment.inventory.products以下列表定义了默认名称的各组成部分:
topicPrefix
由 topic.prefix 连接器配置属性指定的主题前缀。
schemaName
发生该操作的模式名称。
tableName
发生该操作的表名称。
连接器采用类似的命名约定来标记其内部数据库模式历史主题、模式变更主题 以及事务元数据主题。
如果默认主题名称不能满足你的需求,可以配置自定义主题名称。要配置自定义主题名称,需要在逻辑主题路由 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 MySQL 连接器会为每一行的 INSERT、UPDATE 和 DELETE 操作生成一个数据变更事件。每个事件都包含一个键和一个值。键和值的结构取决于发生变更的表。
Debezium 和 Kafka Connect 是围绕事件消息的持续流设计的。然而,这些事件的结构可能会随时间发生变化,这对消费者来说可能难以处理。为了解决这个问题,每个事件都包含其内容的 schema,或者(如果你使用 schema 注册表)包含一个 schema ID,消费者可以用它从注册表中获取 schema。这样每个事件都是自包含的。
下面的骨架 JSON 展示了变更事件的四个基本组成部分。Debezium 在变更消息中表示这些组件的具体方式,取决于你在应用中如何配置 Kafka Connect 转换器。只有当你配置转换器生成 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 结构相对应,其中包含被更改行的实际数据。
默认情况下,连接器将变更事件记录流式传输到名称与其事件源表相同的主题。有关更多信息,请参阅主题名称。
MySQL 连接器确保所有 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 表变更的变更事件都具有相同的事件键模式。因此,Debezium 为该表发出的每个变更事件,其事件消息都具有相同的键结构。只要 customers 表的模式保持不变,键的一致性就会保持稳定。上述表定义会产生如下 JSON 示例所示的键结构:
{
"schema": {
"type": "struct",
"name": "mysql-server-1.inventory.customers.Key",
"optional": false,
"fields": [
{
"field": "id",
"type": "int32",
"optional": false
}
]
},
"payload": {
"id": 1001
}
}下面的列表描述了 Debezium 从 customers 表发出的变更事件消息中,事件键的部分元素:
schema
事件键的 schema 元素指定了定义事件键 payload 部分中结构和数据类型的 Kafka Connect schema。
name
指定定义键载荷(payload)结构的 schema 名称。该 schema 描述发生变更的表的主键结构。事件键的 schema 名称格式如下:connector-name.database-name.table-name.Key。
在上面的事件消息中,schema 名称由以下几部分组成:
mysql-server-1
指定生成此事件的连接器名称。
inventory
指定包含发生变更的表的数据库。
customers
指定发生变更的表。
fields
指定 payload 中的必填字段。fields 元素包含以下属性:
field
字段的名称,例如 id。
type
字段的数据类型,例如 int32。
optional
一个布尔值,用于指定事件键的载荷字段中是否必须有值。当表没有主键时,键的载荷字段中的值是可选的。
payload
指定生成此变更事件的行的键。在上面的示例中,payload 只包含一个字段 id,其值为 1001。
变更事件值
变更事件值包含操作前后整行的状态,以及源元数据和操作类型,其结构为 Kafka Connect 的 Envelope。
变更事件中的值比键稍微复杂一些。与键一样,值也包含 schema 部分和 payload 部分。schema 部分包含描述 payload 部分 Envelope 结构的 schema,其中包括其嵌套字段。创建、更新或删除数据的操作所产生的变更事件,其值载荷都具有信封(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;针对该表的变更,变更事件的 value 部分描述如下:
create 事件
create 事件包含 INSERT 操作的完整行数据,该新行的状态记录在事件 value 负载的 after 字段中。
下面的示例展示了连接器针对在 customers 表中创建数据的操作所生成的变更事件的 value 部分:
{
"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": "mysql-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": "mysql-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.mysql.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": "mysql-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": "mysql",
"name": "mysql-server-1",
"ts_ms": 0,
"ts_us": 0,
"ts_ns": 0,
"snapshot": false,
"db": "inventory",
"table": "customers",
"server_id": 0,
"gtid": null,
"file": "mysql-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 字段指定事件值负载中某个字段的模式。
mysql-server-1.inventory.customers.Value 是负载中 before 和 after 字段的模式。该模式是 customers 表所特有的。
before 和 after 字段的模式名称采用 logicalName.tableName.Value 的形式。这种格式可确保模式名称在数据库内是唯一的。在使用 Avro 转换器 的环境中,拥有唯一的模式名称可确保每个逻辑源中每张表的 Avro 模式都有其独立的演进历史。
"name": "io.debezium.connector.mysql.Source"
io.debezium.connector.mysql.Source 是负载中 source 字段的模式。该模式是 MySQL 连接器所特有的,连接器在其生成的所有事件中都会使用它。
"name": "mysql-server-1.inventory.customers.Envelope"
指定负载整体结构的模式名称。该模式名称由以下几部分组成:
mysql-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
MySQL 服务器 ID(如果可用)。
file
记录该事件的二进制日志名称。
pos
二进制日志位置。
row
事件中的行位置。
thread
创建该事件的 MySQL 线程 ID(仅非快照事件)。
query
在源表上执行该操作的 SQL 命令。
update 事件
在示例 customers 表中,针对更新操作的变更事件值与该表的 create(创建)事件具有相同的模式。同样,事件值的负载(payload)结构也相同。不过,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": "mysql-server-1",
"connector": "mysql",
"ts_ms": 1465581029100,
"ts_us": 1465581029100000,
"ts_ns": 1465581029100000000,
"snapshot": false,
"db": "inventory",
"table": "customers",
"server_id": 223344,
"gtid": null,
"file": "mysql-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 事件来自二进制日志中的不同位置。update 事件值中的 source 字段提供以下元数据:
version
Debezium 版本。
name
连接器名称。
connector
连接器的类型。
ts_ms, ts_us, ts_ns
以毫秒、微秒和纳秒为单位显示时间戳,表示该更改在数据库中发生的时间。
snapshot
指定该事件是否由快照操作产生。
db, table
包含该新行的数据库和表的名称。
server_id
MySQL 服务器 ID(如果可用)。
file
记录该事件的二进制日志文件名。
pos
二进制日志位置。
row
事件中的行号。
thread
创建该事件的 MySQL 线程 ID(仅限非快照事件)。
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事件。 - 一个 Tombstone 事件,用于指定该行的旧键。
- 一个
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": "mysql",
"name": "mysql-server-1",
"ts_ms": 1465581902300,
"ts_us": 1465581902300000,
"ts_ns": 1465581902300000000,
"snapshot": false,
"db": "inventory",
"table": "customers",
"server_id": 223344,
"gtid": null,
"file": "mysql-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
MySQL 服务器 ID(如果可用)。
file
记录该事件的二进制日志文件名。
pos
二进制日志位置。
row
事件中的行。
thread
创建该事件的 MySQL 线程 ID(仅非快照时)。
query
在源表上执行该操作的 SQL 命令。
op
指定操作类型。在此示例中,d 值表示执行了 DELETE 操作,导致删除一行。
ts_ms、ts_us、ts_ns
以毫秒、微秒和纳秒显示时间戳,指示连接器处理该事件的时间。该时间基于运行 Kafka Connect 任务的 JVM 中的系统时钟。
在源对象中,这些字段表示更改在数据库中发生的时间。通过比较 payload.source.ts_ms 的值与 payload.ts_ms 的值,可以确定源数据库更新与 Debezium 之间的延迟。
delete 变更事件记录为消费者提供了处理该行移除所需的信息。之所以包含旧值,是因为某些消费者可能需要这些值才能正确处理移除操作。
MySQL 连接器事件的设计与 Kafka 日志压缩兼容。日志压缩允许删除部分较旧的消息,只要保留每个键的最新消息即可。这样,Kafka 既能回收存储空间,又能确保主题包含可用于重新加载基于键的状态的完整数据集。
墓碑事件
当某行被删除时,delete 事件的值仍然可以配合日志压缩工作,因为 Kafka 能够移除所有具有相同键的较早消息。但是,要让 Kafka 移除所有具有相同键的消息,消息的值必须为 null。为实现这一点,在 Debezium MySQL 连接器发出 delete 事件之后,连接器会发出一个特殊的墓碑(tombstone)事件,该事件具有相同的键,但值为 null。
truncate 事件
truncate 变更事件表示某张表已被截断。truncate 事件的消息键为 null。消息值类似于以下示例:
{
"schema": { ... },
"payload": {
"source": {
"version": "3.6.3.Final",
"name": "mysql-server-1",
"connector": "mysql",
"name": "mysql-server-1",
"ts_ms": 1465581029100,
"ts_us": 1465581029100000,
"ts_ns": 1465581029100000000,
"snapshot": false,
"db": "inventory",
"table": "customers",
"server_id": 223344,
"gtid": null,
"file": "mysql-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
MySQL 服务器 ID(如果可用)。
file
记录该事件的二进制日志文件名。
pos
二进制日志位置。
row
事件中的行。
thread
创建该事件的 MySQL 线程 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 MySQL 连接器使用结构与行所在表相同的事件来表示行的更改。事件中为每个列值包含一个字段。该列的 MySQL 数据类型决定了 Debezium 在事件中如何表示该值。
存储字符串的列在 MySQL 中通过字符集和排序规则来定义。MySQL 连接器在读取 binlog 事件中列值的二进制表示时,会使用该列的字符集。
连接器可以将 MySQL 数据类型映射为字面类型和语义类型。
- 字面类型(Literal type):使用 Kafka Connect 模式类型来表示值的方式。
- 语义类型(Semantic type):Kafka Connect 模式如何捕获字段的含义(模式名称)。
如果默认的数据类型转换不能满足你的需求,你可以为连接器创建自定义转换器。
基本类型
下表展示了连接器如何映射基本的 MySQL 数据类型。
表 5. 基本类型映射说明
MySQL 类型 字面类型 语义类型
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
无
REAL[(M,D)]
FLOAT32
n/a
FLOAT[(P)]
FLOAT32 或 FLOAT64
精度仅用于确定存储大小。精度 P 取值在 0 到 23 之间时,会生成 4 字节的单精度 FLOAT32 列;精度 P 取值在 24 到 53 之间时,会生成 8 字节的双精度 FLOAT64 列。
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 格式,精度为微秒。MySQL 允许小数秒精度 fsp 的取值范围为 0-6。
.
时间类型
除 TIMESTAMP 数据类型外,MySQL 的时间类型取决于 time.precision.mode 连接器配置属性的值。对于默认值指定为 CURRENT_TIMESTAMP 或 NOW 的 TIMESTAMP 列,在 Kafka Connect 架构中使用 1970-01-01 00:00:00 作为默认值。
MySQL 允许 DATE、DATETIME 和 TIMESTAMP 列取零值,因为有时零值比空值更合适。当列定义允许空值时,MySQL 连接器将零值表示为空值;当列不允许空值时,则表示为 epoch 天数。
不带时区的时间值
DATETIME 类型表示本地日期和时间,例如 "2018-01-13 09:48:27"。可以看到,其中不含时区信息。此类列会根据其精度使用 UTC 转换为 epoch 毫秒或微秒。TIMESTAMP 类型表示不带时区信息的时间戳。写入时 MySQL 会将其从服务器(或会话)的当前时区转换为 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。默认情况下,时区将从服务器查询获取。
运行 Kafka Connect 和 Debezium 的 JVM 的时区设置不会影响这些转换。
关于时间值相关属性的更多细节,请参阅 MySQL 连接器配置属性文档。
time.precision.mode=adaptive_time_microseconds(默认)
MySQL 连接器根据列的数据类型定义来确定字面类型和语义类型,从而使事件准确表示数据库中的值。所有时间字段均以微秒为单位。只有取值范围在 00:00:00.000000 到 23:59:59.999999 之间的正值 TIME 字段才能被正确捕获。
表 6. time.precision.mode=adaptive_time_microseconds 时的映射关系
MySQL 类型 字面类型 语义类型
DATE
INT32
io.debezium.time.Date
表示自 epoch 以来的天数。
TIME[(fsp)]
INT64
io.debezium.time.MicroTime
以微秒为单位表示时间值,不包含时区信息。MySQL 允许的秒小数精度 fsp 范围为 0-6。
DATETIME, DATETIME(0), DATETIME(1), DATETIME(2), DATETIME(3)
INT64
io.debezium.time.Timestamp
表示从纪元(epoch)起经过的毫秒数,不包含时区信息。
DATETIME(4), DATETIME(5), DATETIME(6)
INT64
io.debezium.time.MicroTimestamp
表示从纪元(epoch)起经过的微秒数,不包含时区信息。
time.precision.mode=connect
MySQL 连接器使用 Kafka Connect 中定义的逻辑类型。这种方式的精确度低于默认方式,如果数据库列的秒小数精度大于 3,事件的精确度可能会降低。只能处理 00:00:00.000 到 23:59:59.999 范围内的值。只有在能够确保表中的 TIME 值始终不超过受支持范围时,才应设置 time.precision.mode=connect。connect 设置预计将在 Debezium 的未来版本中被移除。
表 7. time.precision.mode=connect 时的映射
MySQL 类型 字面量类型 语义类型
DATE
INT32
org.apache.kafka.connect.data.Date
表示从纪元(epoch)起经过的天数。
TIME[(fsp)]
INT64
org.apache.kafka.connect.data.Time
表示自午夜起经过的微秒数所对应的时间值,不包含时区信息。
DATETIME[(fsp)]
INT64
org.apache.kafka.connect.data.Timestamp
表示从纪元(epoch)起经过的毫秒数,不包含时区信息。
time.precision.mode=isostring
MySQL 连接器将日期、时间和日期时间值表示为 UTC 时区下符合 ISO-8601 格式的字符串。由于这些值以字符串形式表示,因此该模式保留了 connect 模式可能会丢失的秒小数精度。
表 8. time.precision.mode=isostring 时的映射
MySQL 类型 字面量类型 语义类型
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
按照 ISO 8601 标准,以 UTC 格式表示日期时间值,例如 2019-07-09T02:28:57.123456Z。
| | time.precision.mode 连接器属性还接受 microseconds 和 nanoseconds 两个值,尽管上述映射表中目前并未展示这两种模式。MySQL 的 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 时的映射关系
| MySQL 类型 | 字面量类型 | 语义类型 |
|---|---|---|
NUMERIC[(M[,D])] | BYTES | org.apache.kafka.connect.data.Decimalscale 架构参数包含一个整数,表示小数点移动了多少位。 |
DECIMAL[(M[,D])] | BYTES | org.apache.kafka.connect.data.Decimalscale 架构参数包含一个整数,表示小数点移动了多少位。 |
decimal.handling.mode=double
| MySQL 类型 | 字面类型 | 语义类型 |
|---|---|---|
NUMERIC[(M[,D])] | FLOAT64 | 不适用 |
DECIMAL[(M[,D])] | FLOAT64 | 不适用 |
表 10. decimal.handling.mode=double 时的映射关系
decimal.handling.mode=string
| MySQL 类型 | 字面类型 | 语义类型 |
|---|---|---|
NUMERIC[(M[,D])] | STRING | 不适用 |
DECIMAL[(M[,D])] | STRING | 不适用 |
表 11. decimal.handling.mode=string 时的映射关系
布尔值
MySQL 对 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 MySQL 连接器支持以下空间数据类型。
表 12. 空间类型映射说明
MySQL 类型 字面量类型 语义类型
GEOMETRY
LINESTRING
POLYGON
MULTIPOINT
MULTILINESTRING
MULTIPOLYGON
GEOMETRYCOLLECTION
STRUCT
io.debezium.data.geometry.Geometry
包含一个具有两个字段的结构:
srid (INT32:空间参考系统 ID,用于定义该结构中存储的几何对象类型wkb (BYTES):几何对象的二进制表示,采用众所周知二进制(wkb)格式编码。更多详细信息请参阅开放地理空间联盟。
向量类型
向量类型
目前,Debezium MySQL 连接器支持以下向量数据类型。
| MySQL 类型 | 字面量类型 | 语义类型 |
|---|---|---|
VECTOR | ARRAY (FLOAT32) | io.debezium.data.FloatVector |
表 13. 向量类型映射说明
自定义转换器
默认情况下,Debezium MySQL 连接器为 MySQL 数据类型提供了多个 CustomConverter 实现。这些自定义转换器根据连接器配置,为特定数据类型提供替代映射。要向连接器添加 CustomConverter,请按照自定义转换器文档中的说明进行操作。
TINYINT(1) 转换为布尔值
在连接器快照期间,Debezium MySQL 连接器默认从 JDBC 驱动程序获取列类型,该驱动程序会将 TINYINT(1) 类型赋给 BOOLEAN 列。Debezium 随后使用这些 JDBC 列类型来定义快照事件的模式。当连接器从快照阶段切换到流式传输阶段后,默认映射产生的变更事件模式可能导致 BOOLEAN 列的映射不一致。为确保 MySQL 统一输出 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 MySQL source 连接器集成,MySQL 连接器在快照阶段和流式阶段会以不同的方式发出某些列属性。为了使 JDBC sink 连接器能够一致地消费来自快照阶段和流式阶段的变更,你必须在 MySQL source 连接器配置中包含 JdbcSinkDataTypesConverter 转换器,如以下示例所示:
示例: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.types | DATE,DATETIME,TIMESTAMP | 以逗号分隔的 DATE、DATETIME、TIMESTAMP 子集,用于控制该转换器应用于哪些 JDBC 时间类型。 |
fallback.date | NULL | 零值日期(zero-date)DATE 行所输出的值。设为 NULL(或不设置)时,列模式会提升为可选,并输出 null;否则输出 ISO-8601 日期字面量(yyyy-MM-dd)。 |
fallback.datetime | NULL | 零值日期 DATETIME 行所输出的值。设为 NULL(或不设置)时,列模式会提升为可选,并输出 null;否则输出 UTC 日期时间字面量(yyyy-MM-dd HH:mm:ss)。 |
fallback.timestamp | NULL | 零值日期 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 保证。 |
|---|
设置 MySQL
在安装和运行 Debezium 连接器之前,需要先完成一些 MySQL 的设置任务。
创建用户
Debezium MySQL 连接器需要一个 MySQL 用户账户。该 MySQL 用户必须在 Debezium MySQL 连接器捕获变更的所有数据库上拥有相应的权限。
前提条件
- 一个 MySQL 服务器。
- SQL 命令的基础知识。
操作步骤
创建 MySQL 用户:
mysql> CREATE USER 'user'@'localhost' IDENTIFIED BY 'password';为该用户授予所需的权限:
mysql> GRANT SELECT, RELOAD, SHOW DATABASES, REPLICATION SLAVE, REPLICATION CLIENT ON *.* TO 'user' IDENTIFIED BY 'password';
有关所需权限的说明,请参阅用户权限说明。
如果使用 Amazon RDS 或 Amazon Aurora 这类不允许全局读锁的托管服务,则会使用表级锁来创建一致性快照。在这种情况下,你还需要为你创建的用户授予 LOCK TABLES 权限。更多详情请参阅快照。 |
|
|---|---|
| 3. 完成用户权限的设置: |
mysql> FLUSH PRIVILEGES;用户权限说明
下表介绍了必须授予 Debezium 连接器用户账户的 MySQL 权限,并说明了每项权限的必要性。
表 14. 用户权限说明
关键字 说明
SELECT
使连接器能够从数据库的表中查询数据。仅在执行快照时使用。
RELOAD
使连接器能够使用 FLUSH 语句清除或重新加载内部缓存、刷新表或获取锁。仅在执行快照时使用。
SHOW DATABASES
使连接器能够通过执行 SHOW DATABASE 语句查看数据库名称。仅在执行快照时使用。
REPLICATION SLAVE
使连接器能够连接并读取 MySQL 服务器的 binlog。
REPLICATION CLIENT
使连接器能够使用以下语句:
SHOW MASTER STATUSSHOW SLAVE STATUSSHOW BINARY LOGS
连接器始终需要此权限。
ON
指定权限适用的数据库。
TO 'user'
指定要授予这些权限的用户。
IDENTIFIED BY 'password'
指定该用户的 MySQL 密码。
启用 binlog
你必须为 MySQL 复制启用二进制日志记录。二进制日志以副本能够传播这些更改的方式记录事务更新。
前提条件
- 一个 MySQL 服务器。
- 适当的 MySQL 用户权限。
操作步骤
检查
log-bin选项是否已启用:// for MySQL 5.x mysql> SELECT variable_value as "BINARY LOGGING STATUS (log-bin) ::" FROM information_schema.global_variables WHERE variable_name='log_bin'; // for MySQL 8.x mysql> SELECT variable_value as "BINARY LOGGING STATUS (log-bin) ::" FROM performance_schema.global_variables WHERE variable_name='log_bin';如果 binlog 处于
OFF状态,请将下表中的属性添加到 MySQL 服务器的配置文件中:server-id = 223344 # 查询变量名为 server_id,例如 SELECT variable_value FROM information_schema.global_variables WHERE variable_name='server_id'; log_bin = mysql-bin binlog_format = ROW binlog_row_image = FULL binlog_expire_logs_seconds = 864000再次检查 binlog 状态以确认您的更改:
// for MySQL 5.x mysql> SELECT variable_value as "BINARY LOGGING STATUS (log-bin) ::" FROM information_schema.global_variables WHERE variable_name='log_bin'; // for MySQL 8.x mysql> SELECT variable_value as "BINARY LOGGING STATUS (log-bin) ::" FROM performance_schema.global_variables WHERE variable_name='log_bin';
- 如果你在 Amazon RDS 上运行 MySQL,必须为数据库实例启用自动备份,二进制日志才会生成。如果数据库实例未配置为执行自动备份,即使你应用了前面步骤中描述的设置,binlog 也是禁用的。
MySQL binlog 配置属性说明
下表描述了必须设置的 MySQL 服务器配置属性,用于启用二进制日志并使其可供 Debezium 连接器使用。
| 属性 | 说明 |
|---|---|
server-id |
server-id 的值在 MySQL 集群中对每个服务器和复制客户端都必须唯一。 |
log_bin |
log_bin 的值是 binlog 文件序列的基本名称。 |
binlog_format |
binlog-format 必须设置为 ROW 或 row。 |
binlog_row_image |
binlog_row_image 必须设置为 FULL 或 full。 |
binlog_expire_logs_seconds |
binlog_expire_logs_seconds 对应已弃用的系统变量 expire_logs_days。这是自动删除 binlog 文件的秒数。默认值为 2592000,即 30 天。请根据环境需求设置该值。有关更多信息,请参阅 MySQL 清除 Debezium 使用的 binlog 文件。 |
表 15. MySQL binlog 配置属性说明
启用 GTID
全局事务标识符(Global Transaction Identifiers,GTID)可唯一标识集群中服务器上发生的事务。虽然 Debezium MySQL 连接器并不强制要求使用 GTID,但使用 GTID 可以简化复制,并让你更容易确认主服务器和副本服务器是否一致。
GTID 在 MySQL 5.6.5 及更高版本中可用。更多详细信息请参阅 MySQL 文档。
前提条件
- 一个 MySQL 服务器。
- SQL 命令的基础知识。
- 能够访问 MySQL 配置文件。
操作步骤
启用
gtid_mode:mysql> gtid_mode=ON启用
enforce_gtid_consistency:mysql> enforce_gtid_consistency=ON确认更改:
mysql> show global variables like '%GTID%';
结果
+--------------------------+-------+
| Variable_name | Value |
+--------------------------+-------+
| enforce_gtid_consistency | ON |
| gtid_mode | ON |
+--------------------------+-------+表 16. GTID 选项说明
| 选项 | 说明 |
|---|---|
gtid_mode | 布尔值,用于指定 MySQL 服务器是否启用 GTID 模式。 |
- ON = 启用 | |
- OFF = 禁用 | |
enforce_gtid_consistency | 布尔值,用于指定服务器是否强制 GTID 一致性,即只允许执行能够以事务安全方式记录日志的语句。使用 GTID 时必须启用该选项。 |
- ON = 启用 | |
- OFF = 禁用 |
配置会话超时
对大型数据库进行初始一致性快照时,已建立的连接可能会在读取表的过程中超时。您可以通过在 MySQL 配置文件中配置 interactive_timeout 和 wait_timeout 来防止这种情况发生。
前提条件
- 一个 MySQL 服务器。
- 基本的 SQL 命令知识。
- 访问 MySQL 配置文件的权限。
操作步骤
配置
interactive_timeout:mysql> interactive_timeout=<持续时间(秒)>配置
wait_timeout:mysql> wait_timeout=<duration-in-seconds>
表 17. MySQL 会话超时选项说明
选项 说明
interactive_timeout
服务器在关闭交互式连接之前,等待该连接上出现活动的秒数。 +1
更多信息请参阅 MySQL 文档。
wait_timeout
服务器在关闭非交互式连接之前,等待该连接上出现活动的秒数。更多信息请参阅 MySQL 文档。
启用查询日志事件
你可能希望查看每个 binlog 事件对应的原始 SQL 语句。在 MySQL 配置文件中启用 binlog_rows_query_log_events 选项即可实现这一点。
该选项适用于 MySQL 5.6 及更高版本。
前提条件
- 一个 MySQL 服务器。
- SQL 命令的基础知识。
- 访问 MySQL 配置文件的权限。
操作步骤
在 MySQL 中启用
binlog_rows_query_log_events:mysql> binlog_rows_query_log_events=ON
binlog_rows_query_log_events 被设置为一个启用/禁用在 binlog 条目中包含原始 SQL 语句的值。
ON= 启用OFF= 禁用
验证 binlog 行值选项
在数据库中检查 binlog_row_value_options 变量的设置。要使连接器能够消费 UPDATE 事件,必须将该变量设置为 PARTIAL_JSON 以外的值。
前置条件
- 一个 MySQL 服务器。
- 基本的 SQL 命令知识。
- 访问 MySQL 配置文件的权限。
操作步骤
检查当前变量值
mysql> show global variables where variable_name = 'binlog_row_value_options';
结果
+--------------------------+-------+
| Variable_name | Value |
+--------------------------+-------+
| binlog_row_value_options | |
+--------------------------+-------+如果该变量的值被设置为
PARTIAL_JSON,请运行以下命令将其取消设置:mysql> set @@global.binlog_row_value_options="" ;
部署
要部署 Debezium MySQL 连接器,需要安装 Debezium MySQL 连接器归档文件、配置连接器,并通过将配置添加到 Kafka Connect 来启动连接器。
前提条件
- 已安装 Apache Kafka 和 Kafka Connect。
- 已安装 MySQL Server,并且已完成使其与 Debezium 连接器配合使用的设置。
操作步骤
- 下载 Debezium MySQL 连接器插件。
- 将文件解压到 Kafka Connect 环境中。
- 将包含 JAR 文件的目录添加到 Kafka Connect 的
plugin.path中。 - 配置连接器并将配置添加到 Kafka Connect 集群中。
- 重启 Kafka Connect 进程,以加载新的 JAR 文件。
如果你使用的是不可变容器,可以查看 Debezium 的容器镜像,其中提供了已预装 MySQL 连接器并可直接运行的 Apache Kafka MySQL 和 Kafka Connect 镜像。
你从 quay.io 获取的 Debezium 容器镜像未经过严格测试或安全分析,仅供测试和评估之用。这些镜像不适用于生产环境。为降低生产部署中的风险,请仅部署由可信供应商积极维护并已针对潜在漏洞进行彻底测试的容器。 |
|---|
你也可以在 Kubernetes 和 OpenShift 上运行 Debezium。
MySQL 连接器配置示例
以下示例展示了一个连接器实例的配置,该实例从位于 192.168.99.100 的 3306 端口上的 MySQL 数据库服务器捕获数据,我们在逻辑上将其命名为 fullfillment。通常,你需要在 JSON 文件中通过设置连接器可用的配置属性来配置 Debezium MySQL 连接器。
你可以选择仅为数据库中部分模式和表的子集生成事件。此外,你还可以忽略、遮蔽或截断包含敏感数据的列、超过指定大小的列,或下游应用程序不需要支持的列。
{
"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
指定 MySQL 服务器的地址或主机名。
database.port
指定连接器用于访问 MySQL 服务器的端口号。
database.user
指定连接器用于访问数据库的 MySQL 用户账号名称。该账号必须具备足够的权限,详见用户权限说明表格。
database.password
指定 Debezium 用于访问数据库的 MySQL 用户账号密码。
database.server.id
指定连接器的唯一 ID。
topic.prefix
指定 MySQL 服务器或集群的主题前缀。
database.include.list
指定希望连接器从中捕获数据的、位于指定服务器上的数据库。
schema.history.internal.kafka.bootstrap.servers
指定连接器用于向数据库 schema 历史主题写入和恢复 DDL 语句的 Kafka broker。
schema.history.internal.kafka.topic
指定数据库 schema 历史主题的名称。该主题仅供内部使用,不应由使用者(consumer)使用。
include.schema.changes
指定连接器是否应为 DDL 变更生成事件,并将其发送到 fulfillment schema 变更主题,以供使用者使用。
有关 Debezium MySQL 连接器可设置的完整配置属性列表,请参阅 MySQL 连接器配置属性。
你可以通过 POST 命令将此配置发送给正在运行的 Kafka Connect 服务。该服务会记录该配置,并启动一个连接器任务来执行以下操作:
- 连接到 MySQL 数据库。
- 读取处于捕获模式的表的变更数据表。
- 将变更事件记录流式传输到 Kafka 主题。
添加连接器配置
要启动并运行 MySQL 连接器,请配置连接器配置,并将该配置添加到 Kafka Connect 集群。
前置条件
- MySQL 已设置为可与 Debezium 连接器配合使用。
- Debezium MySQL 连接器已安装。
操作步骤
- 创建 MySQL 连接器的配置。
- 使用 Kafka Connect REST API 将该连接器配置添加到你的 Kafka Connect 集群。
结果
连接器启动后,它会对其所配置的 MySQL 数据库执行一致性快照。随后,连接器开始为行级操作生成数据变更事件,并将变更事件记录流式传输到 Kafka 主题。
连接器属性
Debezium MySQL 连接器提供了众多配置属性,你可以利用它们来使连接器的行为符合应用的需求。许多属性都有默认值。
MySQL 连接器配置属性的相关信息组织如下:
数据库模式历史连接器配置属性,用于控制 Debezium 如何处理它从数据库模式历史主题中读取的事件。
必需的 Debezium MySQL 连接器配置属性
以下列表描述了部署 MySQL 连接器时必需的配置属性,除非另有默认值说明。
默认值
long
说明
指定连接器在变更事件中如何表示 BIGINT UNSIGNED 列。
可设置为以下选项之一:
long
使用 Java long 数据类型表示 BIGINT UNSIGNED 列的值。虽然 long 类型不提供最高的精度,但在大多数消费者中实现起来很容易。在大多数环境中,这是首选设置。
precise
使用 java.math.BigDecimal 数据类型表示值。连接器使用 Kafka Connect 的 org.apache.kafka.connect.data.Decimal 数据类型,以编码二进制格式表示值。如果连接器通常处理大于 2^63 的值,请设置此选项。long 数据类型无法表示该量级的值。
默认值
bytes
说明
指定连接器在变更事件中如何表示二进制列(如 blob、binary、varbinary)的值。
可设置以下选项之一:
bytes
以字节数组表示二进制数据。
base64
以 base64 编码的字符串表示二进制数据。
base64-url-safe
以 base64-url-safe 编码的字符串表示二进制数据。
hex
以十六进制(base16)编码的字符串表示二进制数据。
默认值
空字符串
说明
一个可选的、以逗号分隔的正则表达式列表,用于匹配需要从变更事件记录值中排除的列的完全限定名称。源记录中的其他列照常被捕获。列的完全限定名称格式为 databaseName.tableName.columnName。
在匹配列名时,Debezium 会将你指定的正则表达式作为锚定正则表达式来应用。也就是说,指定的表达式会与列的整个名称字符串进行匹配,而不会匹配列名中可能出现的子串。如果在配置中包含此属性,请不要同时设置 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 会将你指定的正则表达式作为锚定正则表达式来应用。也就是说,指定的表达式会与列的整个名称字符串进行匹配,而不会匹配列名中可能出现的子串。在生成的变更事件记录中,指定列的值将被替换为化名。
化名由应用指定的 hashAlgorithm 和 salt 后得到的哈希值构成。所使用的哈希函数能够保证在列值被替换为化名的同时,仍维持引用完整性。支持的哈希函数详见 Java 加密体系结构标准算法名称文档中的 MessageDigest 部分。
在下面的示例中,CzQMA0cB5K 是随机选取的盐值。
column.mask.hash.SHA-256.with.salt.CzQMA0cB5K = inventory.orders.customerName, inventory.shipment.customerName如有必要,假名会自动缩短至与列长度一致。连接器配置可以包含多个属性,用于指定不同的哈希算法和盐值。
根据所使用的 hashAlgorithm、所选的 salt 以及实际数据集的不同,生成的数据集可能无法被完全掩盖。
哈希策略版本 2 可确保在不同位置或系统中进行哈希的值保持一致性。
默认值
无默认值
描述
一个可选的、以逗号分隔的正则表达式列表,用于匹配基于字符的列的完全限定名称。如果希望连接器对一组列的值进行掩盖处理(例如,当这些列包含敏感数据时),请设置此属性。将 length 设置为正整数,可使用属性名称中 length 指定的星号(*)字符数量替换指定列中的数据。将 length 设置为 0(零),可使用空字符串替换指定列中的数据。
列的完全限定名称遵循以下格式:databaseName.tableName.columnName。为匹配列名,Debezium 会将您指定的正则表达式作为锚定正则表达式应用。也就是说,指定的表达式会与列的整个名称字符串进行匹配,而不会匹配列名中可能出现的子字符串。
您可以在单个配置中指定多个具有不同长度的属性。
默认值
无默认值
描述
一个可选的、以逗号分隔的正则表达式列表,用于匹配列的完全限定名称,对于这些列,您希望连接器发出表示列元数据的附加参数。设置此属性后,连接器会向事件记录的架构中添加以下字段:
__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 会将你指定的正则表达式作为锚定正则表达式来应用。也就是说,该表达式会与列的整个名称字符串进行匹配,而不会匹配列名中可能出现的子字符串。
你可以在单个配置中指定多个具有不同长度的属性。
默认值
30000(30 秒)
说明
一个正整数值,用于指定连接器在连接请求超时之前,等待与 MySQL 数据库服务器建立连接的最长时间(以毫秒为单位)。
默认值
无默认值
说明
连接器的 Java 类名称。对于 MySQL 连接器,始终指定 io.debezium.connector.mysql.MySqlConnector。
默认值
空字符串
说明
一个可选的、以逗号分隔的正则表达式列表,用于匹配你不希望连接器捕获其变更的数据库名称。连接器会捕获名称未列在 database.exclude.list 中的任何数据库的变更。
为了匹配数据库名称,Debezium 会将你指定的正则表达式作为锚定正则表达式来应用。也就是说,该表达式会与数据库的整个名称字符串进行匹配,而不会匹配数据库名中可能出现的子字符串。
如果你在配置中包含了此属性,请不要同时设置 database.include.list 属性。
默认值
无默认值
说明
MySQL 数据库服务器的 IP 地址或主机名。
默认值
空字符串
说明
一个可选的、以逗号分隔的正则表达式列表,用于匹配连接器从中捕获变更的数据库名称。对于名称不在 database.include.list 中的任何数据库,连接器不会捕获其变更。默认情况下,连接器会捕获所有数据库的变更。
要匹配数据库的名称,Debezium 会将你指定的正则表达式作为锚定正则表达式来应用。也就是说,指定的表达式会与数据库的整个名称字符串进行匹配,而不会匹配数据库名称中可能出现的子串。
如果在配置中包含此属性,请不要同时设置 database.exclude.list 属性。
默认值
com.mysql.cj.jdbc.Driver
描述
指定连接器使用的驱动程序类的名称。
设置此属性可配置随连接器一同打包的驱动程序之外的其他驱动程序。
默认值
无默认值
描述
连接器用于连接 MySQL 数据库服务器的 MySQL 用户的密码。
默认值
3306
描述
MySQL 数据库服务器的整数端口号。
默认值
jdbc:mysql
描述
指定驱动程序连接字符串用于连接数据库的 JDBC 协议。
默认值
无默认值
描述
此数据库客户端的数字 ID。指定的 ID 必须在 MySQL 集群中当前运行的所有数据库进程中唯一。为了启用从 binlog 读取,连接器会使用此唯一 ID 作为另一个服务器加入 MySQL 数据库集群。
默认值
无默认值
描述
连接器用于连接 MySQL 数据库服务器的 MySQL 用户的名称。
默认值
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
跳过问题事件,且不记录任何内容。
默认值
无默认值
描述
指定应如何调整字段名称,以兼容连接器所使用的消息转换器。
可设置以下选项之一:
none
不做任何调整。
avro
将 Avro 名称中无效的字符替换为下划线字符。
avro_unicode
将下划线字符或 Avro 名称中不可用的字符替换为相应的 unicode 表示,例如 _uxxxx。
下划线字符(_)表示转义序列,类似于 Java 中的反斜杠 |
|---|
有关更多信息,请参阅:Avro 命名。
默认值
无默认值
描述
以逗号分隔的正则表达式列表,用于匹配连接器在 MySQL 服务器上查找 binlog 位置时所使用的 GTID 集合中的源域 ID。设置此属性后,连接器只会使用其源 UUID 不匹配任何指定 exclude 模式的 GTID 范围。
为了匹配 GTID 的值,Debezium 会将你指定的正则表达式作为锚定正则表达式应用。也就是说,指定的表达式将与 GTID 的域标识符进行匹配。
如果设置了此属性,请勿同时设置 gtid.source.includes 属性。
默认值
无默认值
描述
以逗号分隔的正则表达式列表,用于匹配连接器在 MySQL 服务器上查找 binlog 位置时所使用的 GTID 集合中的源域 ID。设置此属性后,连接器只会使用其源 UUID 匹配指定 include 模式之一的 GTID 范围。
为了匹配 GTID 的值,Debezium 会将你指定的正则表达式作为锚定正则表达式应用。也就是说,指定的表达式将与 GTID 的域标识符进行匹配。
如果设置了此属性,请勿同时设置 gtid.source.excludes 属性。
默认值
false
描述
布尔值,指定连接器在恢复期间是否忽略已存储偏移量中的 GTID 集合。设置为 true 时,连接器将从已存储的 binlog 文件和位置开始,而不是尝试基于 GTID 的定位。
流式传输恢复后,连接器会正常捕获 GTID,并在后续的 offset 中存储刷新后的 GTID 状态。
默认值
false
说明
一个布尔值,用于指定连接器发出的变更事件中是否包含生成该变更的 SQL 查询。
将此属性设置为 true 可能会暴露你通过其他设置明确排除或屏蔽的表或字段的信息。 |
|---|
要启用此属性,必须将数据库属性 binlog_annotate_row_events 设置为 ON。
设置此属性对快照过程生成的事件没有影响。快照事件不包含原始 SQL 查询。
有关配置数据库以为每个日志事件返回原始 SQL 语句的更多信息,请参阅启用查询日志事件。
默认值
true
说明
一个布尔值,用于指定连接器是否将数据库架构的变更发布到与主题前缀同名的 Kafka 主题中。连接器会记录每一次架构变更,其键包含数据库名,值则是描述该架构更新的 JSON 结构。这种记录架构变更的机制与连接器内部对数据库架构历史变更的记录相互独立。
inconsistent.schema.handling.mode
默认值
fail
说明
指定当二进制日志事件引用了内部架构表示中不存在的表时,连接器应如何响应。也就是说,内部表示与数据库不一致。
可设置为以下选项之一:
fail
连接器抛出异常,报告有问题的事件及其二进制日志偏移量,随后停止运行。
warn
连接器将有问题的事件及其二进制日志偏移量记录到日志中,然后跳过该事件。
skip
连接器跳过有问题的事件,并且不会在日志中记录该事件。
默认值
无默认值
说明
一个表达式列表,用于指定连接器为发布到指定表对应 Kafka 主题的变更事件记录构造自定义消息键时所使用的列。默认情况下,Debezium 使用表的主键列作为其发出记录的消息键。若要替代默认设置,或为缺少主键的表指定键,你可以基于一个或多个列配置自定义消息键。
要为某个表设置自定义消息键,请列出该表的表名,其后跟用作消息键的列名。每个列表项采用以下格式:
<完整限定表名>:<键列>,<键列>
要基于多个列名设置表键,请在列名之间插入逗号。
每个完整限定表名都是以下格式的正则表达式:
<数据库名>.<表名>
该属性可以包含多个表的条目。在列表中使用分号分隔各个表的条目。
以下示例为 inventory.customers 和 purchase.orders 表设置消息键:
inventory.customers:pk1,pk2;(.*).purchaseorders:pk3,pk4
对于 inventory.customer 表,将 pk1 和 pk2 列指定为消息键。对于任意数据库中的 purchaseorders 表,将 pk3 和 pk4 列指定为消息键。
用于创建自定义消息键的列数没有限制。但最好只使用指定唯一键所需的最少列数。
默认值
无默认值
说明
连接器的唯一名称。如果尝试用相同的名称注册多个连接器,注册将失败。所有 Kafka Connect 连接器都需要此属性。
默认值
无默认值
说明
指定连接器如何调整模式名称,以兼容连接器所使用的消息转换器。
请设置以下选项之一:
none
不进行调整。
avro
将 Avro 名称中无效的字符替换为下划线字符。
avro_unicode
将下划线字符或不能用于 Avro 名称的字符替换为相应的 Unicode 转义序列,例如 _uxxxx。
_ 是转义序列,类似于 Java 中的反斜杠 |
|---|
默认值
false
说明
指定当连接器未检测到所包含列发生变化时,是否为该记录发出消息。如果列列在 column.include.list 中,或未列在 column.exclude.list 中,则该列被视为已包含。将该值设置为 true 可阻止连接器在所包含列没有变化时捕获记录。
默认值
空字符串
说明
一个可选的、以逗号分隔的正则表达式列表,用于匹配你不希望连接器捕获变更的表的完整限定表标识符。连接器会捕获 table.exclude.list 中未包含的任何表的变更。每个标识符的形式为 数据库名.表名。
要匹配列名,Debezium 会将你指定的正则表达式作为锚定正则表达式来应用。也就是说,指定的表达式会与表的整个名称字符串进行匹配,而不会匹配表名中可能存在的子字符串。
如果设置了此属性,请不要同时设置 table.include.list 属性。
默认值
空字符串
描述
一个可选的、以逗号分隔的正则表达式列表,用于匹配你希望捕获其更改的表的完全限定标识符。连接器不会捕获未包含在 table.include.list 中的任何表的更改。每个标识符的形式为 databaseName.tableName。默认情况下,连接器会捕获其配置为从中捕获更改的每个数据库中所有非系统表的更改。
要匹配表名,Debezium 会将你指定的正则表达式作为锚定正则表达式来应用。也就是说,指定的表达式会与表的整个名称字符串进行匹配,而不会匹配表名中可能存在的子字符串。
如果设置了此属性,请不要同时设置 table.exclude.list 属性。
默认值
1
描述
为此连接器创建的最大任务数。由于 MySQL 连接器始终使用单个任务,更改默认值不会产生任何影响。
默认值
adaptive_time_microseconds
描述
指定连接器用于表示时间、日期和时间戳值的精度类型。
请设置以下选项之一:
adaptive_time_microseconds
连接器会根据数据库列的类型,使用毫秒、微秒或纳秒精度值,完全按照数据库中的原样捕获日期、datetime 和 timestamp 值;但 TIME 类型的字段除外,它们始终以微秒精度捕获。
adaptive
(已弃用)连接器会根据列的数据类型,使用毫秒、微秒或纳秒精度值,完全按照数据库中的原样捕获时间和时间戳值。
connect
连接器始终使用 Kafka Connect 内置的 Time、Date 和 Timestamp 表示形式来表示时间和时间戳值,无论数据库列的精度如何,这些表示形式都使用毫秒精度。
isostring
连接器使用 io.debezium.time.IsoDate、io.debezium.time.IsoTime 和 io.debezium.time.IsoTimestamp 语义类型,将时间、日期和 datetime 值表示为 UTC 时区下的 ISO-8601 格式字符串。
默认值
true
描述
指定 delete 事件之后是否跟随一个墓碑事件。源记录被删除后,连接器可以发出墓碑事件(默认行为),以便在为该主题启用日志压缩时,让 Kafka 能够完全删除与被删除行的键相关的所有事件。
可设置以下选项之一:
true
连接器通过发出一个 delete 事件及其后的墓碑事件来表示删除操作。
false
连接器仅发出 delete 事件。
默认值
无默认值
描述
指定 Debezium 从其捕获变更的 MySQL 数据库服务器或集群命名空间的字符串。由于主题前缀用于命名接收该连接器发出的所有事件的 Kafka 主题,因此该前缀在所有连接器之间必须唯一。值只能包含字母数字字符、连字符、点和下划线。
| 设置此属性后,请勿更改其值。如果更改了该值,连接器重启后将不再继续向原有主题发送事件,而是基于新值的名称向其他主题发送后续事件。连接器还将无法恢复其数据库 schema 历史主题。 |
|---|
Debezium MySQL 连接器高级配置属性
以下列表介绍了 MySQL 连接器的高级配置属性,这些属性可用于针对特定环境微调行为,包括 schema 历史处理、信号通道、事件过滤和自定义转换。这些属性的默认值通常无需更改,因此你不必在连接器配置中指定它们。
默认值
0
描述
binlog 读取器使用的前瞻缓冲区的大小。默认设置 0 表示禁用缓冲。
在某些情况下,MySQL binlog 中可能包含由 ROLLBACK 语句回滚的未提交数据。典型示例包括在同一个事务中使用保存点,或将临时表和普通表的变更混在一起。
当检测到事务开始时,Debezium 会向前推进 binlog 位置,并查找 COMMIT 或 ROLLBACK,从而确定是否要流式传输该事务中的变更。binlog 缓冲区的大小定义了 Debezium 在查找事务边界时能够缓冲的最大变更数量。如果事务的大小超过缓冲区,Debezium 就必须回退并重新读取那些未能放入缓冲区的事件,然后再继续流式传输。
| 此功能处于孵化阶段,欢迎提供反馈。该功能预计尚未完全成熟。 |
|---|
默认值
0
说明
在服务器超时之前,等待从 binlog 连接完成读取的秒数。值为 0 表示连接器使用 MySQL 服务器的默认值。
在高延迟网络环境中,服务器的 net_read_timeout 默认值可能过低,导致服务器以 EOFException 过早关闭 binlog 流式传输连接。增加此值可防止 binlog 流式传输期间出现意外断开。
该值是通过 SET net_read_timeout=<value> 在 binlog 流式传输连接的会话级别上设置的,不会影响服务器的全局设置。 |
|---|
默认值
0
说明
在服务器超时之前,等待向 binlog 连接完成写入的秒数。值为 0 表示连接器使用 MySQL 服务器的默认值。
当服务器通过 binlog 传输大量数据时,服务器的 net_write_timeout 默认值可能过低,导致服务器以 EOFException 关闭连接。增加此值可防止流式传输大型事务期间出现意外断开。
该值是通过 SET net_write_timeout=<value> 在 binlog 流式传输连接的会话级别上设置的,不会影响服务器的全局设置。 |
|---|
默认值
true
说明
一个布尔值,指定是否应使用单独的线程来确保与 MySQL 服务器或集群的连接保持存活。
默认值
无默认值
说明
用于列出连接器可使用的自定义转换器实例的符号名称,这些名称以逗号分隔。例如 boolean。
此属性是连接器使用自定义转换器的必要条件。
对于为连接器配置的每个转换器,你都必须添加一个 .type 属性,用于指定实现转换器接口的类的完全限定名。.type 属性使用以下格式:
<converterSymbolicName>.type
例如
boolean.type: io.debezium.connector.binlog.converters.TinyIntOneToBooleanConverter若要进一步控制已配置转换器的行为,可以添加一个或多个配置参数,以便向转换器传递值。要将这些额外的配置参数与某个转换器关联,需在参数名前加上该转换器的符号名作为前缀。
例如,要定义一个 selector 参数,用于指定 boolean 转换器所处理的列的子集,可添加以下属性:
boolean.selector=db1.table1.*, db1.table2.column1默认值
无默认值
描述
定义用于自定义 MBean 对象名称的标签,通过添加元数据来提供上下文信息。请指定一个以逗号分隔的键值对列表。每个键表示 MBean 对象名称的一个标签,相应的值表示该键的取值,例如 k1=v1,k2=v2。
连接器会将指定的标签附加到基础 MBean 对象名称之后。标签有助于你对指标数据进行组织和分类。你可以定义标签来标识特定的应用实例、环境、区域、版本等。
有关更多信息,请参阅自定义 MBean 名称。
默认值
.*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 的默认屏蔽模式;该指定值不会扩展或增强默认模式。 |
|---|
默认值
无默认值
描述
以分号分隔的 SQL 语句列表,在建立到数据库的 JDBC 连接(不是读取事务日志的那个连接)时执行。若要在 SQL 语句中将分号指定为字符而非分隔符,请使用两个分号(;;)。
连接器可能会自行决定建立 JDBC 连接,因此此属性仅用于配置会话参数,不用于执行 DML 语句。
默认值
600000(10 分钟)
描述
指定连接器等待查询完成的时间,单位为毫秒。
将该值设置为 0(零)即可移除超时限制。
默认值
无默认值
说明
一个可选设置,用于指定密钥库文件的位置。密钥库文件可用于客户端与 MySQL 服务器之间的双向认证。
database.ssl.keystore.password
默认值
无默认值
说明
密钥库文件的密码。仅在配置了 database.ssl.keystore 时才指定密码。
默认值
preferred
说明
指定连接器是否使用加密连接。
可用的设置如下:
disabled
指定使用未加密的连接。
preferred
如果服务器支持安全连接,连接器将建立加密连接;如果服务器不支持安全连接,连接器将回退为使用未加密的连接。
required
连接器建立加密连接。如果无法建立加密连接,连接器将失败。
verify_ca
连接器的行为与设置 required 选项时相同,但它还会根据已配置的证书颁发机构(CA)证书验证服务器的 TLS 证书。如果服务器的 TLS 证书与任何有效的 CA 证书都不匹配,连接器将失败。
verify_identity
连接器的行为与设置 verify_ca 选项时相同,但它还会验证服务器证书是否与远程连接的主机相匹配。
默认值
无默认值
说明
用于验证服务器证书的信任库文件的位置。
database.ssl.truststore.password
默认值
无默认值
说明
信任库文件的密码。用于检查信任库的完整性并解锁信任库。
默认值
default
说明
指定用于解析 MySQL DDL 语句的 ANTLR 语法。
可设置以下选项之一:
default
(默认)使用 Oracle MySQL 语法,该语法处于积极维护状态,并支持 MySQL 8.0+ 的特性。此语法基于官方 MySQL 语法规范,推荐在生产环境中使用。
legacy
使用 Positive Technologies(PT)的 MySQL 语法。此选项用于保持与现有部署的向后兼容。PT 语法可能不支持所有 MySQL 8.0+ 的特性。
在大多数情况下,你应该使用 default 语法。只有在迁移过程中遇到现有 DDL 语句的兼容性问题时,才使用 legacy 语法。
默认值
true
说明
布尔值,指示连接器是否将 2 位年份表示转换为 4 位。当需要完全由数据库执行转换时,将该值设置为 false。
MySQL 用户可以插入 2 位或 4 位的年份值。2 位的值会被映射到 1970 - 2069 范围内的年份。默认情况下,由连接器执行该转换。
默认值
-1
说明
指定连接器在操作出现可重试错误(例如连接错误)后如何响应。
可设置为以下选项之一:
-1
不限制。无论之前失败了多少次,连接器始终自动重新启动并重试该操作。
0
禁用。连接器立即失败,且从不重试该操作。需要人工干预才能重新启动连接器。
>0
连接器会自动重新启动,直到达到指定的最大重试次数。在下一次失败后,连接器停止运行,需要人工干预才能重新启动它。
event.converting.failure.handling.mode
默认值
warn
说明
指定当列的数据类型与 Debezium 内部 schema 指定的类型不匹配,导致无法转换表记录时,连接器如何响应。可设置为以下选项之一:
fail
抛出一个异常,报告由于字段的数据类型与 schema 类型不匹配而导致转换失败,并指出可能需要以 schema _only_recovery 模式重新启动连接器,以使转换成功。
warn
连接器将转换失败的列对应的事件字段写入 null 值,并向警告日志写入一条消息。
skip
连接器将转换失败的列对应的事件字段写入 null 值,并向调试日志写入一条消息。
event.processing.failure.handling.mode
默认值
fail
说明
指定连接器如何处理在处理事件时发生的失败,例如遇到损坏的事件。可用设置如下:
fail
连接器抛出一个异常,报告有问题的事件及其位置。随后连接器停止运行。
warn
连接器不抛出异常,而是记录有问题的事件及其位置,然后跳过该事件。
ignore
连接器忽略有问题的事件,且不生成日志条目。
默认值
无默认值
说明
指定连接器发送心跳消息时在源数据库上执行的查询。
例如,以下查询会定期记录源数据库中已执行的 GTID 集合的状态。
INSERT INTO gtid_history_table (select @gtid_executed)
默认值
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 每读取一个数据块,就会向信号数据集合写入一个条目,用于记录打开快照窗口的信号。快照完成后,该条目会被移除。不会为关闭快照窗口的信号创建任何条目。设置此选项可防止信号数据集合快速增长。
默认值
2048
说明
一个正整数值,用于指定此连接器每次迭代中应处理的每批事件的最大大小。
默认值
8192
说明
一个正整数值,用于指定阻塞队列可容纳的最大记录数。当 Debezium 读取从数据库流式传输的事件时,它会先将事件放入阻塞队列,然后再将其写入 Kafka。当连接器接收消息的速度快于将其写入 Kafka 的速度,或者 Kafka 变得不可用时,阻塞队列可为从数据库读取变更事件提供背压。连接器定期记录偏移量时,队列中保存的事件将不被计入。请始终将 max.queue.size 设置为大于 max.batch.size 的值。
默认值
0
说明
一个长整型值,用于指定阻塞队列的最大容量(字节)。默认情况下,未对阻塞队列指定容量限制。要指定队列可占用的字节数,请将此属性设置为一个正的长整型值。如果同时设置了 max.queue.size,则当队列的大小达到任一属性所指定的限制时,向队列的写入操作将被阻塞。例如,如果你设置 max.queue.size=1000,且 max.queue.size.in.bytes=5000,那么当队列包含 1000 条记录后,或者当队列中记录的总容量达到 5000 字节后,向队列的写入操作将被阻塞。
min.row.count.to.stream.results
默认值
1000
说明
在快照期间,连接器会查询它被配置为捕获更改的每一张表。连接器使用每次查询的结果来生成读取事件,该事件包含该表中所有行的数据。此属性决定 MySQL 连接器是将某张表的结果放入内存中(速度快但需要大量内存),还是对结果进行流式处理(速度可能较慢,但适用于非常大的表)。此属性的设置指定了连接器在对结果进行流式处理之前,一张表必须包含的最小行数。
要跳过所有表大小检查,并始终在快照期间对所有结果进行流式处理,请将此属性设置为 0。
默认值
无默认值
描述
为连接器启用的通知通道名称列表。默认情况下,可用的通道包括:
sinklogjmx
你还可以选择实现自定义通知通道。
默认值
500(0.5 秒)
描述
一个正整数值,指定连接器在开始处理一批事件之前,等待新更改事件出现的毫秒数。
默认值
false
描述
决定连接器是否生成带有事务边界的事件,并使用事务元数据丰富更改事件信封。如果希望连接器执行此操作,请指定 true。有关更多信息,请参阅事务元数据。
默认值
false
描述
指定连接器是否将水印写入信号数据集合,以跟踪增量快照的进度。将该值设置为 true,可使对数据库具有只读连接的连接器启用一种增量快照水印策略,该策略无需向信号数据集合写入数据。
默认值
无默认值
描述
用于向连接器发送信号的数据集合的完全限定名称。请使用以下格式指定集合名称:
<databaseName>.<tableName>
集合名称区分大小写。
默认值
无默认值
描述
为连接器启用的信号通道名称列表。默认情况下,可用的通道包括:
sourcekafkafilejmx
你还可以选择实现一个自定义信号通道。
默认值
t
说明
以逗号分隔的列表,指定连接器在流式传输期间要跳过的操作类型。
设置以下某个选项来指定要跳过的操作:
c
插入/创建操作。
u
更新操作。
d
删除操作。
t
截断操作。
none
连接器不跳过任何操作。
默认值
无默认值
说明
以毫秒为单位的时间间隔,连接器启动后在执行快照前需要等待的时长。如果你在集群中启动多个连接器,此属性有助于避免快照中断,因为中断可能会导致连接器重新平衡。
默认值
未设置
说明
默认情况下,连接器在执行快照期间会分批读取表内容。设置此属性可指定每批的最大行数。
| 为了保持连接器的性能,最好保持此属性的未设置默认值。此默认配置使 MySQL 能够逐行将结果集流式传输给 Debezium。相反,如果你设置了此属性,可能会导致性能问题,因为 Debezium 会尝试一次性将整个结果集加载到内存中。 |
|---|
snapshot.include.collection.list
默认值
table.include.list 中指定的所有表。
说明
一个可选的、以逗号分隔的正则表达式列表,用于匹配要包含在快照中的表的完全限定名称(<databaseName>.<tableName>)。指定的条目必须在连接器的 table.include.list 属性中列出。
此属性仅在连接器的 snapshot.mode 属性被设置为 no_data 以外的值时才会生效。此属性不会影响增量快照的行为。
为了匹配表的名称,Debezium 会将你指定的正则表达式作为锚定正则表达式应用。也就是说,指定的表达式会与表的完整名称字符串进行匹配,而不会匹配表名中可能存在的子字符串。
默认值
10000
描述
一个正整数,用于指定执行快照时等待获取表锁的最大时间(以毫秒为单位)。如果连接器无法在此时间间隔内获取表锁,快照将失败。
有关更多信息,请参阅介绍 MySQL 连接器如何执行数据库快照的文档。
默认值
minimal
描述
指定连接器是否持有全局 MySQL 读锁以及持锁的时长,该锁可在连接器执行快照期间阻止任何对数据库的更新。
可使用以下设置:
minimal
连接器仅在快照的初始阶段(即读取数据库模式及其他元数据的阶段)持有全局读锁。在快照的下一阶段,连接器会在从每张表中查询所有行时释放该锁。为了以一致的方式执行 SELECT 操作,连接器使用 REPEATABLE READ 事务。虽然释放全局读锁允许其他 MySQL 客户端更新数据库,但使用 REPEATABLE READ 隔离级别可确保快照的一致性,因为连接器在事务持续期间会继续读取相同的数据。
extended
在整个快照期间阻塞所有写操作。如果客户端提交的并发操作与 MySQL 中的 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。
默认值
1
描述
指定连接器执行初始快照时所使用的线程数。该值必须是大于或等于 1 的正整数。设置为大于 1 的值时,连接器会根据主键范围将每个表划分为若干块,并在所有可用线程之间并发处理这些块,从而执行并行快照。
| 基于块的并行快照是一项孵化中的功能。 |
|---|
| 没有主键的表以及使用了快照选择覆盖的表,会作为一个块进行处理,并回退为单线程快照。 |
|---|
snapshot.max.threads.multiplier
默认值
1
描述
一个全局乘数,用于控制并行快照期间每个表创建的块数量。默认情况下,Debezium 为每个线程创建一个块。该值设置得更高时,连接器会在执行快照时为每个表创建更多、更小的块。更小的块能使各线程的负载更加均衡。例如,在 4 个线程且乘数为 2 的情况下,Debezium 会创建 8 个块,而不是 4 个。
要为特定表覆盖该乘数,请将 snapshot.max.threads.multiplier.<fully_qualified_table_name> 设置为所需的值。
| 与提高乘数相比,增加线程数通常是提高快照吞吐量更有效的方式。 |
|---|
| 基于块的并行快照是一项孵化中的功能。 |
|---|
默认值
false
描述
默认值(false)使连接器能够通过将表划分为若干块,并使用独立线程并发处理每个块来加速初始快照。
启用传统行为后,完成其表快照的线程会保持空闲状态,等待其他线程完成。在强制执行连接超时的环境中,空闲连接可能导致连接器在快照结束后无法正常关闭连接。这样即使快照成功捕获了所有数据,也可能引发异常。如果遇到此问题,请将 snapshot.max.threads 设置为 1 并重试快照。 |
|---|
传统的“每个线程一个表”行为已弃用,并将在未来版本中移除。属性 internal.legacy.snapshot.max.threads 是 legacy.snapshot.max.threads 的已弃用别名,不应在新配置中使用。 |
|---|
默认值
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,则使用此设置来指定在 io.debezium.spi.snapshot.Snapshotter 接口中定义的 name() 方法所提供的自定义实现的名称。连接器重启后,Debezium 会调用指定的自定义实现,以确定是否执行快照。有关更多信息,请参阅自定义快照器 SPI。
默认值
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
连接器在执行初始快照时忽略行数。
默认值
v2
描述
Debezium 事件中 source 代码块的架构版本。Debezium 0.10 为统一所有连接器对外暴露的结构,对 source 代码块的结构进行了一些破坏性变更。
将此选项设置为 v1 可生成早期版本所使用的结构。但是,不建议使用此设置,并计划在未来的 Debezium 版本中移除它。
默认值
true
描述
指定连接器是否为流式处理指标收集高级统计指标,例如分位数。
设置为 true 时,连接器会收集以下统计数据:
- 最小值
- 最大值
- 平均值
- P50(中位数)百分位
- P95 百分位
- P99 百分位
目前仅 MilliSecondsBehindSource 指标支持收集分位数。 |
|---|
这些统计数据是使用概率型数据结构(DDSketch)计算的,该结构以 1% 的相对精度提供近似的分位数值。
+ 设置为 false 时,连接器不会收集分位数,分位数 JMX 指标将返回 null 值。连接器仍会继续收集最小值、最大值和平均值。禁用分位数收集可略微降低内存开销。
+ 有关更多信息,请参阅流式处理指标。
默认值
0
描述
指定连接器在完成快照后延迟启动流式处理过程的时间(以毫秒为单位)。设置延迟间隔有助于防止在快照刚刚完成但流式处理尚未开始的瞬间发生故障,从而导致连接器重新启动快照。设置的延迟值应高于为 Kafka Connect 工作进程设置的 offset.flush.interval.ms 属性的值。
默认值
true
描述
一个布尔值,指定是否应忽略内置系统表。此设置与表的包含列表和排除列表无关。默认情况下,系统表中值的变更不会被捕获,Debezium 也不会为系统表的变更生成事件。
默认值
10000
说明
指定在有界并发哈希映射中可以在内存中存储的主题名称数量。连接器使用该缓存来帮助确定与某个数据集合对应的主题名称。
默认值
.(句点)
说明
指定连接器在主题名称的各个组成部分之间插入的分隔符。
默认值
__debezium-heartbeat
说明
指定连接器向其发送心跳消息的主题的名称。主题名称采用以下格式:
topic.heartbeat.prefix.topic.prefix
例如,当您将该属性设置为默认值且主题前缀为 fulfillment 时,主题名称为 __debezium-heartbeat.fulfillment。
如果设置了 topic.heartbeat.name,则忽略此属性。
默认值
空
说明
为连接器发送心跳消息的主题指定一个明确的完整名称,从而覆盖由 topic.heartbeat.prefix 和 topic.prefix 派生出的基于前缀的命名方式。
设置后,所有心跳消息都会被路由到这个确切的主题名称,而不受连接器 topic.prefix 的影响。在运行多个连接器并且希望将心跳事件汇总到单个共享主题中时,此属性非常有用,可以避免创建大量单分区的心跳主题。
例如,将此属性设置为 debezium-heartbeat 后,所有心跳消息都会被路由到名为 debezium-heartbeat 的主题。
如果此属性为空或未设置,连接器将回退到默认行为:topic.heartbeat.prefix.topic.prefix。
默认值
io.debezium.schema.DefaultTopicNamingStrategy
说明
连接器使用的 TopicNamingStrategy 类的名称。指定的策略决定了连接器如何为主题命名,这些主题用于存储数据变更、模式变更、事务、心跳等的事件记录。
默认值
transaction
说明
指定连接器向其发送事务元数据消息的主题的名称。主题名称采用以下模式:
topic.prefix.topic.transaction
例如,如果主题前缀为 fulfillment,则默认主题名称为 fulfillment.transaction。
默认值
false
描述
一个布尔值,用于指定二进制日志客户端的 keepalive 线程是否将 SO_LINGER 套接字选项设置为 0,以便立即关闭陈旧的 TCP 连接。如果连接器在 SSLSocketImpl.close 中出现死锁,请将该值设置为 true。有关更多信息,请参阅 mysql-binlog-connector-java GitHub 仓库中的 Issue 133。
默认值
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 属性。
表 18. 连接器数据库模式历史配置属性
属性 默认值 描述
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 的连接器共享同一个全局驻留池。当多个连接器跟踪结构相似的表时,可最大化去重效果。
传递式 MySQL 连接器配置属性
你可以在连接器配置中设置传递式属性,以自定义 Apache Kafka 生产者和消费者的行为。有关 Kafka 生产者和消费者的完整配置属性范围,请参阅 Kafka 文档。
用于配置生产者和消费者客户端与模式历史主题交互方式的传递式属性
Debezium 依赖 Apache Kafka 生产者将结构变更写入数据库模式历史主题。同样,在连接器启动时,它依赖 Kafka 消费者从数据库模式历史主题读取数据。你通过为以 schema.history.internal.producer.* 和 schema.history.internal.consumer.* 为前缀的一组传递式配置属性赋值,来定义 Kafka 生产者和消费者客户端的配置。这些传递式生产者和消费者数据库模式历史属性控制一系列行为,例如这些客户端如何与 Kafka broker 建立安全连接,如下例所示:
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=test1234Debezium 在将属性传递给 Kafka 客户端之前,会先去除属性名称中的前缀。
有关 Kafka producer 配置属性和 Kafka consumer 配置属性的更多信息,请参阅 Apache Kafka 文档。
用于配置 MySQL 连接器如何与 Kafka signaling topic 交互的透传属性
Debezium 提供了一组 signal.* 属性,用于控制连接器如何与 Kafka signals topic 交互。
下表描述了 Kafka 的 signal 属性。
表 19. Kafka signals 配置属性
属性 默认值 说明
<topic.prefix>-signal
连接器用于监听临时信号的 Kafka topic 名称。
| 如果自动创建 topic被禁用,你必须手动创建所需的 signaling topic。为保证信号顺序,必须存在 signaling topic。signaling topic 必须只有单个分区。 |
|---|
kafka-signal
Kafka 消费者所使用的组 ID 名称。
signal.kafka.bootstrap.servers
无默认值
连接器用于建立与 Kafka 集群初始连接的主机和端口对列表。每个对都指向 Debezium Kafka Connect 进程所使用的 Kafka 集群。
100
一个整数值,用于指定连接器在轮询信号时等待的最大毫秒数。
用于配置 signaling 通道 Kafka 消费者客户端的透传属性
Debezium 连接器支持对 signals Kafka 消费者进行透传配置。透传的 signals 属性以 signal.consumer.* 为前缀。例如,连接器会将诸如 signal.consumer.security.protocol=SSL 的属性传递给 Kafka 消费者。
Debezium 在将这些属性传递给 Kafka signals 消费者之前,会先去除属性名称中的前缀。
用于配置 MySQL 连接器 sink 通知通道的透传属性
下表描述了可用于配置 Debezium sink notification 通道的属性。
| 属性 | 默认值 | 说明 |
|---|---|---|
notification.sink.topic.name | 无默认值 | 接收来自 Debezium 的通知的主题名称。当你将 notification.enabled.channels 属性配置为将 sink 作为启用的通知通道之一时,此属性为必填项。 |
表 20. Sink 通知配置属性
Debezium 连接器透传数据库驱动程序配置属性
Debezium 连接器支持对数据库驱动程序进行透传配置。透传的数据库属性以 driver.* 前缀开头。例如,连接器会将 driver.foobar=false 这样的属性传递给 JDBC URL。
Debezium 在将属性传递给数据库驱动程序之前,会去除这些属性的前缀。
监控
除了 Kafka 和 Kafka Connect 提供的内建 JMX 指标支持之外,Debezium MySQL 连接器还提供三类额外的指标。
- 快照指标提供连接器在执行快照期间的运行信息。
- 流式指标提供连接器读取 binlog 时的运行信息。
- Schema 历史指标提供连接器 schema 历史记录的状态信息。
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 名称
默认情况下,MySQL 连接器对流式指标使用以下 MBean 名称:
debezium.mysql:type=connector-metrics,context=streaming,server=<topic.prefix>如果你将 custom.metric.tags 的值设置为 database=salesdb-streaming,table=inventory,Debezium 将生成如下自定义 MBean 名称:
debezium.mysql:type=connector-metrics,context=streaming,server=<topic.prefix>,database=salesdb-streaming,table=inventory快照指标
MBean 为 debezium.mysql:type=connector-metrics,context=snapshot,server=<topic.prefix>。
下表列出了可用于监控 Debezium 快照操作的 JMX 指标,包括行计数、表进度、持续时间和队列容量。除非快照操作正在进行中,或者自上次连接器启动以来已发生过快照,否则不会公开快照指标。
| 属性 | 类型 | 描述 |
|---|---|---|
LastEvent | string | 连接器读取到的最后一个快照事件。 |
MilliSecondsSinceLastEvent | long | 连接器读取并处理最近一个事件以来经过的毫秒数。 |
NumberOfErroneousEvents | long | 记录连接器在快照操作期间识别为错误的变更事件数量。每当连接器在初始快照、增量快照或临时快照期间遇到无法处理的事件时,该指标都会递增。事件处理失败可能是因为事件格式错误、与模式不兼容,或在转换过程中出现故障。该指标值在连接器任务的整个生命周期内保持有效。如果快照被中断且连接器任务重启,该指标计数会重置为 0。 |
TotalNumberOfEventsSeen | long | 自上次启动或重置以来,此连接器看到的事件总数。 |
NumberOfEventsFiltered | long | 被连接器上配置的包含/排除列表过滤规则过滤掉的事件数量。 |
CapturedTables | string[] | 连接器捕获的表列表。 |
QueueTotalCapacity | int | 用于在快照器和主 Kafka Connect 循环之间传递事件的队列长度。 |
QueueRemainingCapacity | int | 用于在快照器和主 Kafka Connect 循环之间传递事件的队列的剩余容量。 |
TotalTableCount | int | 被纳入快照的表的总数。 |
RemainingTableCount | int | 快照尚待复制的表数量。 |
SnapshotRunning | boolean | 快照是否已启动。 |
SnapshotPaused | boolean | 快照是否已暂停。 |
SnapshotAborted | boolean | 快照是否已中止。 |
SnapshotCompleted | boolean | 快照是否已完成。 |
SnapshotSkipped | boolean | 快照是否已跳过。 |
SnapshotDurationInSeconds | long | 快照迄今为止所花费的总秒数,即使尚未完成也会计入。其中也包括快照处于暂停状态的时间。 |
SnapshotPausedDurationInSeconds | long | 快照被暂停的总秒数。如果快照被多次暂停,暂停时间会累加。 |
RowsScanned | Map<String, Long> | 包含快照中每个表已扫描行数的映射。处理过程中会逐步将各表添加到该映射中。每扫描 10,000 行以及每完成一个表时都会更新。 |
TableChunkCounts | Map<String, Long> | 在使用基于分块的多线程快照时,包含快照中每个表的分块数量的映射。 |
TableChunksCompletedCounts | Map<String, Long> | 在使用基于分块的多线程快照时,包含快照中每个表已完成分块数量的映射。 |
MaxQueueSizeInBytes | long | 队列的最大缓冲字节数。如果将 max.queue.size.in.bytes 设置为正的 long 值,则可使用此指标。 |
CurrentQueueSizeInBytes | long | 队列中记录的当前字节大小。 |
下表列出了连接器执行增量快照时可用的其他 JMX 指标,其中包括可用于跟踪快照进度的分块标识和表边界标识。
| 属性 | 类型 | 描述 |
|---|---|---|
ChunkId | string | 当前快照分块的标识符。 |
ChunkFrom | string | 定义当前分块的主键集合的下界。 |
ChunkTo | string | 定义当前分块的主键集合的上界。 |
TableFrom | string | 当前正在快照的表的主键集合的下界。 |
TableTo | string | 当前正在快照的表的主键集合的上界。 |
Debezium MySQL 连接器还提供了 HoldingGlobalLock 自定义快照指标。该指标是一个布尔值,用于指示连接器当前是否持有全局写锁或表写锁。
流处理指标
Debezium MySQL 连接器提供了三类指标,它们是 Kafka 和 Kafka Connect 内置 JMX 指标支持之外的补充。
- 快照指标提供连接器执行快照期间的运行信息。
- 流处理指标提供连接器读取 binlog 时的运行信息。
- Schema 历史指标提供连接器 schema 历史的状态信息。
Debezium 监控文档提供了如何使用 JMX 暴露这些指标的详细说明。
MBean 为 debezium.mysql:type=connector-metrics,context=streaming,server=<topic.prefix>。
下表列出了可用于监控 Debezium 流式传输操作的 JMX 指标,包括按类型统计的事件计数、与源端相比的滞后、队列容量以及连接状态。
| 属性 | 类型 | 描述 |
|---|---|---|
LastEvent | string | 连接器读取的最后一个流式事件。 |
MilliSecondsSinceLastEvent | long | 自连接器读取并处理最近一个事件以来经过的毫秒数。 |
NumberOfErroneousEvents | long | 记录连接器在流式传输期间识别为错误的变更事件数量。在流式会话的生命周期内,每当连接器遇到无法处理的事件时,都会将此指标值加一。事件可能因格式错误、与模式不兼容,或在转换过程中发生失败而导致处理失败。该指标值在连接器任务的生命周期内持续保留。连接器重启后,该指标计数会重置为 0。 |
TotalNumberOfEventsSeen | long | 自上次启动连接器或重置指标以来,源数据库报告的数据变更事件总数。表示 Debezium 需要处理的数据变更工作负载。 |
TotalNumberOfCreateEventsSeen | long | 自上次启动或重置指标以来,连接器处理的创建事件总数。 |
TotalNumberOfUpdateEventsSeen | long | 自上次启动或重置指标以来,连接器处理的更新事件总数。 |
TotalNumberOfDeleteEventsSeen | long | 自上次启动或重置指标以来,连接器处理的删除事件总数。 |
NumberOfEventsFiltered | long | 被连接器上配置的包含/排除列表过滤规则所过滤掉的事件数量。 |
NumberOfUnchangedEventsSkipped | long | 自上次启动连接器或重置指标以来,因被监控列未发生变化而跳过的更新事件数量。如果 skip.messages.without.change 为 false,则默认值为 -1。 |
CapturedTables | string[] | 连接器所捕获的表的列表。 |
QueueTotalCapacity | int | 用于在流式传输器与主 Kafka Connect 循环之间传递事件的队列长度。 |
QueueRemainingCapacity | int | 用于在流式传输器与主 Kafka Connect 循环之间传递事件的队列的剩余可用容量。 |
Connected | boolean | 指示连接器当前是否已连接到数据库服务器的标志。 |
MilliSecondsBehindSource | long | 最后一个变更事件的时间戳与连接器处理该事件之间相差的毫秒数。该值会包含运行数据库服务器和连接器的机器之间时钟的任何差异。 |
MilliSecondsBehindSourceMinValue | long | 连接器运行期间观测到的相对于源的最小滞后毫秒数。 |
MilliSecondsBehindSourceMaxValue | long | 连接器运行期间观测到的相对于源的最大滞后毫秒数。 |
MilliSecondsBehindSourceAverageValue | double | 连接器运行期间所有观测值计算出的相对于源的平均滞后毫秒数。 |
MilliSecondsBehindSourceP50 | double | 相对于源的滞后毫秒数的第 50 百分位数(中位数)。由于受异常值影响较小,该指标比平均值更能稳健地衡量典型滞后。当 statistics.metrics.enabled 设置为 true(默认值)时可用。 |
MilliSecondsBehindSourceP95 | double | 相对于源的滞后毫秒数的第 95 百分位数。该指标表示 95% 的滞后测量值低于此值,有助于识别尾部延迟并设定 SLA 阈值。当 statistics.metrics.enabled 设置为 true(默认值)时可用。 |
MilliSecondsBehindSourceP99 | double | 相对于源的滞后毫秒数的第 99 百分位数。该指标表示 99% 的滞后测量值低于此值,有助于了解最坏情况下的性能表现。当 statistics.metrics.enabled 设置为 true(默认值)时可用。 |
NumberOfCommittedTransactions | long | 已提交的已处理事务数量。 |
SourceEventPosition | Map<String, String> | 最后接收到的事件的位置坐标。 |
LastTransactionId | string | 最后处理的事务的事务标识符。 |
MaxQueueSizeInBytes | long | 队列的最大缓冲区大小(字节)。当 max.queue.size.in.bytes 设置为正整数值时,此指标可用。 |
CurrentQueueSizeInBytes | long | 队列中记录的当前大小(字节)。 |
Debezium MySQL 连接器还提供以下额外的流处理指标:
| 属性 | 类型 | 描述 |
|---|---|---|
BinlogFilename | string | 连接器最近读取的 binlog 文件名。 |
BinlogPosition | long | 连接器最近读取的 binlog 内位置(以字节为单位)。 |
IsGtidModeEnabled | boolean | 表示连接器当前是否正在跟踪 MySQL 服务器的 GTID 的标志。 |
GtidSet | string | 连接器读取 binlog 时最近处理的 GTID 集合的字符串表示形式。 |
NumberOfSkippedEvents | long | MySQL 连接器跳过的事件数量。通常,事件被跳过是因为来自 MySQL binlog 的事件格式错误或无法解析。 |
NumberOfDisconnects | long | MySQL 连接器断开连接的次数。 |
NumberOfRolledBackTransactions | long | 已处理但被回滚且未被流式传输的事务数量。 |
NumberOfNotWellFormedTransactions | long | 未遵循 BEGIN + COMMIT/ROLLBACK 预期协议的事务数量。在正常情况下,此值应为 0。 |
NumberOfLargeTransactions | long | 未能放入前瞻缓冲区(look-ahead buffer)的事务数量。为了获得最佳性能,此值应显著小于 NumberOfCommittedTransactions 和 NumberOfRolledBackTransactions。 |
表 21. 其他 MySQL 流式指标说明
Schema 历史记录指标
其 MBean 为 debezium.mysql:type=connector-metrics,context=schema-history,server=<topic.prefix>。
下表列出了可用于监控连接器 Schema 历史记录过程的 JMX 指标,包括恢复状态、已应用的 Schema 变更数量以及最近变更的时间戳。
| 属性 | 类型 | 描述 |
|---|---|---|
Status | string | 描述数据库 schema 历史状态的取值之一:STOPPED(已停止)、RECOVERING(正在从存储中恢复历史记录)、RUNNING(正在运行)。 |
RecoveryStartTime | long | 恢复开始的时间,以纪元秒(epoch seconds)表示。 |
ChangesRecovered | long | 恢复阶段读取到的变更数量。 |
ChangesApplied | long | 恢复和运行期间应用的 schema 变更总数。 |
MilliSecondsSinceLastRecoveredChange | long | 自从上次从历史存储中恢复变更以来经过的毫秒数。 |
MilliSecondsSinceLastAppliedChange | long | 自上次应用变更以来经过的毫秒数。 |
LastRecoveredChange | string | 从历史存储中恢复的最后一个变更的字符串表示。 |
LastAppliedChange | string | 最后应用的变更的字符串表示。 |
出现问题时的行为
Debezium 是一个分布式系统,用于捕获多个上游数据库中的所有更改;它绝不会遗漏或丢失任何事件。当系统正常运行或被妥善管理时,Debezium 会为每个更改事件记录提供 精确一次(exactly once)的传递。
如果确实发生故障,系统不会丢失任何事件。但是,在 Debezium 从故障中恢复的过程中,可能会重复投递某些更改事件。在这些异常情况下,Debezium 与 Kafka 一样,为更改事件提供 至少一次(at least once)的传递。
本节的其余部分将描述 Debezium 如何处理各种类型的故障和问题。
配置与启动错误
在以下情况下,连接器在尝试启动时会失败,在日志中报告错误或异常,并停止运行:
- 连接器的配置无效。
- 连接器无法使用指定的连接参数成功连接到 MySQL 服务器。
- 连接器尝试从 binlog 中的某个位置重新启动,而 MySQL 已不再保留该位置的历史记录。
在这些情况下,错误消息中会包含有关问题的详细信息,并可能给出建议的解决方法。修正配置或解决 MySQL 的问题后,请重新启动连接器。
MySQL 变得不可用
如果 MySQL 服务器变得不可用,Debezium MySQL 连接器会因错误而失败,连接器随之停止。当服务器再次可用时,请重新启动连接器。
不过,如果你连接的是高可用 MySQL 集群,则可以立即重新启动连接器。它会连接到集群中的另一台 MySQL 服务器,找到该服务器 binlog 中代表最后一个事务的位置,并从该特定位置开始读取新服务器的 binlog。
Kafka Connect 正常停止
当 Kafka Connect 正常停止时,会有一段短暂的延迟,用于在新的 Kafka Connect 进程上停止并重新启动 Debezium MySQL 连接器任务。
Kafka Connect 进程崩溃
如果 Kafka Connect 崩溃,进程将停止,且任何 Debezium MySQL 连接器任务都会在未记录其最近已处理偏移量的情况下终止。在分布式模式下,Kafka Connect 会在其他进程上重新启动连接器任务。但是,MySQL 连接器会从先前进程记录的最后一个偏移量处继续。因此,替换任务可能会重新生成崩溃前已处理过的某些事件,从而产生重复事件。
每条更改事件消息都包含特定于来源的信息,你可以利用这些信息来识别重复事件,例如:
- 事件来源
- MySQL 服务器的事件时间
- binlog 文件名和位置 +1
- GTID(如果使用了 GTID)
Kafka 变得不可用
Kafka Connect 框架通过 Kafka producer API 将 Debezium 更改事件记录到 Kafka 中。如果 Kafka broker 变得不可用,Debezium MySQL 连接器会暂停,直到连接重新建立,然后连接器会从中断处继续。
MySQL 清除 binlog 文件
如果 Debezium MySQL 连接器停止的时间过长,MySQL 服务器会清除较旧的 binlog 文件,连接器的最后位置可能因此丢失。重新启动连接器时,MySQL 服务器中已没有该起始点,连接器将再次执行初始快照。如果快照功能被禁用,连接器会以错误失败。
关于 MySQL 连接器如何执行初始快照的详细信息,请参阅快照。
评论
登录后参与评论
KnowForge