Informix
Debezium Informix 连接器
Debezium Informix 连接器可以捕获 Informix 数据库中表的行级变更。有关与此连接器兼容的 Informix 数据库版本信息,请参阅 Debezium 发布版本概览。
该连接器深受 IBM Db2 的 Debezium 实现启发,但使用 Informix Change Streams API for Java 来捕获事务数据。变更数据捕获(Change Data Capture)API 从启用了完整行日志记录的数据库中捕获数据,并从当前逻辑日志中捕获事务。该 API 按顺序处理所有事务。
当 Debezium Informix 连接器首次连接到 Informix 数据库时,连接器会为配置为捕获变更的表读取一份一致的快照。默认情况下,连接器会捕获所有非系统表的变更。要自定义快照行为,您可以设置配置属性,以指定要包含在快照中或从快照中排除的表。
快照完成后,连接器开始为提交到处于捕获模式(capture mode)的表中的更新发出变更事件。默认情况下,某个特定表的变更事件会发送到与该表同名的 Kafka 主题。应用和服务随后可以从这些主题中消费变更事件记录。
| 该连接器需要使用 Informix Change Streams API for Java,它作为 Informix JDBC 安装的一部分打包,并与最新的 JDBC 驱动程序一起发布在 Maven Central 上。 |
|---|
Informix 连接器已在 Linux 版 Informix 上通过测试。预计该连接器同样可以在其他平台(如 Windows)上工作,如果您能确认这一点,我们非常希望收到您的反馈。
概述
Debezium Informix 连接器基于 Informix Change Data Capture API,该 API 可在 Informix 中启用变更数据捕获。
数据库管理员必须为使用变更数据捕获 API 准备好数据库和数据库服务器。请参阅为使用变更数据捕获 API 做准备。
将表置于捕获模式后,连接器即可读取每次表更新的变更流记录,并据此生成变更事件。连接器会为每一行的插入、更新和删除操作发出一条变更事件记录。默认情况下,变更事件记录会被发送到与源表同名的 Kafka 主题中。你也可以自定义目标主题的名称。客户端应用读取与所关注的数据库表相对应的 Kafka 主题,并可对每一行级变更事件作出响应。
通常,数据库管理员会在表生命周期的中途将表置于捕获模式。这意味着连接器并不掌握该表所发生全部变更的完整历史。因此,当 Informix 连接器首次连接到某个特定的 Informix 数据库时,它会先对处于捕获模式的每张表执行一次一致性快照。连接器完成快照后,会从快照生成的时间点开始流式传输变更事件。这样,连接器便以处于捕获模式的表的一致性视图为起点,并且不会丢失在执行快照期间发生的任何变更。
Debezium 连接器具有故障容忍能力。当连接器读取并生成变更事件时,它会记录变更流记录的日志序列号(LSN)。LSN 是变更事件在数据库日志中的位置。如果连接器因任何原因停止运行——包括通信故障、网络问题或崩溃——重启后它会从上次中断的位置继续读取变更流。这一行为同样适用于快照。也就是说,如果连接器停止时快照尚未完成,重启后连接器会开始执行一次新的快照。
连接器的工作原理
要优化配置并运行 Debezium Informix 连接器,了解该连接器如何执行快照、流式传输变更事件、确定 Kafka 主题名称以及处理架构变更会很有帮助。
快照
Informix 的复制功能并非为存储数据库变更的完整历史而设计。因此,Debezium Informix 连接器无法从日志中检索数据库的全部历史记录。为了使连接器能够为数据库的当前状态建立基线,连接器首次启动时会对处于捕获模式的表执行一次初始一致性快照。对于快照捕获到的每一项变更,连接器都会向被捕获表对应的 Kafka 主题发出一个 read 事件。
Debezium Informix 连接器执行初始快照所采用的默认工作流程
以下工作流程列出了 Debezium 创建快照时所采取的步骤。这些步骤描述的是 snapshot.mode 配置属性为其默认值 initial 时的快照过程。你可以通过更改 snapshot.mode 属性的值来定制连接器创建快照的方式。如果你配置了不同的快照模式,连接器将使用此工作流程的修改版本来完成快照。
建立与数据库的连接。
确定哪些表处于捕获模式并应包含在快照中。默认情况下,连接器会捕获所有非系统表的数据。快照完成后,连接器会继续为指定的表流式传输数据。如果你希望连接器仅从特定表中捕获数据,可以通过设置
table.include.list或table.exclude.list等属性,将连接器配置为仅捕获部分表或表元素的数据。对处于捕获模式的每张表获取锁。该锁可确保在快照完成之前这些表不会发生结构变更。锁的级别由连接器配置属性
snapshot.isolation.mode的值决定。读取服务器事务日志中最高(最新)的 LSN 位置。
捕获所有表或所有指定要捕获的表的结构。连接器会将其结构信息持久化到内部的数据库结构历史主题中。结构历史提供了变更事件发生时生效的结构信息。
默认情况下,连接器会捕获数据库中每张处于捕获模式的表的结构,包括未配置为捕获的表。如果某些表未配置为捕获,初始快照仅捕获其结构,不会捕获任何表数据。
有关为何快照会为你未包含在初始快照中的表持久化结构信息的更多信息,请参阅了解为何初始快照会捕获所有表的结构。
释放第 3 步中获取的所有锁。其他数据库客户端现在可以向之前被锁定的任何表写入数据。
在第 4 步读取的 LSN 位置处,连接器扫描指定要捕获的表。在扫描过程中,连接器完成以下任务:
确认该表在快照开始之前就已创建。如果表是在快照开始之后创建的,连接器会跳过该表。快照完成后,连接器转入流式处理阶段,此时它会为所有在快照开始之后创建的表发出变更事件。
- 为从表中捕获的每一行数据生成一个
read事件。所有read事件都包含相同的 LSN 位置,即在第 4 步中获取的 LSN 位置。 - 将每个
read事件发送到源表对应的 Kafka 主题。 - 如果适用,释放数据表锁。
- 为从表中捕获的每一行数据生成一个
在连接器偏移量中记录快照已成功完成。
最终的初始快照捕获了被捕获表中每一行的当前状态。以该基准状态为基础,连接器会捕获随后发生的所有变更。
快照过程开始后,如果由于连接器故障、重新平衡或其他原因导致该过程中断,连接器重启后该过程将重新开始。
连接器完成初始快照后,会从第 4 步中读取的位置继续进行流式处理,以确保不会遗漏任何更新。
如果连接器再次因任何原因停止,重启后它会从上次中断的位置继续流式处理变更事件。
表 1. snapshot.mode 连接器配置属性的设置
设置说明
always
连接器每次启动时都会执行快照。快照完成后,连接器开始流式传输后续数据库变更的事件记录。
initial
连接器按照执行初始快照的默认工作流中所述执行数据库快照。快照完成后,连接器开始流式传输后续数据库变更的事件记录。
initial_only
连接器执行数据库快照。快照完成后,连接器停止运行,不会流式传输后续数据库变更的事件记录。
schema_only
已弃用,请参阅 no_data。
no_data
连接器捕获所有相关表的结构,执行默认快照工作流中描述的所有步骤,但不会创建 READ 事件来表示连接器启动时的数据集(第 7.b 步)。
recovery
设置此选项可恢复丢失或损坏的数据库架构历史主题。重启后,连接器会运行一次快照,从源表重建该主题。您也可以通过设置该属性来定期清理意外增长的数据库架构历史主题。
| 不要使用此模式在上次连接器关闭后数据库中已提交了架构变更的情况下执行快照。 |
|---|
when_needed
连接器启动后,仅在检测到以下情况之一时才执行快照:
- 无法检测到任何主题偏移量。
- 先前记录的偏移量指定的日志位置在服务器上不可用。
configuration_based
将快照模式设置为 configuration_based,可通过带有前缀 snapshot.mode.configuration.based 的连接器属性集来控制快照行为。
custom
custom 快照模式允许你注入自己实现的 io.debezium.spi.snapshot.Snapshotter 接口。将 snapshot.mode.custom.name 配置属性设置为你的实现中 name() 方法提供的名称。该名称位于 Kafka Connect 集群的类路径上。如果使用 DebeziumEngine,该名称则包含在连接器 JAR 文件中。有关更多信息,请参阅 自定义快照器 SPI。
有关更多信息,请参阅连接器配置属性表中的 snapshot.mode。
理解为何初始快照会捕获所有表的架构历史
连接器执行的初始快照会捕获两类信息:
表数据
关于连接器 table.include.list 属性中所列各表的 INSERT、UPDATE 和 DELETE 操作信息。
架构数据
描述应用于表的结构变更的 DDL 语句。架构数据会持久化到内部架构历史主题,以及连接器的架构变更主题(如果已配置)。
运行初始快照后,你可能会注意到快照捕获了未指定为采集对象的表的架构信息。默认情况下,初始快照被设计为捕获数据库中每张表的架构信息,而不仅仅是被指定为采集对象的表。连接器要求表的架构必须先存在于架构历史主题中,才能采集该表。通过使初始快照能够为不在原始采集集合中的表捕获架构数据,Debezium 让连接器在日后需要时能随时采集这些表的事件数据。如果初始快照未捕获某张表的架构,你必须先将该架构添加到历史主题中,连接器才能采集该表的数据。
在某些情况下,你可能希望限制初始快照中的架构捕获。当你想要缩短完成快照所需的时间,或者当 Debezium 通过一个可访问多个逻辑数据库的用户账号连接到数据库实例、而你只希望连接器捕获特定逻辑数据库中表的变更时,这会很有用。
补充信息
- 捕获初始快照未包含的表中的数据(无架构变更)
- 捕获初始快照未包含的表中的数据(有架构变更)
- 设置
schema.history.internal.store.only.captured.tables.ddl属性,以指定要从哪些表中捕获架构信息。 - 设置
schema.history.internal.store.only.captured.databases.ddl属性,以指定要从哪些逻辑数据库中捕获架构变更。
捕获初始快照未包含的表中的数据(无架构变更)
在某些情况下,你可能希望连接器捕获某个表中的数据,而该表的架构并未被初始快照捕获。根据连接器的配置,初始快照可能仅为数据库中的特定表捕获表架构。如果历史主题中不存在该表的架构,连接器将无法捕获该表,并会报告缺少架构的错误。
你仍然可能从该表中捕获数据,但必须执行额外步骤来添加该表的架构。
前提条件
- 你希望从某个表中捕获数据,而该表的架构在初始快照期间未被连接器捕获。
- 在连接器读取的最早变更表条目与最晚变更表条目对应的 LSN 之间,该表未发生架构变更。有关捕获已发生结构变更的新表中的数据,请参阅捕获初始快照未包含的表中的数据(有架构变更)。
操作步骤
- 停止连接器。
- 删除由
schema.history.internal.kafka.topic属性指定的内部数据库架构历史主题。 - 清除已配置的 Kafka Connect
offset.storage.topic中的位点。有关如何移除位点的更多信息,请参阅 Debezium 社区常见问题。
| 移除偏移量的操作仅应由具备操作 Kafka Connect 内部数据经验的高级用户执行。此操作具有潜在的破坏性,仅应在万不得已时使用。 |
|---|
对连接器配置进行以下更改:
(可选)将
schema.history.internal.captured.tables.ddl的值设置为false。此设置会使快照捕获所有表的结构,并保证连接器今后能够重建所有表的结构历史。
捕获所有表结构的快照需要更长的完成时间。 将你希望连接器捕获的表添加到
table.include.list中。将
snapshot.mode设置为以下值之一:initial当你重启连接器时,它会对数据库执行全量快照,捕获表数据和表结构。如果你选择此选项,建议将
schema.history.internal.captured.tables.ddl属性的值设置为false,以使连接器能够捕获所有表的结构。no_data当你重启连接器时,它会执行仅捕获表结构的快照。与全量数据快照不同,此选项不会捕获任何表数据。如果你想比全量快照更快地重启连接器,请使用此选项。
重启连接器。连接器将完成由
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 内部数据经验的高级用户才应执行删除偏移量的操作。此操作具有潜在破坏性,仅应在别无选择时使用。 | |
|---|---|
| 4. 按以下步骤为连接器配置中的属性设置值: |
- 将
snapshot.mode属性的值设置为no_data。 - 编辑
table.include.list,添加要捕获的表。 - 重启连接器。
- 等待 Debezium 捕获新表和现有表的架构。连接器停止后,任何表中发生的数据更改都不会被捕获。
- 为确保不丢失数据,请发起增量快照。
流程 2:初始快照,随后进行可选的增量快照
在此流程中,连接器对数据库执行完整的初始快照。与任何初始快照一样,在包含大量大表的数据库中,运行初始快照可能是一项耗时的操作。快照完成后,你可以选择触发增量快照,以捕获连接器离线期间发生的任何更改。
停止连接器。
删除由
schema.history.internal.kafka.topic property属性指定的内部数据库架构历史主题。清除已配置的 Kafka Connect
offset.storage.topic中的偏移量。有关如何删除偏移量的更多信息,请参阅 Debezium 社区常见问题。只有具备操作 Kafka Connect 内部数据经验的高级用户才应执行删除偏移量的操作。此操作具有潜在破坏性,仅应在别无选择时使用。 编辑
table.include.list,添加要捕获的表。按以下步骤为连接器配置中的属性设置值:
将
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 主题被删除,必须重建。
- 由于配置错误或其他问题导致数据损坏。
你可以通过发起所谓的临时快照,为之前已捕获快照的表重新运行快照。临时快照需要使用信号表。你通过向 Debezium 信号表发送信号请求来发起临时快照。
当你对现有表发起临时快照时,连接器会将内容追加到该表已存在的主题中。如果之前存在的主题已被删除,只要启用了自动创建主题,Debezium 就能自动创建主题。
即席快照信号会指定要包含在快照中的表。快照可以捕获数据库的全部内容,也可以只捕获数据库中部分表;还可以只捕获表中部分数据。
通过向信号表发送 execute-snapshot 消息来指定要捕获的表。将 execute-snapshot 信号的类型设置为 incremental 或 blocking,并提供要包含在快照中的表名,如下表所述:
表 2. 即席 execute-snapshot 信号记录示例
字段 默认值
type
incremental
指定要运行的快照类型。目前可以请求 incremental(增量)或 blocking(阻塞)快照。
data-collections
不适用
一个数组,其中包含与要包含在快照中的表的完全限定名相匹配的正则表达式。对于 Informix 连接器,使用以下格式指定表的完全限定名:database.schema.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)。
目前,你可以使用以下任一方法发起增量快照:
| Informix 的 Debezium 连接器不支持在增量快照运行期间进行架构更改。 |
|---|
触发增量快照
要发起增量快照,你可以向源数据库上的信号表发送临时快照信号。快照信号以 SQL INSERT 查询的形式提交。
Debezium 检测到信号表中的更改后,会读取该信号并运行所请求的快照操作。
你提交的查询会指定要包含在快照中的表,并且可以选择指定快照操作的类型。Debezium 目前支持 incremental 和 blocking 两种快照类型。
要指定包含在快照中的表,请提供一个 data-collections 数组,其中列出各个表,或者列出用于匹配表的正则表达式,例如:
{"data-collections": ["public.MyFirstTable", "public.MySecondTable"]}
数据集合名称区分大小写。增量快照信号的 data-collections 数组没有默认值。如果 data-collections 数组为空,Debezium 会将空数组解释为无需执行任何操作,因此不会执行快照。
如果要包含在快照中的表名含有点号(.)、空格或其他非字母数字字符,则必须用双引号对表名进行转义。例如,要包含位于 db1 数据库的 public 架构中、名称为 My.Table 的表,请使用以下格式:"db1.public.\"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.myschema.debezium_signal (id, type, data)
values ('ad-hoc-1',
'execute-snapshot',
'{"data-collections": ["db1.schema1.table1", "db1.schema1.table2"],
"type":"incremental",
"additional-conditions":[{"data-collection": "db1.schema1.table1" ,"filter":"color=\'blue\'"}]}');命令中 id、type 和 data 参数的值对应于信号表的字段。以下列表描述了上例中的参数:
INSERT INTO database.schema.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 运行临时增量快照。
使用物理行标识符作为代理键
某些数据库提供物理行标识符,这是一种表示行在磁盘上物理位置的伪列。这些标识符可以显著提升增量快照分块的性能。
物理行标识符在以下场景中特别有用:
具有复合主键的表
当表没有单列代理键,而是使用多个列作为主键时,分块查询会变得复杂且低效。
索引利用效率低下
数据库查询优化器通常无法有效利用增量快照分块查询所产生的条件析取(OR 条件)来使用索引,从而导致全表扫描。
通过使用物理行标识符作为代理键,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.myschema.debezium_signal (id, type, data) VALUES('ad-hoc-1', 'execute-snapshot', '{"data-collections": ["db1.schema1.products"],"type":"incremental", "additional-conditions":[{"data-collection": "db1.schema1.products", "filter": "color=blue"}]}');additional-conditions 参数还支持你传入基于多个列的条件。例如,使用上例中的 products 表,你可以提交一个查询,触发增量快照,且仅包含满足 color=blue 和 quantity>10 的那些条目的数据:
INSERT INTO db1.myschema.debezium_signal (id, type, data) VALUES('ad-hoc-1', 'execute-snapshot', '{"data-collections": ["db1.schema1.products"],"type":"incremental", "additional-conditions":[{"data-collection": "db1.schema1.products", "filter": "color=blue AND quantity>10"}]}');以下示例展示了连接器捕获的增量快照事件的 JSON 内容。
示例 1. 增量快照事件消息
{
"before":null,
"after": {
"pk":"1",
"value":"New data"
},
"source": {
...
"snapshot":"incremental"
},
"op":"r",
"ts_ms":"1620393591654",
"ts_us":"1620393591654547",
"ts_ns":"1620393591654547920",
"transaction":null
}以下列表描述了前面增量快照事件消息示例中的部分字段:
snapshot
指定快照的类型。
op
指定操作的类型。对于快照事件,op 字段的值为 r,因为快照是一次 READ 操作。
使用 Kafka 信号通道触发增量快照
要使用 Kafka 信号通道触发一次即时增量快照,请向已配置的 Kafka 信号主题发送 execute-snapshot 消息。
Kafka 消息的键必须与连接器配置选项 topic.prefix 的值相匹配。
消息的值是一个包含 type 和 data 字段的 JSON 对象。
信号类型为 execute-snapshot,而 data 字段必须包含以下字段:
表 3. 执行快照的数据字段
| 字段 | 默认值 |
|---|---|
type |
incremental |
要执行的快照类型。目前 Debezium 支持 incremental 和 blocking 两种类型。更多详情请参见下一节。 |
|
data-collections |
不适用 |
| 一个由逗号分隔的正则表达式数组,用于匹配要包含在快照中的表的全限定名。 | |
| 请使用与 signal.data.collection 配置选项要求相同的格式指定名称。数据集合名称区分大小写。 | |
additional-conditions |
不适用 |
| 一个可选的附加条件数组,用于指定连接器评估的条件,从而指定要包含在快照中的记录子集。 | |
| 每个附加条件都是一个对象,用于指定即时快照捕获数据时的过滤条件。每个附加条件可以设置以下参数: | |
data-collection |
过滤器所应用的表的全限定名。你可以为每张表应用不同的过滤器。 |
filter |
指定数据库记录中必须存在的列值,快照才会包含该记录,例如 "color='blue'"。 |
为 filter 参数赋的值,与为阻塞快照设置 snapshot.select.statement.overrides 属性时在 SELECT 语句的 WHERE 子句中指定的值类型相同。 |
示例 2. 一条 execute-snapshot Kafka 消息
Key = `test_connector`
Value = `{"type":"execute-snapshot","data": {"data-collections": ["{collection-container}.table1", "{collection-container}.table2"], "type": "INCREMENTAL"}}`使用 additional-conditions 的临时增量快照
Debezium 使用 additional-conditions 字段来选取表中的一部分内容。
通常情况下,当 Debezium 运行快照时,会执行如下 SQL 查询:
SELECT * FROM <tableName> ….
当快照请求中包含 additional-conditions 属性时,该属性的 data-collection 和 filter 参数会被追加到 SQL 查询中,例如:
SELECT * FROM <data-collection> WHERE <filter> ….
例如,假设有一个 products 表,包含 id(主键)、color 和 brand 列,如果希望快照只包含 color='blue' 的内容,那么在请求快照时,可以添加 additional-conditions 属性来过滤内容:
Key = `test_connector`
Value = `{"type":"execute-snapshot","data": {"data-collections": ["db1.schema1.products"], "type": "INCREMENTAL", "additional-conditions": [{"data-collection": "db1.schema1.products" ,"filter":"color='blue'"}]}}`你还可以使用 additional-conditions 属性传入基于多个列的条件。例如,沿用前面示例中的 products 表,如果你希望快照仅包含 products 表中 color='blue' 且 brand='MyBrand' 的内容,可以发送如下请求:
Key = `test_connector`
Value = `{"type":"execute-snapshot","data": {"data-collections": ["db1.schema1.products"], "type": "INCREMENTAL", "additional-conditions": [{"data-collection": "db1.schema1.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.myschema.debezium_signal (id, type, data)
values ('ad-hoc-1',
'stop-snapshot',
'{"data-collections": ["db1.schema1.table1", "db1.schema1.table2"],
"type":"incremental"}');信号命令中 id、type 和 data 参数的值对应于信号表的字段。
以下列表描述了前面信号示例中的各个字段:
database.schema.debezium_signal
指定源数据库上信号表的完全限定名称。
ad-hoc-1
信号 id 参数的值。这个任意字符串提供了一个标签,用于区分不同的信号请求,并将日志消息与信号表中的条目关联起来。Debezium 不会使用这个字符串。
stop-snapshot
type 参数,用于标识该信号所要触发的操作。
data-collections
信号 data 字段中的可选组件,用于指定一个表名数组或用于匹配表名的正则表达式数组,以从快照中移除这些表。该数组列出的正则表达式按照 database.schema.table 格式的完全限定名称来匹配表。
如果在 data 字段中省略此组件,信号将停止当前正在进行的整个增量快照。
incremental
信号 data 字段中的必需组件,用于指定要停止的快照操作类型。目前唯一有效的选项是 incremental。如果未指定 type 值,信号将无法停止增量快照。
使用 Kafka 信号通道停止增量快照
要使用 Kafka 信号通道停止正在进行的增量快照,请向配置的 Kafka 信号主题发送 stop-snapshot 消息。
Kafka 消息的键必须与 topic.prefix 连接器配置选项的值匹配。
消息的值是一个包含 type 和 data 字段的 JSON 对象。
信号类型为 stop-snapshot,且 data 字段必须包含以下字段:
| 字段 | 默认值 | 值说明 |
|---|---|---|
type | incremental | 要执行的快照类型。目前 Debezium 仅支持 incremental 类型。详情请参阅下一节。 |
data-collections | 不适用 | 一个可选的数组,包含以逗号分隔的正则表达式,用于匹配表的全限定名;或者一个表名数组、用于匹配表名的正则表达式,以从快照中移除这些表。指定表名时请使用 database.schema.table 格式。 |
表 4. 执行快照数据字段
以下示例展示了一条典型的 stop-snapshot Kafka 消息:
Key = `test_connector`
Value = `{"type":"stop-snapshot","data": {"data-collections": ["db1.schema1.table1", "db1.schema1.table2"], "type": "INCREMENTAL"}}`自定义快照器 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"}]}可能的重复数据
从发出触发快照的信号,到流式传输停止并开始快照之间,可能存在一定的延迟。由于这种延迟,快照完成后,连接器可能会发出一些与快照所捕获的记录重复的事件记录。
更改流记录
在完成完整快照后,当 Debezium Informix 连接器首次启动时,连接器开始消费处于捕获模式的源表的更改流记录。连接器执行以下操作:
- 从当前 LSN 读取可用的更改记录。
- 按事务 Id 对记录进行分组,并按每条记录的更改 LSN 排序。
- 随着事务的提交处理记录。
- 将开始、提交和更改 LSN 作为偏移量传递给 Kafka Connect。
- 存储连接器传递给 Kafka Connect 的最高提交 LSN 和最低的未提交开始 LSN。
重新启动后,连接器会从上次中断处的偏移量(开始、提交和更改 LSN)继续发出更改事件。在恢复正常活动的过程中,连接器按顺序执行以下步骤:
- 读取在最后存储的最低未提交开始 LSN 与当前 LSN 之间创建的更改记录。
- 按事务 Id 对记录进行分组,并按每个事件的更改 LSN 排序。
- 丢弃已处理过的事务(提交 LSN 低于最后存储的提交 LSN)。
- 丢弃最后一个未完整处理的事务中已处理过的记录(如果有的话)(更改 LSN 低于最后存储的更改 LSN,且提交 LSN 等于最后存储的提交 LSN)。
- 处理任何未完整处理的事务的剩余记录。
- 随着事务的提交继续处理记录。
主题名称
默认情况下,Informix 连接器会将表中发生的所有 INSERT、UPDATE 和 DELETE 操作的更改事件写入特定于该表的单个 Apache Kafka 主题。连接器使用以下约定命名更改事件主题:
topicPrefix.schemaName.tableName
以下列表提供了默认名称各组成部分的定义:
topicPrefix
由 topic.prefix 连接器配置属性指定的主题前缀。
schemaName
发生该操作的数据库架构的名称。
tableName
发生该操作的数据库表的名称。
例如,考虑一个 Informix 安装,其中 mydatabase 数据库在 myschema 架构下包含以下表:
productsproducts_on_handcustomersorders连接器会将事件发出到以下 Kafka 主题:mydatabase.myschema.productsmydatabase.myschema.products_on_handmydatabase.myschema.customersmydatabase.myschema.orders
连接器对内部数据库模式历史主题、模式变更主题和事务元数据主题采用类似的命名约定。
如果默认的主题名称不能满足你的需求,可以配置自定义主题名称。要配置自定义主题名称,需要在逻辑主题路由 SMT 中指定正则表达式。有关使用逻辑主题路由 SMT 自定义主题命名的更多信息,请参阅主题路由。
模式历史主题
当数据库客户端查询数据库时,该客户端使用数据库当前的模式。但是,数据库模式随时可能发生变化,这意味着连接器必须能够识别在记录每次插入、更新或删除操作时的模式是什么。此外,连接器不一定能将当前模式应用于每个事件。如果某个事件相对陈旧,它可能是在应用当前模式之前被记录的。
为确保在模式变更之后发生的事件能够得到正确处理,Debezium Informix 连接器会基于 Informix 变更数据表的结构存储新模式的快照,这些变更数据表与其关联的数据表结构相同。连接器将表模式信息以及导致模式变更的操作的 LSN 一起存储在数据库模式历史 Kafka 主题中。连接器使用存储的模式表示来生成变更事件,这些事件能够正确反映在每次插入、更新或删除操作发生时表的结构。
无论连接器是在崩溃后还是在正常停止后重启,它都会从上次读取的位置继续读取 Informix 变更数据表中的条目。基于连接器从数据库模式历史主题中读取的模式信息,连接器会应用在重启位置时存在的表结构。
如果你更新了处于捕获模式的 Informix 表的模式,务必同时更新相应变更表的模式。你必须是拥有提升权限的 Informix 数据库管理员,才能更新数据库模式。有关如何在 Debezium 环境中更新 Informix 数据库模式的更多信息,请参阅模式历史演进。
数据库模式历史主题仅供连接器内部使用。此外,连接器还可以选择将模式变更事件发送到另一个面向消费者应用程序的主题。
其他资源
- 接收 Debezium 事件记录的主题默认名称。
模式更改主题
你可以配置 Debezium Informix 连接器,使其生成模式更改事件,用于描述应用于数据库中表的模式变更。
在源数据库中发生以下操作后,Debezium 会向模式更改主题发送一条消息:
- 你启用了 Debezium 从新表捕获变更。
- 你为此前由 Debezium 捕获变更的表禁用了捕获功能。
连接器将模式更改事件写入名为 <topicPrefix> 的 Kafka 模式更改主题,其中 <topicPrefix> 是在 topic.prefix 连接器配置属性中指定的主题前缀。
模式更改事件的模式包含以下元素:
name
模式更改事件消息的名称。
type
变更事件消息的类型。
version
模式的版本。版本是一个整数,每次模式发生更改时递增。
fields
变更事件消息中包含的字段。
示例:Informix 连接器模式更改主题的模式
以下示例展示了一个以 JSON 格式表示的典型模式。
{
"schema": {
"type": "struct",
"fields": [
{
"type": "string",
"optional": false,
"field": "databaseName"
}
],
"optional": false,
"name": "io.debezium.connector.informix.SchemaChangeKey",
"version": 1
},
"payload": {
"databaseName": "inventory"
}
}连接器发送到 schema change 主题的消息包含一个负载,其中包含以下元素:
databaseName
语句所应用到的数据库名称。databaseName 的值用作消息键。
pos
语句在事务日志中出现的位置。
tableChanges
schema 变更后整个表结构的结构化表示。tableChanges 字段包含一个数组,其中列出了该表的每个列。由于结构化表示以 JSON 或 Avro 格式呈现数据,消费者无需先通过 DDL 解析器处理即可轻松读取消息。
| 对于处于捕获模式的表,连接器不仅会将 schema 变更历史存储到 schema change 主题中,还会存储到内部数据库 schema 历史主题中。内部数据库 schema 历史主题仅供连接器使用,不打算供消费应用直接使用。请确保需要 schema 变更通知的应用仅从 schema change 主题获取该信息。 |
|---|
切勿对数据库 schema 历史主题进行分区。数据库 schema 历史主题必须保持连接器向其发出的事件记录的全局一致顺序,才能正常运行。
为确保该主题不会被拆分到多个分区中,可通过以下任一方法设置该主题的分区数:
- 如果手动创建数据库 schema 历史主题,请将分区数指定为
1。 - 如果使用 Apache Kafka 代理自动创建数据库 schema 历史主题,请将 Kafka
num.partitions配置选项的值设置为1。
| 连接器向其 schema change 主题发出的消息格式尚处于孵化阶段,可能会在不另行通知的情况下发生变更。 |
|---|
示例:发送到 Informix 连接器 schema change 主题的消息
以下示例展示了 schema change 主题中的一条消息。该消息包含表结构的逻辑表示。
{
"schema": {
...
},
"payload": {
"source": {
"version": "3.6.3.Final",
"connector": "informix",
"name": "informix",
"ts_ms": 1588252618953,
"snapshot": "true",
"db": "testdb",
"schema": "informix",
"table": "customers",
"commit_lsn": "0",
"change_lsn": "0",
"txId": null,
"begin_lsn": "0"
},
"ts_ms": 1588252618953,
"databaseName": "testdb",
"schemaName": "informix",
"ddl": null,
"tableChanges": [
{
"type": "CREATE",
"id": "\"testdb\".\"informix\".\"customers\"",
"table": {
"defaultCharsetName": null,
"primaryKeyColumnNames": [
"id"
],
"columns": [
{
"name": "id",
"jdbcType": 4,
"nativeType": null,
"typeName": "int identity",
"typeExpression": "int identity",
"charsetName": null,
"length": 10,
"scale": 0,
"position": 1,
"optional": false,
"autoIncremented": false,
"generated": false
},
{
"name": "first_name",
"jdbcType": 12,
"nativeType": null,
"typeName": "varchar",
"typeExpression": "varchar",
"charsetName": null,
"length": 255,
"scale": null,
"position": 2,
"optional": false,
"autoIncremented": false,
"generated": false
},
{
"name": "last_name",
"jdbcType": 12,
"nativeType": null,
"typeName": "varchar",
"typeExpression": "varchar",
"charsetName": null,
"length": 255,
"scale": null,
"position": 3,
"optional": false,
"autoIncremented": false,
"generated": false
},
{
"name": "email",
"jdbcType": 12,
"nativeType": null,
"typeName": "varchar",
"typeExpression": "varchar",
"charsetName": null,
"length": 255,
"scale": null,
"position": 4,
"optional": false,
"autoIncremented": false,
"generated": false
}
],
"attributes": [
{
"customAttribute": "attributeValue"
}
]
}
}
]
}
}以下列表说明了前面的模式更改消息中的各个字段:
ts_ms
可选字段,显示连接器处理事件的时间。该时间基于运行 Kafka Connect 任务的 JVM 中的系统时钟。
在 source 对象中,ts_ms 表示在数据库中做出更改的时间。要确定源数据库中发生更改与 Debezium 处理该更改之间的时间延迟,请比较 payload.source.ts_ms 和 payload.ts_ms 的值。
databaseName、schemaName
标识包含该更改的数据库和模式。
ddl
对于 Informix 连接器,该字段始终为 null。对于其他连接器,该字段包含导致模式更改的 DDL。Informix 连接器无法获得该 DDL。
tableChanges
一个包含一项或多项的数组,这些项包含由 DDL 命令生成的模式更改。
tableChanges.type
描述更改的类型。该字段包含以下值之一:
CREATE | 创建了表。 |
|---|---|
ALTER | 修改了表。 |
DROP | 删除了表。 |
tableChanges.id
被创建、修改或删除的表的完整标识符。
tableChanges.table
表示应用更改之后的表元数据。
tableChanges.table.primaryKeyColumnNames
构成该表主键的列的列表。
tableChanges.table.columns
已更改表中每一列的元数据。
tableChanges.table.attributes
每次表更改的自定义属性元数据。
在连接器发送到模式更改主题的消息中,消息键是包含该模式更改的数据库的名称。在下面的示例中,payload 字段包含该键:
{
"schema": {
"type": "struct",
"fields": [
{
"type": "string",
"optional": false,
"field": "databaseName"
}
],
"optional": false,
"name": "io.debezium.connector.informix.SchemaChangeKey",
"version": 1
},
"payload": {
"databaseName": "testdb"
}
}事务元数据
Debezium 可以生成表示事务边界的事件,并对变更事件消息进行补充。
Debezium 接收事务元数据的限制
Debezium 仅对在你部署连接器之后发生的事务注册并接收元数据。在你部署连接器之前发生的事务,其元数据不可用。
Debezium 会为每个事务中的 BEGIN 和 END 分隔符生成事务边界事件。事务边界事件包含以下字段:
status
BEGIN 或 END。
id
唯一事务标识符的字符串表示形式,由 Informix 事务 ID 本身与给定操作的 LSN 组成,以冒号分隔,即格式为 txID:LSN。
ts_ms
事务边界事件(BEGIN 或 END 事件)在数据源处发生的时间。如果数据源未向 Debezium 提供事件时间,则该字段表示 Debezium 处理该事件的时间。
event_count(针对 END 事件)
该事务产生的事件总数。
data_collections(针对 END 事件)
一个由 data_collection 和 event_count 元素组成的对的数组,用于表示连接器为源自某个数据集合的变更所发出的事件数量。
示例
{
"status": "BEGIN",
"id": "571:53195829",
"ts_ms": 1486500577125,
"event_count": null,
"data_collections": null
}
{
"status": "END",
"id": "571:53195832",
"ts_ms": 1486500577691,
"event_count": 2,
"data_collections": [
{
"data_collection": "testdb.informix.tablea",
"event_count": 1
},
{
"data_collection": "testdb.informix.tableb",
"event_count": 1
}
]
}默认情况下,连接器将事务事件发送到 <topic.prefix>.transaction 主题。您可以通过修改 topic.transaction 属性的值来覆盖默认设置。
数据更改事件的丰富
启用事务元数据后,连接器会在变更事件的 Envelope 中加入一个新的 transaction 字段。该字段以复合字段的形式提供每个事件的相关信息:
id
唯一事务标识符的字符串表示。
total_order
该事件在该事务生成的所有事件中的绝对位置。
data_collection_order
该事件在该事务发出的所有事件中,按数据集合划分的位置。
下面是消息的示例:
{
"before": null,
"after": {
"pk": "2",
"aa": "1"
},
"source": {
...
},
"op": "c",
"ts_ms": "1580390884335",
"ts_us": "1580390884335641",
"ts_ns": "1580390884335641387",
"transaction": {
"id": "571:53195832",
"total_order": "1",
"data_collection_order": "1"
}
}数据变更事件
Debezium Informix 连接器为每个行级的 INSERT、UPDATE 和 DELETE 操作生成一个数据变更事件。每个事件都包含一个键和一个值,键和值的结构取决于发生变更的表。
Debezium 和 Kafka Connect 是为处理连续的事件消息流而设计的。然而,由于这些事件的结构可能随时间发生变化,消费者在处理某些 Debezium 事件时可能会遇到困难。为了应对这一挑战,每个事件都被设计为自包含的。也就是说,事件要么包含其内容的架构,要么在使用了模式注册表(schema registry)的环境中,包含一个模式 ID,消费者可据此从注册表中获取对应的模式。
以下示例中的 JSON 结构展示了一个典型的 Debezium 事件记录如何表示变更事件的四个基本组件。事件的具体表示形式取决于你为应用程序所配置的 Kafka Connect 转换器。只有在配置转换器生成 schema 字段时,变更事件中才会出现该字段。同样地,只有在配置转换器生成事件键和事件负载时,变更事件中才会包含它们。如果你使用 JSON 转换器,并配置它生成变更事件的全部四个基本部分,那么变更事件的结构如下:
{
"schema": {
...
},
"payload": {
...
},
"schema": {
...
},
"payload": {
...
},
}以下列表描述了前面基本变更事件中的字段:
`schema
第一个 schema 字段是事件键的一部分。它指定了 Kafka Connect 模式,用于描述事件键中 payload 部分的内容结构。换句话说,对于发生变更的表,第一个 schema 字段描述的是主键的结构;如果未定义主键,则描述表的唯一键的结构。
payload
第一个 payload 字段是事件键的一部分。它的结构由前面的 schema 字段描述,用于指定发生变更的行的键。
schema
第二个 schema 字段是事件值的一部分。它指定了用于描述事件值 payload 结构的 Kafka Connect 模式。换句话说,第二个 schema 描述的是被修改行的结构。通常,事件值模式包含嵌套模式。
payload
第二个 payload 字段是事件值的一部分。它具有事件值 schema 字段中描述的结构,并包含被修改行的实际数据。
默认情况下,连接器将变更事件记录发送到主题名与其来源表同名的主题中。有关更多信息,请参阅主题名称。
Debezium Informix 连接器确保所有 Kafka Connect 模式名称都符合 Avro 模式名称格式。符合 Avro 模式名称格式意味着逻辑服务器名称以拉丁字母或下划线开头,即 a-z、`A-Z` 或 _。逻辑服务器名称中的其余每个字符,以及数据库和表名称中的每个字符,必须是拉丁字母、数字或下划线,即 a-z、A-Z、0-9 或 \_。如果存在无效字符,则将其替换为下划线字符。
使用下划线替换无效字符可能会导致意外的冲突。例如,当逻辑服务器、数据库或表的名称包含一个或多个无效字符,且这些字符是唯一区分该名称与同类型其他实体名称的字符时,就可能出现冲突。
命名冲突也可能发生,因为 Informix 中数据库、模式和表的名称区分大小写。在某些情况下,连接器可能会将来自多个表的事件记录发送到同一个 Kafka 主题中。
变更事件键
变更事件的键包含被修改表的键的模式以及被修改行的实际键。该模式及其对应的载荷都包含一个字段,用于表示连接器创建事件时被修改表的 PRIMARY KEY(或唯一约束)中的每一列。
请看下面的 customers 表,其后是该表变更事件键的示例。
示例表
CREATE TABLE customers (
id INTEGER IDENTITY(1001,1) NOT NULL PRIMARY KEY,
first_name VARCHAR(255) NOT NULL,
last_name VARCHAR(255) NOT NULL,
email VARCHAR(255) NOT NULL UNIQUE
);示例变更事件键
当 Debezium 从 customers 表捕获变更时,会发出一个包含事件键模式的变更事件记录。只要 customers 表的定义保持不变,Debezium 从 customers 表捕获的每次变更所产生的事件记录都具有相同的键结构。下面的示例展示了该事件结构的 JSON 表示:
{
"schema": {
"type": "struct",
"fields": [
{
"type": "int32",
"optional": false,
"field": "ID"
}
],
"optional": false,
"name": "mydatabase.myschema.customers.Key"
},
"payload": {
"ID": 1004
}
}以下列表描述了前文变更事件键 JSON 对象中的各个字段:
schema
表示事件键的 schema 字段,显示描述事件键 payload 结构的 Kafka Connect schema。
schema.fields
为 payload 定义的字段定义数组。每个字段定义包含字段的名称、类型以及是否为必需字段。
schema.fields[].type
指定 payload 中某个字段的数据类型。在本例中,int32 表示 32 位整数。
schema.fields[].optional
指示该字段是否可以包含 null 值。在本例中,false 值表示该字段是必需字段,不能为 null。当表没有主键时,键的 payload 字段中的值是可选的。
schema.fields[].field
指定 payload 中字段的名称。在本例中,字段名称为 ID。
schema.name
指定定义键 payload 结构的 schema 名称。该 schema 描述了发生变更的表的主键结构。键 schema 的名称遵循以下格式:
<connector-name>.<database-name>.<table-name>.Key。
上例中的 schema 名称由以下各部分组成:
connector-name
mydatabase:生成此事件的连接器名称。
database-name
myschema:包含发生变更的表的数据库 schema。
table-name
customers:被更新的表的名称。
payload
指定发生变更事件的表行的键。在上例中,键包含单个 ID 字段,其值为 1004。
尽管 column.exclude.list 和 column.include.list 连接器配置属性允许你只捕获表列的子集,但主键或唯一键中的所有列始终会包含在事件的键中。 |
|---|
| 如果表没有主键或唯一键,则变更事件的键为 null。没有主键或唯一键约束的表中的行无法被唯一标识。 |
|---|
变更事件值
变更事件中的值比键稍显复杂。与事件键一样,值包含一个 schema 元素和一个 payload 元素。schema 元素包含描述 payload 元素 Envelope 结构的 schema,其中包括其嵌套字段。创建、更新或删除数据的操作的变更事件都具有采用封套(envelope)结构的值 payload。
考虑在前面的变更事件键示例中所使用的 customers 示例表:
示例表
CREATE TABLE customers (
id INTEGER IDENTITY(1001,1) NOT NULL PRIMARY KEY,
first_name VARCHAR(255) NOT NULL,
last_name VARCHAR(255) NOT NULL,
email VARCHAR(255) NOT NULL UNIQUE
);Debezium 从 customers 表捕获的每个变更事件的 value 元素都使用相同的模式。每个事件值的负载会因事件类型的不同而有所差异:
create 事件
create 事件表示行级 INSERT 操作,它会在事件负载的 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": "mydatabase.myschema.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": "mydatabase.myschema.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": false,
"field": "schema"
},
{
"type": "string",
"optional": false,
"field": "table"
},
{
"type": "string",
"optional": true,
"field": "commit_lsn"
},
{
"type": "string",
"optional": true,
"field": "change_lsn"
},
{
"type": "string",
"optional": true,
"field": "txId"
},
{
"type": "string",
"optional": true,
"field": "begin_lsn"
}
],
"optional": false,
"name": "io.debezium.connector.informix.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": "mydatabase.myschema.customers.Envelope"
},
"payload": {
"before": null,
"after": {
"id": 1005,
"first_name": "john",
"last_name": "doe",
"email": "[email protected]"
},
"source": {
"version": "3.6.3.Final",
"connector": "informix",
"name": "myconnector",
"ts_ms": 1559729468470,
"ts_us": 1559729468470000,
"ts_ns": 1559729468470000000,
"snapshot": false,
"db": "mydatabase",
"schema": "myschema",
"table": "customers",
"commit_lsn": "627404540760620",
"change_lsn": "627404540485812",
"txId": "157",
"begin_lsn": "627404540372400"
},
"op": "c",
"ts_ms": 1559729471739,
"ts_us": 1559729471739241,
"ts_ns": 1559729471739241367
}
}以下列表描述了前面 create 事件值中的各个字段:
schema
表示变更事件值的 schema 字段,它显示描述事件 payload 结构的 Kafka Connect schema。对于连接器为特定表生成的每个变更事件,其变更事件值的 schema 都是相同的。
schema.type
指定 schema 类型。struct 值表示该 schema 定义了一个包含多个字段的结构化数据类型。
schema.fields
payload 的字段定义数组。每个字段定义描述 payload 中的一个顶层字段,包括 before、after、source、op 以及时间戳字段。
schema.fields.name
io.debezium.connector.informix.Source 是 payload 中 source 字段的 schema。此 schema 是 Informix 连接器特有的。连接器为其生成的所有事件都使用它。
schema.optional
指示变更事件值是否必须包含 payload。false 值表示 payload 是必需的。
schema.name
指定定义变更事件 payload 结构的 schema 名称。变更事件值的 schema 名称遵循以下格式:
<connector-name>.<database-name>.<table-name>.Envelope。
在此示例中,mydatabase.myschema.customers.Envelope 表示 customers 表的 envelope schema。
payload
变更事件的实际数据。payload 中的信息说明了该事件如何更改了表中某一行的数据。它提供了变更前后该行的状态,以及源元数据、操作类型和时间戳。
payload.before
一个可选字段,指定事件发生前该行的状态。对于 create(插入)操作,此字段始终为 null,因为在插入之前该行并不存在。
payload.after
一个可选字段,指定事件发生后该行的状态。对于 create 操作,此字段包含新插入行中所有列的值。在此示例中,它显示的新行为 id 为 1005、first_name 为 john、last_name 为 doe、email 为 [email protected]。
payload.source
一个必需字段,描述事件的源元数据。source 结构显示此次变更的 Informix 元数据,从而提供可追溯性。你可以使用 source 元素中的信息来比较同一主题内或不同主题之间的事件,以了解此事件是发生在其他事件之前、之后,还是与其他事件属于同一次提交。该字段包含有关变更发生的数据库、表和事务上下文的信息。
payload.source.connector
生成该事件的连接器类型。在此示例中,informix 表示该事件是由 Debezium Informix 连接器发出的。
payload.op
一个必需的字符串字段,描述导致该事件的操作类型。在此示例中,c 表示 create(插入)操作。
payload.ts_ms
连接器处理该事件时的时间戳(以自 Unix 纪元起的毫秒数表示)。该时间基于运行 Kafka Connect 任务的 JVM 中的系统时钟。
payload.ts_us
连接器处理该事件时的时间戳(以自 Unix 纪元起的微秒数表示)。
payload.ts_ns
连接器处理该事件时的时间戳(以自 Unix 纪元起的纳秒数表示)。
update 事件
更新事件表示行级的 UPDATE 操作,它在 before 字段中记录变更前行的状态,并在 after 字段中记录变更后行的新状态。
在示例 customers 表中,更新事件的变更事件值与该表的 create 事件具有相同的 schema。同样,update 事件值的负载结构与 create 事件中值的负载结构相对应。但是,update 事件与 create 事件的 value 负载中并不包含相同的值。以下示例展示了连接器针对 customers 表中的更新所生成的事件记录的变更事件值:
{
"schema": { ... },
"payload": {
"before": {
"id": 1005,
"first_name": "john",
"last_name": "doe",
"email": "[email protected]"
},
"after": {
"ID": 1005,
"first_name": "john",
"last_name": "doe",
"email": "[email protected]"
},
"source": {
"version": "3.6.3.Final",
"connector": "informix",
"name": "myconnector",
"ts_ms": 1559729995937,
"ts_us": 1559729995937000,
"ts_ns": 1559729995937000000,
"snapshot": false,
"db": "mydatabase",
"schema": "myschema",
"table": "customers",
"commit_lsn": "627404540760620",
"change_lsn": "627404540485812",
"txId": "157",
"begin_lsn": "627404540372400"
},
"op": "u",
"ts_ms": 1559729998706,
"ts_us": 1559729998706742,
"ts_ns": 1559729998706742877
}
}以下列表描述了前面 update 事件值中的各个字段:
payload.before
一个可选字段,表示事件发生前行的状态。存在时,它包含一个对象,其中的字段表示变更前的列值。在本示例中,它显示了更新前的行状态,包括原始电子邮件地址 [email protected]。
对于 create(插入)操作,该字段为 null,因为事件发生前该行并不存在。 |
|---|
payload.after
一个可选字段,表示事件发生后行的状态。存在时,它包含一个对象,其中的字段表示变更后的列值。在本示例中,它显示了更新后的行状态,包括新的电子邮件地址 [email protected]。
对于 delete 操作,该字段为 null,因为事件发生后该行已不复存在。 |
|---|
payload.source
一个必填字段,用于描述事件的来源元数据。该字段包含有关变更发生的数据库、表和事务上下文的信息。
payload.op
一个必填的字符串字段,用于描述引发该事件的操作类型。值 u 表示该行因更新而发生了变更。
payload.ts_ms
连接器处理该事件时的时间戳(自 Unix 纪元以来的毫秒数)。该时间基于运行 Kafka Connect 任务的 JVM 中的系统时钟。它表示 Debezium 创建变更事件消息的时间,而非变更在数据库中发生的时间。
payload.ts_us
连接器处理该事件时的时间戳(自 Unix 纪元以来的微秒数)。
payload.ts_ns
连接器处理该事件时的时间戳(自 Unix 纪元以来的纳秒数)。
如果你更新了行的主键或唯一键所包含的列,就会改变行的键值。键值变更后,Debezium 会发出以下事件:
- 一个
DELETE事件 - 一个 墓碑事件,其中包含该行的旧键
- 一个包含该行新键的事件。
delete 事件
delete 事件表示行级 DELETE 操作,它将被删除行的最终状态捕获到 before 字段中,以便消费者能够识别并处理该移除操作。
表的 delete 变更事件中的 value 具有与同一表的 create 和 update 事件中类似的 schema 部分。当用户在示例 customers 表上执行 delete 操作后,Debezium 会发出如下例所示的事件消息:
{
"schema": { ... },
},
"payload": {
"before": {
"id": 1005,
"first_name": "john",
"last_name": "doe",
"email": "[email protected]"
},
"after": null,
"source": {
"version": "3.6.3.Final",
"connector": "informix",
"name": "myconnector",
"ts_ms": 1559730445243,
"ts_us": 1559730445243000,
"ts_ns": 1559730445243000000,
"snapshot": false,
"db": "mydatabase",
"schema": "myschema",
"table": "customers",
"commit_lsn": "627404540760620",
"change_lsn": "627404540485812",
"txId": "157",
"begin_lsn": "627404540372400"
},
"op": "d",
"ts_ms": 1559730450205,
"ts_us": 1559730450205104,
"ts_ns": 1559730450205104870
}
}以下列表描述了前面 delete 事件值中的各个字段:
payload.before
一个可选字段,表示事件发生前行的状态。对于 delete 操作,该字段包含行被删除前的最终状态。在本示例中,它显示了删除时的行值,包括 id 1005 和 email [email protected]。
payload.after
一个可选字段,指定事件发生后行的状态。对于删除操作,该字段始终为 null,因为删除后该行已不存在。
payload.source
一个必需字段,描述事件的源元数据。该字段包含有关发生更改的数据库、表和事务上下文的信息。
payload.op
一个必需的字符串字段,描述操作类型。值 d 表示该行已被删除。
payload.ts_ms
连接器处理事件时的时间戳(自 Unix 纪元起的毫秒数)。该时间基于运行 Kafka Connect 任务的 JVM 中的系统时钟。它表示 Debezium 创建更改事件消息的时间,而非数据库中发生更改的时间。
payload.ts_us
连接器处理事件时的时间戳(自 Unix 纪元起的微秒数)。
payload.ts_ns
连接器处理事件时的时间戳(自 Unix 纪元起的纳秒数)。
delete 更改事件记录为消费者提供了处理行移除所需的信息。该记录包含先前的值,以支持可能需要这些值来处理移除操作的消费者。
Informix 连接器事件被设计为可与 Kafka 日志压缩 配合使用。日志压缩允许移除一些较旧的消息,只要为每个键保留至少一条最新消息即可。保留最新消息使 Kafka 能够回收存储空间,同时确保主题包含可用于重新加载基于键的状态的完整数据集。
当某行被删除时,delete 事件值仍然可与日志压缩配合使用,因为 Kafka 可以移除具有相同键的所有较早消息。但是,要让 Kafka 移除具有相同键的所有消息,消息值必须为 null。为实现这一点,在 Debezium 的 Informix 连接器发出 delete 事件之后,连接器会发出一个具有相同键但值为 null 的特殊墓碑(tombstone)事件。
数据类型映射
有关 Informix 支持的数据类型的完整说明,请参阅 Informix 文档中的数据类型部分。
以下数据类型不支持数据捕获:
- 简单大对象(
TEXT和BYTE数据类型) - 用户定义的数据类型
- 集合数据类型(
SET、MULTISET、LIST和ROW数据类型)
有关更多信息,请参阅 Informix 文档中的 Data for capture。
Informix 连接器通过发出事件来表示行的变更,这些事件的结构与发生变更事件的源表结构一致。事件记录中包含每个列值对应的字段。为了用源列填充这些字段的值,连接器使用默认映射,将值从原始的 Informix 数据类型转换为 Kafka Connect 模式类型或语义类型。连接器为以下 Informix 数据类型提供了默认映射:
如果默认的数据类型转换不能满足你的需求,你可以为连接器创建自定义转换器。
基本类型
下表描述了连接器如何将每种 Informix 数据类型映射为事件字段中的字面类型和语义类型。
- 字面类型描述值如何使用 Kafka Connect 模式类型来表示:
INT8、INT16、INT32、INT64、FLOAT32、FLOAT64、BOOLEAN、STRING、BYTES、ARRAY、MAP和STRUCT。 - 语义类型描述 Kafka Connect 模式如何通过字段的 Kafka Connect 模式名称来体现该字段的含义。
表 5. Informix 基本数据类型的映射
Informix 数据类型 字面类型(模式类型) 语义类型(模式名称)及说明
BIGINT
INT64
不适用
BIGSERIAL
INT64
不适用
BLOB
BYTES
不适用
BOOLEAN
BOOLEAN
不适用
CHAR[(N)]
STRING
不适用
CLOB
STRING
不适用
DATE
INT32
io.debezium.time.Date
不带时区信息的日期
DATETIME
INT64
io.debezium.time.Timestamp
不带时区信息的时间戳
DECIMAL
BYTES
org.apache.kafka.connect.data.Decimal
DOUBLE
FLOAT64
不适用
FLOAT
FLOAT64
不适用
INTEGER
INT32
不适用
LVARCHAR[(N)]
STRING
不适用
NUMERIC
BYTES
org.apache.kafka.connect.data.Decimal
REAL
FLOAT32
不适用
SERIAL
INT32
不适用
SMALLINT
INT16
不适用
SMALLFLOAT
FLOAT32
不适用
TINYINT
INT16
取值范围为 0 到 255 的 8 位无符号整数值,因此需要存储为 int16
VARCHAR[(N)]
STRING
不适用
如果存在,列的默认值会传播到对应字段的 Kafka Connect 模式中。除非指定了显式的列值,否则变更事件中会包含该字段的默认值。因此,通常很少需要从模式中获取默认值。传递默认值有助于在使用 Avro 作为序列化格式并结合 Confluent 模式注册表时满足兼容性规则。
时间类型
Informix 根据 time.precision.mode 连接器配置属性的值来映射时间类型。以下各节说明了这些映射关系:
time.precision.mode=adaptive
为确保事件精确地表示数据库中的值,当 time.precision.mode 配置属性设置为默认值 adaptive 时,连接器会根据列的数据类型定义来确定字面类型和语义类型。
表 6. time.precision.mode 为 adaptive 时的映射
Informix 数据类型 | 字面类型(schema 类型) | 语义类型(schema 名称)及说明
DATE | INT32 | io.debezium.time.Date,表示自纪元以来的天数。
DATETIME | INT64 | io.debezium.time.Timestamp,表示自纪元以来的毫秒数,不包含时区信息。
time.precision.mode=connect
当 time.precision.mode 配置属性设置为 connect 时,连接器使用 Kafka Connect 逻辑类型。此设置适用于只能处理 Kafka Connect 内置逻辑类型、无法处理可变精度时间值的使用者。但是,由于 Informix 支持数十微秒的精度,如果连接器配置为使用 connect 时间精度,而数据库列的小数秒精度值大于 3,则连接器生成的事件会导致精度丢失。
表 7. time.precision.mode 为 connect 时的映射
Informix 数据类型 | 字面类型(schema 类型) | 语义类型(schema 名称)及说明
DATE | INT32 | org.apache.kafka.connect.data.Date,表示自纪元以来的天数。
DATETIME | INT64 | org.apache.kafka.connect.data.Timestamp,表示自纪元以来的毫秒数,不包含时区信息。
time.precision.mode=isostring
将 time.precision.mode 属性设置为 isostring,可将连接器配置为以 UTC 时区将时间值映射为 ISO-8601 格式的字符串。应用此设置后,连接器会使用语义类型 io.debezium.time.IsoDate 和 io.debezium.time.IsoTimestamp 来映射日期值和时间戳值。
表 8. time.precision.mode 为 isostring 时的映射
Informix 数据类型 | 字面类型(schema 类型) | 语义类型(schema 名称)及说明
DATE | STRING | io.debezium.time.IsoDate,按 ISO 8601 标准以 UTC 格式表示日期值,例如 2017-09-15Z。
DATETIME | STRING | io.debezium.time.IsoTimestamp,按 ISO 8601 标准以 UTC 格式表示时间戳值,例如 2019-07-09T02:28:57.123456Z。
INTERVAL
INTERVAL 类型不受 Informix Change Stream 客户端支持。
时间戳类型
DATETIME 类型表示不带时区信息的时间戳。此类列会基于 UTC 转换为等价的 Kafka Connect 值。例如,DATETIME 值 "2018-06-20 15:13:16.94514" 由值为 "1529507596000" 的 io.debezium.time.Timestamp 表示。
运行 Kafka Connect 和 Debezium 的 JVM 所使用的时区不会影响此转换。
十进制类型
下表描述了连接器如何将 Informix 十进制数据类型映射到变更事件字段中的 Kafka Connect 字面量类型和语义类型。
Informix 数据类型 字面量类型(架构类型) 语义类型(架构名称)与说明
NUMERIC[(P[,S])]
BYTES
org.apache.kafka.connect.data.Decimal
scale 架构参数包含一个整数,表示小数点向右移动的位数。connect.decimal.precision 架构参数包含一个整数,表示给定十进制值的精度。
DECIMAL[(P[,S])]
BYTES
org.apache.kafka.connect.data.Decimal
scale 架构参数包含一个整数,表示小数点向右移动的位数。connect.decimal.precision 架构参数包含一个整数,表示给定十进制值的精度。
设置 Informix
为了让 Debezium 捕获已提交到 Informix 表中的变更事件,拥有必要权限的 Informix 数据库管理员必须为变更数据捕获配置数据库。
执行以下任务,为使用变更数据捕获 API 做准备:
- 以数据库用户
informix的身份,运行$INFORMIXDIR/etc目录中的syscdcv1.sql脚本。这将安装syscdcv1数据库。 - 以用户
informix的身份建立与syscdcv1数据库的连接,验证该数据库是否存在。 - 将
DB_LOCALE环境变量设置为与要从中捕获数据的数据库的区域设置相同。
| 关于针对变更数据捕获优化 Informix 的具体指导超出了本文档的范围。 |
|---|
部署
要部署 Debezium Informix 连接器,需要安装 Debezium Informix 连接器归档包,配置连接器,并通过将配置添加到 Kafka Connect 来启动连接器。
先决条件
- 已安装 Apache Kafka 和 Kafka Connect。
- 已安装 Informix,并且已为表启用捕获模式,以便将数据库准备为可与 Debezium 连接器配合使用。
操作步骤
- 从 Maven Central 下载 Debezium Informix 连接器插件归档包。
- 将 JAR 文件解压到你的 Kafka Connect 环境中。
- 从 Maven Central 下载 Informix JDBC 驱动 和 Informix Change Stream 客户端,并将下载的 JAR 文件复制到包含 Debezium Informix 连接器 JAR 文件的目录中(即
debezium-connector-informix-3.6.3.Final.jar所在的目录)。
| 由于许可证要求,Debezium Informix 连接器归档包中不包含 Debezium 连接 Informix 数据库所需的 Informix JDBC 驱动和 Change Stream 客户端。要使连接器能够访问数据库,你必须将该驱动和客户端库添加到连接器环境中。 |
|---|
- 将包含 JAR 文件的目录添加到 Kafka Connect 的
plugin.path中。 - 重启 Kafka Connect 进程,以加载新的 JAR 文件。
如果你使用的是不可变容器,可以查阅 Debezium 的容器镜像,其中提供了已预装 Informix 连接器、可直接运行的 Apache Kafka 和 Kafka Connect 镜像。
你从 quay.io 获取的 Debezium 容器镜像未经过严格测试或安全分析,仅供测试和评估使用。这些镜像不适用于生产环境。为了降低生产部署中的风险,请只部署由可信供应商积极维护并经过潜在漏洞全面测试的容器。 |
|---|
你也可以在 Kubernetes 和 OpenShift 上运行 Debezium。
后续步骤
Informix 连接器配置示例
以下示例展示了某个连接器实例的配置,该实例从 192.168.99.100 上端口为 9088 的 Informix 服务器捕获数据,逻辑名称为 fullfillment。通常,你会通过在 JSON 文件中设置连接器可用的配置属性来配置 Debezium Informix 连接器。
你可以选择仅为数据库中部分模式和表生成事件。此外,你还可以忽略、遮蔽或截断包含敏感数据的列、超过指定大小的列,或者你不需要的列。
{
"name": "informix-connector",
"config": {
"connector.class": "io.debezium.connector.informix.InformixConnector",
"database.hostname": "192.168.99.100",
"database.port": "9088",
"database.user": "informix",
"database.password": "in4mix",
"database.dbname": "mydatabase",
"topic.prefix": "fullfillment",
"table.include.list": "mydatabase.myschema.customers",
"schema.history.internal.kafka.bootstrap.servers": "kafka:9092",
"schema.history.internal.kafka.topic": "schemahistory.fullfillment"
}
}以下列表说明了前面配置示例中的各个字段:
name
指定连接器注册到 Kafka Connect 服务时所使用的名称。
connector.class
指定 Debezium 连接器类的名称。
database.hostname
指定 Informix 实例的地址。
database.port
指定 Informix 实例的端口号。
database.user
指定 Informix 用户的名称。
database.password
指定 Informix 用户的密码。
database.dbname
指定要捕获其更改的数据库名称。
topic.prefix
指定 Informix 实例或集群的逻辑名称。该名称构成一个命名空间,用于连接器写入的所有 Kafka 主题名称、Kafka Connect 模式名称,以及使用 Avro 连接器 时对应 Avro 模式的命名空间。
table.include.list
指定 Debezium 应捕获其更改的所有表的列表。
schema.history.internal.kafka.bootstrap.servers
指定该连接器用于向数据库架构历史主题写入和恢复 DDL 语句的 Kafka 代理列表。
schema.history.internal.kafka.topic
指定连接器写入和恢复 DDL 语句的数据库架构历史主题的名称。该主题仅供内部使用,不打算由使用者直接使用。
有关可为 Debezium Informix 连接器设置的完整配置属性列表,请参阅 Informix 连接器属性。
你可以使用 POST 命令将此配置发送到正在运行的 Kafka Connect 服务。该服务会记录配置并启动一个连接器任务,执行以下操作:
- 连接到 Informix 数据库。
- 读取处于捕获模式的表的变更数据表。
- 将变更事件记录流式传输到 Kafka 主题。
添加连接器配置
要开始运行 Informix 连接器,请创建一个连接器配置并将其添加到 Kafka Connect 集群。
前置条件
- 已启用 Informix 复制,以便为处于捕获模式的表公开变更数据。
- 已安装 Informix 连接器。
操作步骤
- 为 Informix 连接器创建配置。
- 使用 Kafka Connect REST API 将该连接器配置添加到 Kafka Connect 集群。
结果
连接器启动后,它会对配置为捕获更改的 Informix 数据库表执行一致性快照。随后,连接器开始生成行级操作的数据变更事件,并将变更事件记录流式传输到 Kafka 主题。
连接器属性
Debezium Informix 连接器提供了众多配置属性,可用于实现适合你应用的连接器行为。许多属性都有默认值。属性的相关信息按如下方式组织:
数据库模式历史连接器配置属性,用于控制 Debezium 如何处理从数据库模式历史主题中读取的事件。
Debezium Informix 连接器必需配置属性
除非提供了默认值,否则以下配置属性为必填。
表 9. 必需的连接器配置属性
属性 默认值 描述
无默认值
连接器的唯一名称。你只能使用指定的名称注册一次连接器,后续的注册尝试将失败。所有 Kafka Connect 连接器都需要此属性。
无默认值
连接器的 Java 类名称。对于 Informix 连接器,始终使用 io.debezium.connector.informix.InformixConnector 这个值。
1
此连接器可以创建的最大任务数。Informix 连接器始终使用单个任务,因此不会使用此值,默认值始终可接受。
无默认值
Informix 数据库服务器的 IP 地址或主机名。
9088
Informix 数据库服务器的整数端口号。
无默认值
用于连接 Informix 数据库服务器的 Informix 数据库用户名。
无默认值
连接 Informix 数据库服务器时使用的密码。
无默认值
要从中流式传输变更的 Informix 数据库的名称。
无默认值
主题前缀,为承载 Debezium 正在捕获变更的数据库的特定 Informix 数据库服务器提供命名空间。该前缀在所有其他连接器之间必须唯一,因为它会用作此连接器接收记录的所有 Kafka 主题的名称前缀。数据库服务器的逻辑名称中只能使用字母数字字符、连字符、点和下划线。
| 请勿更改此属性的值。如果更改了该名称值,重启后连接器将不再继续向原主题发送事件,而是将后续事件发送到基于新值命名的主题。连接器还将无法恢复其数据库结构历史主题。 |
|---|
无默认值
一个可选的、以逗号分隔的正则表达式列表,用于匹配你希望连接器捕获其变更的表的完全限定标识符。设置此属性后,连接器仅捕获指定表中的变更。每个标识符的形式为 databaseName.schemaName.tableName。默认情况下,连接器会捕获每个非系统表中的变更。
为了匹配表名,Debezium 会将你指定的正则表达式作为锚定正则表达式应用。也就是说,指定的表达式会与该表的整个标识符进行匹配;它不会匹配表名中可能出现的子串。
如果在配置中包含该属性,请不要同时设置 table.exclude.list 属性。
无默认值
一个可选的、以逗号分隔的正则表达式列表,用于匹配你不想捕获其变更的表的完全限定标识符。连接器会捕获排除列表之外的每个非系统表的变更。每个标识符的格式为 databaseName.schemaName.tableName。
为了匹配表名,Debezium 会将你指定的正则表达式作为锚定正则表达式来应用。也就是说,指定的表达式会与表的整个标识符进行匹配,而不会匹配可能出现在表名中的子字符串。
如果在配置中包含该属性,请不要同时设置 table.include.list 属性。
无默认值
一个可选的、以逗号分隔的正则表达式列表,用于匹配需要包含在变更事件记录值中的列的完全限定名。列的完全限定名遵循以下格式之一:databaseName.tableName.columnName,或 databaseName.schemaName.tableName.columnName。
为了匹配列名,Debezium 会将你指定的正则表达式作为锚定正则表达式来应用。也就是说,指定的表达式会与列的整个名称字符串进行匹配,而不会匹配可能出现在列名中的子字符串。如果在配置中包含该属性,请不要同时设置 column.exclude.list 属性。
无默认值
一个可选的、以逗号分隔的正则表达式列表,用于匹配需要从变更事件值中排除的列的完全限定名。列的完全限定名遵循以下格式之一:databaseName.tableName.columnName,或 databaseName.schemaName.tableName.columnName。
为了匹配列名,Debezium 会将你指定的正则表达式作为锚定正则表达式来应用。也就是说,指定的表达式会与列的整个名称字符串进行匹配,而不会匹配可能出现在列名中的子字符串。主键列始终包含在事件的键中,即使你在此属性中设置的值试图将其排除。如果在配置中包含该属性,请不要设置 column.include.list 属性。
column.mask.hash.hashAlgorithm.with.salt.salt
不适用
一个可选的、以逗号分隔的正则表达式列表,用于匹配基于字符的列的完全限定名。列的完全限定名遵循以下两种格式之一:databaseName.tableName.columnName,或 databaseName.schemaName.tableName.columnName。
在匹配列名时,Debezium 会将你指定的正则表达式作为锚定正则表达式来应用。也就是说,指定的表达式会与列的整个名称字符串进行匹配,而不会匹配列名中可能出现的子字符串。在生成的变更事件记录中,指定列的值会被替换为化名。
化名由应用指定的 hashAlgorithm(哈希算法)和 salt(盐值)后得到的哈希值构成。所使用的哈希函数在保留引用完整性的同时,将列值替换为化名。受支持的哈希函数在 Java 密码学架构标准算法名称文档的 MessageDigest 章节中有所说明。
在下面的示例中,CzQMA0cB5K 是随机选择的盐值。
column.mask.hash.SHA-256.with.salt.CzQMA0cB5K = inventory.orders.customerName, inventory.shipment.customerName如有需要,哈希别名会自动缩短至与列长度相匹配。连接器配置可以包含多个属性,用于指定不同的哈希算法和盐值。
根据所使用的 hashAlgorithm、所选的 salt 以及实际数据集的不同,生成的数据集可能无法被完全遮蔽。
adaptive
指定连接器用于表示时间、日期和时间戳值的数字精度。可指定以下值之一:
adaptive
根据表列的数据类型,连接器使用毫秒、微秒或纳秒精度值,来精确表示源表中存在的时间和时间戳值。
connect
连接器始终使用默认的 Kafka Connect 格式来表示 Time、Date 和 Timestamp 值,该格式使用毫秒精度,而不考虑源表中为该列配置的精度。有关更多信息,请参阅时间类型。
true
指定 delete 事件之后是否跟随一个墓碑(tombstone)事件。可指定以下值之一:
true
对于每次删除操作,连接器会发出一个 delete 事件,以及随后的一个墓碑事件。选择此选项可确保 Kafka 能够删除与被删除行的键相关的所有事件。如果禁用了墓碑事件,而目标主题启用了日志压缩,Kafka 可能无法识别并删除所有共享该键的事件。
false
连接器仅发出一个 delete 事件。
true
一个布尔值,指定连接器是否将数据库架构的变更发布到与主题前缀同名的 Kafka 主题中。连接器以包含数据库名称的键、以及描述该架构更新的 JSON 结构作为值,来记录每一次架构变更。这种记录架构变更的机制,独立于连接器内部对数据库架构历史变更的记录。
column.truncate.to.length.chars
无
一个可选的、以逗号分隔的正则表达式列表,用于匹配基于字符的列的完全限定名称。如果您希望在列中的数据超出属性名称中 length 所指定的字符数时截断这些数据,可设置此属性。将 length 设置为正整数值,例如 column.truncate.to.20.chars。
列的完全限定名称遵循以下格式之一:databaseName.tableName.columnName,或 databaseName.schemaName.tableName.columnName。
为匹配列名,Debezium 会将你指定的正则表达式作为锚定正则表达式来应用。也就是说,指定的表达式会与列的整个名称字符串进行匹配;该表达式不会匹配列名中可能出现的子字符串。
你可以在单个配置中指定多个长度不同的属性。
n/a
一个可选的、以逗号分隔的正则表达式列表,用于匹配基于字符的列的完全限定名称。如果你想让连接器对一组列的值进行掩码处理,例如这些列包含敏感数据时,可设置此属性。将 length 设置为正整数,即可用属性名称中 length 所指定数量的星号(*)字符替换指定列中的数据。将 length 设置为 0(零),即可用空字符串替换指定列中的数据。
列的完全限定名称遵循以下格式之一:databaseName.tableName.columnName,或 databaseName.schemaName.tableName.columnName。
为匹配列名,Debezium 会将你指定的正则表达式作为锚定正则表达式来应用。也就是说,指定的表达式会与列的整个名称字符串进行匹配;该表达式不会匹配列名中可能出现的子字符串。
你可以在单个配置中指定多个长度不同的属性。
n/a
一个可选的、以逗号分隔的正则表达式列表,用于匹配你希望连接器发出表示列元数据的额外参数的那些列的完全限定名称。设置此属性后,连接器会向事件记录的模式中添加以下字段:
__debezium.source.column.type__debezium.source.column.length__debezium.source.column.scale
这些参数会传递原始数据类型名称以及适用的类型属性,例如可变宽度类型的长度和数值类型的精度。启用连接器发出这些额外数据,有助于在目标数据库中正确设置特定数值列或字符列的大小。
列的完全限定名称遵循以下格式之一:databaseName.tableName.columnName,或 databaseName.schemaName.tableName.columnName。
要匹配列名,Debezium 会将你指定的正则表达式作为锚定正则表达式来应用。也就是说,指定的表达式会与列的整个名称字符串进行匹配;该表达式不会匹配列名中可能出现的子字符串。
datatype.propagate.source.type
无
一个可选的、以逗号分隔的正则表达式列表,用于指定数据库中为列定义的数据类型的全限定名称。设置此属性后,对于数据类型匹配的列,连接器发出的事件记录会在其模式(schema)中包含以下额外字段:
__debezium.source.column.type__debezium.source.column.length__debezium.source.column.scale
这些参数会传播原始数据类型名称以及适用的类型属性,例如可变宽度类型的长度和数值类型的精度(scale)。启用连接器发出这些额外数据,有助于在目标数据库中正确设置特定数值列或字符列的大小。
列的全限定名称遵循以下格式之一:databaseName.tableName.typeName,或 databaseName.schemaName.tableName.typeName。
要匹配数据类型名,Debezium 会将你指定的正则表达式作为锚定正则表达式来应用。也就是说,指定的表达式会与数据类型的整个名称字符串进行匹配;该表达式不会匹配类型名中可能出现的子字符串。
有关 Informix 特有数据类型名称的列表,请参阅 Informix 数据类型映射。
空字符串
一个表达式列表,用于指定连接器为发布到指定表对应 Kafka 主题的变更事件记录构造自定义消息键时所使用的列。
默认情况下,Debezium 使用表的主键列作为其发出记录的消息键。若要替代默认设置,或为缺少主键的表指定键,你可以基于一个或多个列配置自定义消息键。
要为表建立自定义消息键,请先列出该表,然后列出用作消息键的列。列表中的每个条目采用以下格式:
<fully-qualified_tableName>:<keyColumn>,<keyColumn>
要基于多个列名构建表键,请在各列名之间插入逗号。每个全限定表名都是采用以下格式的正则表达式:
<databaseName>.<schemaName>.<tableName>
该属性可以列出多个表的条目。使用分号分隔列表中不同表的条目。
以下示例为 inventory.customers 和 purchaseorders 表设置消息键:
inventory.customers:pk1,pk2;(.*).purchaseorders:pk3,pk4
在上例中,列 pk1 和 pk2 被指定为 inventory.customer 表的消息键。对于任意 schema 中的 purchaseorders 表,列 pk3 和 pk4 用作消息键。
none
指定如何调整 schema 名称,以兼容连接器所使用的消息转换器。可以指定以下值之一:
none
不作调整。
avro
将对 Avro 类型无效的字符替换为下划线(_)。
avro_unicode
将下划线或对 Avro 类型无效的字符替换为相应的 unicode 编码,例如 _uxxxx。
注意:_ 是转义序列,等同于 Java 中的反斜杠。
none
指定如何调整字段名称,以兼容连接器所使用的消息转换器。可以指定以下值之一:
none
不作调整。
avro
将对 Avro 类型无效的字符替换为下划线(_)。
avro_unicode
将下划线或对 Avro 类型无效的字符替换为相应的 unicode 编码,例如 _uxxxx。
注意:_ 是转义序列,等同于 Java 中的反斜杠。
有关 Avro 兼容性的更多信息,请参见 Avro 命名。
连接器高级配置属性
以下高级配置属性的默认值适用于大多数场景,因此很少需要在连接器配置中指定。
表 10. 连接器高级配置属性
| 属性 | 默认值 | 说明 |
|---|---|---|
converters | 无默认值 | 列出连接器可使用的自定义转换器实例的符号名称(以逗号分隔)。例如,isbn。必须设置 converters 属性,连接器才能使用自定义转换器。对于为连接器配置的每个转换器,还必须添加一个 .type 属性,用于指定实现该转换器接口的类的完全限定名。.type 属性使用以下格式:<converterSymbolicName>.type。例如, |
isbn.type: io.debezium.test.IsbnConverter如果要进一步控制已配置的转换器的行为,可以添加一个或多个配置参数,以便向转换器传递值。要将任意附加配置参数与某个转换器关联,请在参数名前加上该转换器的符号名作为前缀。例如:
isbn.schema.name: io.debezium.informix.type.Isbninitial
指定连接器启动时执行快照的判据:
always
连接器每次启动都会执行快照。快照包含被捕获表的结构和数据。指定此值后,每次连接器启动时,都会用被捕获表中数据的完整表示来填充主题。快照完成后,连接器开始流式传输后续数据库变更的事件记录。
initial
连接器执行数据库快照,方式如执行初始快照的默认工作流中所述。快照完成后,连接器开始流式传输后续数据库变更的事件记录。
initial_only
仅当该逻辑服务器名称没有记录任何 offset 时,连接器才执行数据库快照。快照完成后,连接器停止运行,不会转入流式传输后续数据库变更的事件记录。
schema_only
已弃用,请参阅 no_data。
no_data
连接器运行快照以捕获所有相关表的结构,执行默认快照工作流中描述的所有步骤,但不会创建 READ 事件来表示连接器启动时的数据集(第 7.b 步)。
recovery
设置此选项可恢复丢失或损坏的数据库架构历史主题。重启后,连接器会运行快照,从源表重建该主题。你也可以设置此属性,定期清理异常增长的数据库架构历史主题。
如果在上次连接器关闭之后数据库中已提交了架构变更,请不要使用 recovery 模式执行快照。 |
|---|
when_needed
连接器启动后,仅在检测到以下情况之一时才执行快照:
- 无法检测到任何主题 offset。
- 之前记录的 offset 指定了服务器上不可用的日志位置。
configuration_based
使用此选项,你可以通过一组以 snapshot.mode.configuration.based 为前缀的连接器属性来控制快照行为。
custom
custom 快照模式允许你注入自己实现的 io.debezium.spi.snapshot.Snapshotter 接口。将 snapshot.mode.custom.name 配置属性设置为你的实现中 name() 方法所提供的名称。
有关更多信息,请参阅自定义快照器 SPI。
snapshot.mode.configuration.based.snapshot.data
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.configuration.based.snapshot.on.schema.error
false
如果 snapshot.mode 设置为 configuration_based,可设置此属性以指定当结构历史主题不可用时,连接器是否在快照中包含表结构。
snapshot.mode.configuration.based.snapshot.on.data.error
false
如果 snapshot.mode 设置为 configuration_based,此属性指定当连接器在事务日志中找不到最后提交的偏移量时,是否尝试对表数据进行快照。
将该值设置为 true,可指示连接器执行新的快照。
无默认值
如果 snapshot.mode 设置为 custom,请使用此设置指定自定义实现的名称。该名称由 io.debezium.spi.snapshot.Snapshotter 接口中定义的 name() 方法提供。连接器重启后,Debezium 会调用指定的自定义实现,以确定是否执行快照。
有关更多信息,请参阅自定义快照器 SPI。
exclusive
控制连接器是否持有表锁以及持锁的时长。表锁可在快照期间防止其他数据库客户端执行某些表操作。您可以设置以下值:
exclusive
当 snapshot.isolation.mode 为 REPEATABLE_READ 或 EXCLUSIVE 时,此选项控制连接器在执行架构快照期间对表施加锁的方式。连接器会持有表锁,但仅在快照的初始阶段(即连接器读取数据库架构和其他元数据的阶段)保证对表的独占访问。在快照的后续阶段,连接器使用无需加锁的 flashback 查询来选取每个表中的所有行。
share
当 snapshot.isolation.mode 为 REPEATABLE_READ 或 EXCLUSIVE 时,此选项控制连接器在执行架构快照期间对表施加锁的方式。连接器会持有读表锁,但仅在快照的初始阶段(即连接器读取数据库架构和其他元数据的阶段)保证对表的读取访问。在快照的后续阶段,连接器使用无需加锁的 flashback 查询来选取每个表中的所有行。如果你希望快照在不施加任何锁的情况下运行,请设置以下选项 none。
none
阻止连接器在快照期间获取任何表锁。仅当快照创建期间不可能发生架构更改时,才使用此设置。
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。
select_all
指定连接器在执行快照期间如何查询数据。可设置为以下选项之一:
select_all
连接器默认执行 select all 查询,并可根据列包含和排除列表配置有选择地调整所选取的列。
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。
repeatable_read
在快照期间,指定事务隔离级别,以及连接器锁定处于捕获模式的表的持续时间。请指定以下值之一:
read_uncommitted
在初始快照期间不阻止其他事务更新表行。此模式不提供任何数据一致性保证;部分数据可能会丢失或损坏。
read_committed
在初始快照期间不阻止其他事务更新表行。新记录可能出现两次:一次在初始快照中,一次在流式传输阶段。不过,此一致性级别适用于数据镜像。
repeatable_read
在初始快照期间阻止其他事务更新表行。新记录可能出现两次:一次在初始快照中,一次在流式传输阶段。不过,此一致性级别适用于数据镜像。
exclusive
使用可重复读隔离级别,但对所有待读取的表加排他锁。此模式在初始快照期间阻止其他事务更新表行。只有 exclusive 模式能保证完全一致性:初始快照与流式日志构成一条线性历史。
5
指定对变更流客户端的读取调用超时行为的正整数值。请指定以下值之一:
<0
不超时。
0
如果没有可用数据,则立即返回。
>=1
指定连接器在超时前等待数据的秒数。
65536
指定 Informix Change Stream Client 处理的每批记录的最大大小的正整数值。
64
用于指定 Informix Change Stream Client 处理的 CDC 记录的最大数量的正整数值。
true
指定在流式传输关闭时,Informix 是否应停止对被监视表进行整行日志记录的布尔值。
true
指定 Debezium Informix 连接器是否为空事务发出记录的布尔值。
event.processing.failure.handling.mode
fail
指定连接器在处理事件期间如何处理异常。可指定以下值之一:
fail
连接器记录有问题事件的偏移量并停止处理。
warn
连接器记录有问题事件的偏移量,并继续处理下一个事件。
skip
连接器跳过有问题的事件,并继续处理下一个事件。
500(0.5 秒)
用于指定连接器在开始处理一批事件之前,等待新变更事件出现的毫秒数的正整数值。
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 字节时,向队列的写入操作就会被阻塞。
0
指定连接器向 Kafka 主题发送心跳消息的时间间隔,单位为毫秒。
心跳消息有助于确认连接器仍在处理事务日志,并确保它将最新的偏移量提交到 Kafka。
将此属性设置为一个正整数即可启用并调度心跳消息。默认情况下,连接器不会发出心跳消息。
在没有心跳消息的情况下,即使数据库日志接收到大量变更,只要这些变更不涉及被捕获的表,连接器就没有机会提交最新的偏移量。因此,如果连接器重启,它将从一个过期的偏移量处恢复读取日志,这可能导致重新发送大量积压的变更事件消息。通过启用心跳,即使受监控的表中发生的变更很少,也能确保连接器将最新的偏移量发送到 Kafka。
无默认值
指定连接器在发送心跳消息时在源数据库上执行的查询。
这有助于解决以下情况:当低流量数据库与高流量数据库位于同一主机上时,从低流量数据库捕获变更会使 Debezium 无法在物理日志文件轮转之前更新其最后已提交/重启 LSN。为了解决此问题,可在低流量数据库中创建一个心跳表,并将此属性设置为向该表插入记录的语句,例如:
INSERT INTO test_heartbeat_table (text) VALUES ('test_heartbeat')
这样,连接器就能从低流量数据库接收变更,并在物理日志文件轮转之前更新其最后已提交的/重启 LSN。
无默认值
连接器在启动后执行快照前应等待的时间间隔(以毫秒为单位)。如果在集群中启动多个连接器,此属性有助于避免快照中断,因为快照中断可能导致连接器重新平衡。
true
指定连接器是否为流式处理指标收集高级统计数据,例如分位数。设置为 true 时,连接器收集以下统计数据:
- 最小值
- 最大值
- 平均值
- P50(中位数)百分位数
- P95 百分位数
- P99 百分位数
目前仅 MilliSecondsBehindSource 指标支持收集分位数。 |
|---|
统计数据使用一种概率数据结构(DDSketch)计算得出,该结构可提供相对精度为 1% 的近似分位数值。设置为 false 时,连接器不收集分位数,相应的分位数 JMX 指标将返回 null 值。连接器仍会继续收集最小值、最大值和平均值。禁用分位数收集可略微降低内存开销。有关更多信息,请参阅流式处理指标。
0
指定连接器在完成快照后延迟开始流式处理过程的时间(以毫秒为单位)。设置延迟间隔有助于防止连接器在快照刚完成但流式处理尚未开始时因发生故障而重新执行快照。
请将延迟值设置为高于为 Kafka Connect 工作进程配置的 offset.flush.interval.ms 属性的值。
snapshot.include.collection.list
table.include.list 中指定的所有表
一个可选的、以逗号分隔的正则表达式列表,用于匹配要包含在快照中的表的完全限定名称(databaseName.schemaName.tableName)。指定的项必须出现在连接器的 table.include.list 属性中。此属性仅在连接器的 snapshot.mode 属性设置为 no_data 以外的值时才生效。
此属性不会影响增量快照的行为。
要匹配表的名称,Debezium 会将你指定的表达式作为锚定正则表达式来应用。也就是说,指定的表达式会与表的完整名称字符串进行匹配;它不会匹配表名中可能出现的子字符串。
2000
在快照期间,连接器按批次读取表内容。此属性指定每个批次中的最大行数。
10000
指定执行快照时获取表锁所需等待的最长时间(以毫秒为单位)。如果连接器在此时间段内无法获取表锁,快照将失败。有关更多信息,请参阅连接器如何执行快照。
可以指定以下设置之一:
一个大于 0 的整数
即连接器等待获取表锁的毫秒数。如果连接器在指定的时间段结束前无法获取锁,快照将失败。
0
如果连接器无法获取锁,快照将立即失败。
-1
连接器将无限期等待以获取锁。
snapshot.select.statement.overrides
无默认值
指定连接器使用自定义 SELECT 语句来确定哪些行应包含在快照中的表。
此属性仅影响快照。它不适用于连接器在流式传输阶段从日志中读取的事件。
此属性由两部分组成,它们协同工作:
主属性
以 <databaseName>.<schemaName>.<tableName> 格式书写的完全限定表名的逗号分隔列表。此列表标识了你要为其指定自定义快照查询的表。
例如,
"snapshot.select.statement.overrides": "mydatabase.inventory.products,mydatabase.customers.orders"
如果完全限定的数据库、Schema 或表名包含特殊字符,例如空格、方括号([ 或 ])或句点字符(.),请将字符串用双引号括起来,以防止连接器将这些特殊字符解释为分隔符。如果完全限定的表名不包含空格或特殊字符,则表名周围的双引号是可选的。 |
|---|
次属性
对于在主属性中列出的每个表,你都必须定义对应的 snapshot.select.statement.overrides.<databaseName>.<schemaName>.<tableName> 属性,用于指定快照期间要运行的自定义 SELECT 语句。该 SELECT 语句决定了表中哪些行会被包含在快照中。
例如,要为 mydatabase.customer.orders 表指定 SELECT 语句,请添加以下属性:
snapshot.select.statement.overrides.mydatabase.customers.orders
| 如果某个表已在主属性中列出,但缺少相应的次级属性,连接器会记录一条警告日志,并对该表使用默认的快照行为。 |
|---|
配置示例
以下示例展示了如何配置 snapshot.select-statement 属性,以便对 mydatabase.customers.orders 表执行快照,且仅包含未被软删除的记录;也就是说,软删除字段 delete_flag 的值为 0。
"snapshot.select.statement.overrides": "mydatabase.customer.orders",
"snapshot.select.statement.overrides.mydatabase.customer.orders": "SELECT * FROM customers.orders WHERE delete_flag = 0 ORDER BY id DESC"false
确定连接器是否生成带有事务边界的事件,并用事务元数据充实更改事件信封。如果您希望连接器执行这些操作,请将该值设置为 true。有关更多信息,请参阅事务元数据。
t
以逗号分隔的操作类型列表,表示连接器在流式传输期间跳过的操作。您可以指定以下值:
c
连接器不为插入(create)操作发出事件。
u
连接器不为更新操作发出事件。
d
连接器不为删除操作发出事件。
t
连接器不为截断操作发出事件。
none
连接器为所有操作类型发出事件。
无默认值
用于向连接器发送信号的数据集合的完全限定名称。集合名称区分大小写。
请使用以下格式指定集合名称:
<databaseName>.<schemaName>.<tableName>
source
为连接器启用的信号通道名称列表。默认情况下,可以使用以下通道:
sourcekafkafilejmx
您还可以选择实现自定义信号通道。
无默认值
为连接器启用的通知通道名称列表。默认情况下,可以使用以下通道:
sinklogjmx
您还可以选择实现自定义通知通道。
incremental.snapshot.chunk.size
1024
连接器在增量快照分块期间获取并读入内存的最大行数。增大分块大小可提高效率,因为快照执行的查询次数更少,但每次查询的数据量更大。不过,更大的分块大小也需要更多内存来缓冲快照数据。请将分块大小调整为在您的环境中能提供最佳性能的值。
incremental.snapshot.watermarking.strategy
insert_insert
指定连接器在增量快照期间使用的水印机制,用于对重复事件去重——这些事件可能被增量快照捕获,并在流处理恢复后再次被捕获。您可以指定以下选项之一:
insert_insert
当您发送信号以启动增量快照时,对于 Debezium 在快照期间读取的每个分块,它都会向信号数据集合写入一个条目,以记录打开快照窗口的信号。快照完成后,Debezium 会插入第二个条目,记录关闭窗口的信号。
insert_delete
当您发送信号以启动增量快照时,对于 Debezium 读取的每个分块,它都只向信号数据集合写入单个条目,以记录打开快照窗口的信号。快照完成后,该条目会被移除。不会为关闭快照窗口的信号创建任何条目。设置此选项可防止信号数据集合快速增长。
io.debezium.schema.SchemaTopicNamingStrategy
连接器用于为数据变更、模式变更、事务、心跳及其他类型事件构建主题名称的 TopicNamingStrategy 类的名称。
.
指定连接器用于构建主题名称的分隔符。
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。
transaction
指定连接器追加到其发送事务元数据消息的主题名称之后的字符串。生成的主题名称遵循以下模式:
topic.prefix.transaction
例如,如果主题前缀为 fulfillment,根据该属性的默认值,连接器会为事务元数据主题指定以下名称:fulfillment.transaction。
1
指定连接器执行初始快照时使用的线程数。该值必须是大于或等于 1 的正整数。当设置为大于 1 的值时,连接器会执行并行快照,按主键范围将每个表划分为多个块(chunk),并在所有可用线程之间并发处理这些块。
| 没有主键的表以及使用快照选择覆盖(snapshot select override)的表,将作为一个块处理,并回退到单线程快照。 |
|---|
snapshot.max.threads.multiplier
1
控制并行快照期间每个表创建的块数量的全局乘数。默认情况下,Debezium 为每个线程创建一个块。乘数越大,连接器在执行快照时会为每个表创建更多、更小的块。较小的块有助于让线程负载更加均衡。例如,在 4 个线程且乘数为 2 的情况下,Debezium 会创建 8 个块,而不是 4 个。
要为特定表覆盖此乘数,请将 snapshot.max.threads.multiplier.<fully_qualified_table_name> 设置为所需的值。
| 与增大倍数相比,增加线程数通常对提高快照吞吐量更有效。 |
|---|
false
默认值(false)允许连接器通过将表拆分为多个块,并使用独立线程并发处理每个块来加速初始快照。若要恢复为每个线程处理一张表的传统并行快照行为,请将该值设置为 true。
启用传统行为后,完成自身表快照的线程在等待其他线程完成时会处于空闲状态。在强制执行连接超时的环境中,空闲连接可能会导致连接器在快照结束后无法正常关闭连接,进而引发异常,即使快照已成功捕获所有数据也是如此。如果遇到此问题,请将 snapshot.max.threads 设置为 1 并重试快照。 |
|---|
属性 internal.legacy.snapshot.max.threads 是 legacy.snapshot.max.threads 的已废弃别名,不应在新配置中使用。 |
|---|
无默认值
定义用于自定义 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 的默认屏蔽模式;指定的值不会扩展或增强默认模式。 |
|---|
-1
指定连接器在失败之前对可重试错误(如连接错误)的最大重试次数。请设置以下值之一:
-1 | 不限次数。 |
|---|---|
0 | 已禁用。不允许重试。 |
> 0 | 最大重试次数。 |
600000(10 分钟)
指定连接器等待查询完成的时间,以毫秒为单位。将该值设置为 0(零)可移除超时限制。
true
此属性指定 Debezium 是否为其发出的消息添加带有前缀 __debezium.context. 的上下文标头。
这些标头是 OpenLineage 集成所必需的,它们提供的元数据使下游处理系统能够追踪并识别变更事件的来源。
该属性会添加以下标头:
__debezium.context.connectorLogicalName
Debezium 连接器的逻辑名称。
__debezium.context.taskId
连接器任务的唯一标识符。
__debezium.context.connectorName
Debezium 连接器的名称。
Debezium Informix 连接器数据库 Schema 历史配置属性
Debezium 提供了一组 schema.history.internal.* 属性,用于控制连接器与 Schema 历史主题的交互方式。
下表描述了用于配置 Debezium 连接器的 schema.history.internal 属性。
表 11. 连接器数据库 Schema 历史配置属性
属性 默认值 描述
schema.history.internal.kafka.topic
无默认值
连接器存储数据库 Schema 历史的 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 捕获变更事件的逻辑数据库和模式中的表记录架构结构。
false
连接器记录所有逻辑数据库的架构结构。
schema.history.internal.memory.optimization
off
控制 Debezium 如何使用驻留器(interner)在内存中对相同的架构对象(表、列、属性)进行去重。请指定以下值之一:
off
不执行任何去重(默认值)。
on
每个连接器使用其独立的隔离驻留池。可在不影响其他连接器的情况下,减少单个连接器的堆内存占用。
shared
所有配置为 shared 的连接器共享一个全局驻留池。当许多连接器跟踪结构相似的表时,可最大化去重效果。
传递式 Informix 连接器配置属性
连接器支持传递式属性,使 Debezium 能够指定自定义配置选项,以便对 Apache Kafka 生产者和消费者的行为进行微调。有关 Kafka 生产者和消费者全部配置属性的信息,请参阅 Kafka 文档。
用于配置生产者和消费者客户端与架构历史主题交互方式的传递式属性
Debezium 依赖 Apache Kafka 生产者将架构更改写入数据库架构历史主题。同样,连接器启动时,它依赖 Kafka 消费者从数据库架构历史主题中读取数据。你可以通过为一组以 schema.history.internal.producer.* 和 schema.history.internal.consumer.* 前缀开头的直通配置属性赋值,来定义 Kafka 生产者和消费者客户端的配置。这些直通的生产者和消费者数据库架构历史属性控制着一系列行为,例如这些客户端如何与 Kafka 代理建立安全连接,如以下示例所示:
schema.history.internal.producer.security.protocol=SSL
schema.history.internal.producer.ssl.keystore.location=/var/private/ssl/kafka.server.keystore.jks
schema.history.internal.producer.ssl.keystore.password=test1234
schema.history.internal.producer.ssl.truststore.location=/var/private/ssl/kafka.server.truststore.jks
schema.history.internal.producer.ssl.truststore.password=test1234
schema.history.internal.producer.ssl.key.password=test1234
schema.history.internal.consumer.security.protocol=SSL
schema.history.internal.consumer.ssl.keystore.location=/var/private/ssl/kafka.server.keystore.jks
schema.history.internal.consumer.ssl.keystore.password=test1234
schema.history.internal.consumer.ssl.truststore.location=/var/private/ssl/kafka.server.truststore.jks
schema.history.internal.consumer.ssl.truststore.password=test1234
schema.history.internal.consumer.ssl.key.password=test1234Debezium 在将属性传递给 Kafka 客户端之前,会先从属性名中移除该前缀。
有关 Kafka 生产者配置属性和 Kafka 消费者配置属性的更多信息,请参阅 Apache Kafka 文档。
用于配置 Informix 连接器如何与 Kafka 信号主题交互的透传属性
Debezium 提供了一组 signal.* 属性,用于控制连接器与 Kafka 信号主题的交互方式。
下表描述了 Kafka signal 属性。
表 12. Kafka 信号配置属性
属性 | 默认值 | 说明
<topic.prefix>-signal
连接器用于监控临时信号的 Kafka 主题名称。
| 如果禁用了自动创建主题,则必须手动创建所需的信号主题。为了保证信号的顺序,必须使用信号主题。信号主题必须只有一个分区。 |
|---|
kafka-signal
Kafka 消费者所使用的组 ID 名称。
signal.kafka.bootstrap.servers
无默认值
连接器用于建立与 Kafka 集群初始连接的主机和端口对列表。每一对都指向 Debezium Kafka Connect 进程所使用的 Kafka 集群。
100
一个整数值,用于指定连接器在轮询信号时等待的最大毫秒数。
用于配置信号通道 Kafka 消费者客户端的透传属性
Debezium 连接器支持对信号 Kafka 消费者进行透传配置。透传信号属性以 signal.consumer.* 前缀开头。例如,连接器会将 signal.consumer.security.protocol=SSL 之类的属性传递给 Kafka 消费者。
Debezium 在将这些属性传递给 Kafka 信号消费者之前,会先移除属性名中的前缀。
用于配置 Informix 连接器接收器通知通道的传递属性
下表描述了可用于配置 Debezium 接收器 notification 通道的属性。
| 属性 | 默认值 | 描述 |
|---|---|---|
notification.sink.topic.name | 无默认值 | 接收来自 Debezium 通知的主题名称。当您将 notification.enabled.channels 属性配置为包含 sink 作为启用的通知通道之一时,此属性为必填项。 |
表 13. 接收器通知配置属性
Debezium 连接器传递数据库驱动程序配置属性
Debezium 连接器支持对数据库驱动程序进行传递配置。传递的数据库属性以 driver.* 前缀开头。例如,连接器会将 driver.foobar=false 之类的属性传递给 JDBC URL。
Debezium 在将属性传递给数据库驱动程序之前会去除属性的前缀。
监控
Debezium Informix 连接器提供三种类型的指标,这些指标是对 Apache Kafka 和 Kafka Connect 提供的 JMX 指标内建支持的补充。
- 快照指标提供连接器执行快照期间的运行信息。
- 流处理指标提供连接器捕获变更并流式传输变更事件记录时的运行信息。
- Schema 历史指标提供连接器 Schema 历史的状态信息。
Debezium 监控文档详细介绍了如何通过 JMX 暴露这些指标。
自定义 MBean 名称
Debezium 连接器通过连接器的 MBean 名称暴露指标。这些指标针对每个连接器实例,提供有关连接器快照、流处理和 Schema 历史进程行为的数据。
默认情况下,当你部署一个正确配置的连接器时,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 名称
默认情况下,Informix 连接器为流式指标使用以下 MBean 名称:
debezium.informix:type=connector-metrics,context=streaming,server=<topic.prefix>如果你将 custom.metric.tags 的值设置为 database=salesdb-streaming,table=inventory,Debezium 将生成如下的自定义 MBean 名称:
debezium.informix:type=connector-metrics,context=streaming,server=<topic.prefix>,database=salesdb-streaming,table=inventory快照指标
MBean 为 debezium.informix: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 | 当前正在快照的表的主键集合的上界。 |
流处理指标
MBean 为 debezium.informix: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 | 队列
Schema 历史记录指标
MBean 为 debezium.informix:type=connector-metrics,context=schema-history,server=<topic.prefix>。
下表列出了可用于监控连接器 Schema 历史记录过程的 JMX 指标,包括恢复状态、已应用的 Schema 变更数量,以及最近一次变更的时间戳。
| 属性 | 类型 | 描述 |
|---|---|---|
Status | string | 描述数据库模式历史状态的取值之一:STOPPED、RECOVERING(正在从存储中恢复历史记录)、RUNNING。 |
RecoveryStartTime | long | 恢复开始的时间,以纪元秒为单位。 |
ChangesRecovered | long | 恢复阶段读取的变更数量。 |
ChangesApplied | long | 恢复期间和运行时应用的模式变更总数。 |
MilliSecondsSinceLastRecoveredChange | long | 自从上次从历史存储中恢复变更以来经过的毫秒数。 |
MilliSecondsSinceLastAppliedChange | long | 自上次应用变更以来经过的毫秒数。 |
LastRecoveredChange | string | 最近一次从历史存储中恢复的变更的字符串表示。 |
LastAppliedChange | string | 最近一次应用的变更的字符串表示。 |
模式演进
Debezium Informix 连接器虽然能够捕获模式变更,但要更新模式,你必须与数据库管理员协作,以确保连接器能持续生成变更事件。
| 当你对某个表发起模式更新时,必须等待更新过程完成,然后才能对同一表再次执行新的模式更新。如果可能,最好将所有 DDL 放在同一批次中执行,只执行一次模式更新过程。 |
|---|
离线模式更新
Informix 不支持在捕获变更的同时进行在线模式更新。你必须先停止 Debezium Informix 连接器,然后再执行模式更新。
| 由于必须停止 Debezium 才能完成模式更新过程,为尽量减少对下游应用的影响,最好在计划的维护窗口期间执行此操作。 |
|---|
前提条件
- 一个或多个处于捕获模式的表需要进行模式更新。
操作步骤
- 挂起更新数据库的应用。
- 等待 Debezium 连接器完成所有尚未传输的变更事件记录的流式传输。
- 停止 Debezium 连接器。
- 对源表模式应用所有变更。
- 恢复更新数据库的应用。
- 重新启动 Debezium 连接器。
评论
登录后参与评论
KnowForge