源连接器

SQL Server

师成师成· 更新于 2026-09-28· 阅读 353 分钟· 0 次阅读

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

Debezium SQL Server 连接器

想帮助我们进一步打磨和改进它吗?了解方式。

Debezium SQL Server 连接器捕获 SQL Server 数据库模式中发生的行级变更。

有关与此连接器兼容的 SQL Server 版本信息,请参阅 Debezium 版本概览。

Debezium SQL Server 连接器首次连接到 SQL Server 数据库或集群时,会对数据库中的模式进行一致性快照。初始快照完成后,连接器会持续捕获已启用 CDC 的 SQL Server 数据库中提交的 INSERT、UPDATE 或 DELETE 操作所产生的行级变更。连接器为每个数据变更操作生成事件,并将其流式传输到 Kafka 主题。连接器会将某个表的所有事件流式传输到专门的 Kafka 主题中。应用和服务随后可以消费该主题中的数据变更事件记录。

概述

Debezium SQL Server 连接器基于 变更数据捕获 功能,该功能适用于 SQL Server 2016 Service Pack 1 (SP1) 及更高版本 的 Standard 版或 Enterprise 版。SQL Server 捕获进程会监视指定的数据库和表,并将变更存储到专门创建的、带有存储过程外观(facade)的 变更表 中。

要使 Debezium SQL Server 连接器能够捕获数据库操作的变更事件记录,必须首先在 SQL Server 数据库上启用变更数据捕获。必须在数据库以及要捕获的每个表上都启用 CDC。在源数据库上设置好 CDC 之后,连接器便可以捕获数据库中发生的行级 INSERT、UPDATE 和 DELETE 操作。连接器会将每个源表的事件记录写入专门用于该表的 Kafka 主题。每个被捕获的表都对应一个主题。客户端应用读取其所关注的数据库表对应的 Kafka 主题,并可以对其从这些主题消费的行级事件作出响应。

连接器首次连接到 SQL Server 数据库或集群时,会对配置为捕获变更的所有表的模式进行一致性快照,并将该状态流式传输到 Kafka。快照完成后,连接器会持续捕获随后发生的行级变更。通过首先建立所有数据的一致视图,连接器可以在不丢失快照期间所做任何变更的情况下继续读取。

Debezium SQL Server 连接器具备容错能力。当连接器读取变更并生成事件时,它会定期在数据库日志中记录事件的位置(LSN / 日志序列号)。如果连接器因任何原因停止(包括通信故障、网络问题或崩溃),重启后连接器会从上次读取的位置继续读取 SQL Server 的 CDC 表。

偏移量会定期提交,而不是在变更事件发生时立即提交。因此,在故障中断之后,可能会产生重复事件。

容错机制同样适用于快照。也就是说,如果连接器在快照过程中停止,重启后连接器会重新开始新的快照。

SQL Server 连接器的工作原理

要优化配置并运行 Debezium SQL Server 连接器,了解连接器如何执行快照、流式传输变更事件、确定 Kafka 主题名称以及使用元数据是很有帮助的。

快照

SQL Server CDC 并非为存储数据库变更的完整历史而设计。为了让 Debezium SQL Server 连接器能够建立数据库当前状态的基准,它会使用称为*快照(snapshotting)*的过程。初始快照会捕获数据库中表的结构和数据。

Debezium SQL Server 连接器执行初始快照所使用的默认工作流

Debezium SQL Server 连接器在开始流式传输变更事件之前,会执行初始快照以捕获数据库的当前状态。您可以通过修改 snapshot.mode 配置属性的值来自定义快照工作流。

以下工作流列出了 Debezium 创建快照所采取的步骤。这些步骤描述的是 snapshot.mode 配置属性设置为默认值 initial 时的快照过程。如果您配置了不同的快照模式,连接器将以该工作流的修改版本来完成快照。

  1. 与数据库建立连接。

  2. 确定要捕获的表。默认情况下,连接器会捕获所有非系统表。若要让连接器捕获表或表元素的子集,你可以设置若干 include 和 exclude 属性来过滤数据,例如 table.include.list 或 table.exclude.list。

  3. 对已启用 CDC 的 SQL Server 表获取锁,以防止在创建快照期间发生结构变更。锁的级别由 snapshot.isolation.mode 配置属性决定。

  4. 读取服务器事务日志中的最大日志序列号(LSN)位置。

  5. 捕获所有非系统表的结构,或所有被指定为捕获目标的表的结构。连接器会将此信息持久化到其内部数据库模式历史主题中。模式历史提供了变更事件发生时生效的结构信息。

    默认情况下,连接器会捕获数据库中每个处于捕获模式的表的模式,包括那些未配置为捕获目标的表。如果表未配置为捕获目标,初始快照只会捕获其结构,而不会捕获任何表数据。有关为什么快照会为未包含在初始快照中的表保留模式信息的详细说明,请参阅了解为什么初始快照会捕获所有表的模式。
  6. 如有必要,释放第 3 步中获取的锁。其他数据库客户端现在可以向此前被锁定的任意表写入数据。

  7. 连接器在第 4 步读取的 LSN 位置处扫描要捕获的表。在扫描过程中,连接器会完成以下任务:

  8. 确认该表在快照开始之前就已创建。如果该表是在快照开始之后创建的,连接器会跳过该表。当快照完成、连接器转入流式传输后,它会为快照开始之后创建的所有表发出变更事件。

    1. 为从表中捕获的每一行生成一个 read 事件。所有 read 事件都包含相同的 LSN 位置,即第 4 步中获取的 LSN 位置。
    2. 将每个 read 事件发送到该表对应的 Kafka 主题。
  9. 在连接器偏移量中记录快照的成功完成。

由此生成的初始快照捕获了已启用 CDC 的表中每一行的当前状态。以此基准状态为基础,连接器会持续捕获后续发生的变化。

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

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

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

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

设置 说明

always

每次启动连接器时都执行快照。快照完成后,连接器开始为后续的数据库变更流式传输事件记录。

initial

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

initial_only

连接器执行数据库快照后,在流式传输任何变更事件记录之前停止,不允许捕获任何后续变更事件。

schema_only

已弃用,请参见 no_data。

no_data

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

recovery

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

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

when_needed

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

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

when_needed_no_data

连接器启动后,会像 when_needed 选项一样执行快照。与 no_data 选项一样,它不会创建 READ 事件来表示连接器启动时的数据集(第 7.b 步)。

configuration_based

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

custom

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

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

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

默认情况下,初始快照会捕获数据库中每张表的架构信息,而不仅仅是被指定捕获的表。理解这一行为有助于你控制快照范围,并在后续将表添加到捕获集时避免架构历史错误。

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

表数据

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

架构数据

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

连接器要求表的架构必须先存在于架构历史主题中,才能捕获该表。通过让初始快照捕获不属于原始捕获集的表的架构数据,Debezium 使连接器能够在日后需要时方便地捕获这些表的事件数据。如果初始快照没有捕获某张表的架构,你必须先将该架构添加到历史主题中,连接器才能从该表捕获数据。

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

其他资源

捕获初始快照未捕获的表中的数据(无模式变更)

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

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

前提条件

  • 你希望从一个连接器在初始快照期间未捕获其模式的表中捕获数据。
  • 在连接器读取的最早和最晚变更表条目对应的 LSN 之间,该表没有发生任何模式变更。有关捕获发生过结构变更的新表中的数据的信息,请参阅捕获初始快照未捕获的表中的数据(有模式变更)。

操作步骤

  1. 停止连接器。
  2. 删除由 schema.history.internal.kafka.topic 属性指定的内部数据库模式历史主题。
  3. 清除已配置的 Kafka Connect offset.storage.topic 中的偏移量。有关如何删除偏移量的更多信息,请参阅 Debezium 社区常见问题解答。
移除偏移量的操作仅应由熟悉 Kafka Connect 内部数据操作的高级用户执行。此操作具有潜在的破坏性,只能作为最后的手段使用。
  1. 对连接器配置进行以下更改:

    1. (可选)将 schema.history.internal.store.only.captured.tables.ddl 的值设置为 false。此设置会使快照捕获所有表的架构,并保证连接器今后能够重建所有表的架构历史。

      捕获所有表架构的快照需要更长的时间才能完成。
    2. 将你要让连接器捕获的表添加到 table.include.list 中。

    3. 将 snapshot.mode 设置为以下值之一:

      initial

      重启连接器时,连接器会对数据库执行完整快照,捕获表数据和表结构。如果选择此选项,建议将 schema.history.internal.store.only.captured.tables.ddl 属性的值设置为 false,以便连接器能够捕获所有表的架构。

      no_data

      重启连接器时,连接器只执行捕获表架构的快照。与完整数据快照不同,此选项不会捕获任何表数据。如果你希望比完整快照更快地重启连接器,可以使用此选项。

  2. 重启连接器。连接器将完成 snapshot.mode 所指定类型的快照。

  3. (可选)如果连接器执行的是 no_data 快照,则在快照完成后,发起增量快照以捕获你新增表中的数据。连接器在继续从这些表流式传输实时更改的同时执行快照。运行增量快照会捕获以下数据更改:

    • 对于连接器之前已捕获的表,增量快照会捕获连接器停机期间发生的更改,即从连接器停止到本次重启之间这段时间内的更改。
    • 对于新添加的表,增量快照会捕获所有现有的表行。

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

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

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

前置条件

  • 你想要从一张在初始快照期间连接器未捕获其架构的表中获取数据。
  • 该表已被应用了架构变更,导致待捕获的记录结构不统一。

操作步骤

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

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

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

如果初始快照未保存你想要捕获的表的架构,请完成以下操作步骤之一:

操作步骤 1:先执行架构快照,再执行增量快照

在此操作步骤中,连接器首先执行架构快照。之后你可以发起增量快照,使连接器能够同步数据。

  1. 停止连接器。
  2. 删除由 schema.history.internal.kafka.topic 属性指定的内部数据库架构历史主题。
  3. 清除已配置的 Kafka Connect offset.storage.topic 中的位点信息。有关如何移除位点的更多信息,请参阅 Debezium 社区常见问题。
移除偏移量仅应由熟悉 Kafka Connect 内部数据操作的高级用户执行。此操作具有潜在的破坏性,只能作为最后的手段。
  1. 按照以下步骤为连接器配置中的属性设置值:

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

  3. 等待 Debezium 捕获新表和现有表的架构。连接器停止之后发生的任何表的数据变更都不会被捕获。

  4. 为确保不丢失数据,请发起增量快照。

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

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

  1. 停止连接器。

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

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

    移除偏移量仅应由熟悉 Kafka Connect 内部数据操作的高级用户执行。此操作具有潜在的破坏性,只能作为最后的手段。
  4. 编辑 table.include.list,添加你想要捕获的表。

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

  6. 将 snapshot.mode 属性的值设置为 initial。

    1. (可选)将 schema.history.internal.store.only.captured.tables.ddl 设置为 false。
  7. 重启连接器。连接器会对整个数据库执行快照。快照完成后,连接器会切换到流式传输模式。

  8. (若要捕获连接器离线期间发生的任何数据变更,可发起增量快照。

基于分块的并行快照

基于分块的并行快照通过将较小的工作单元分布到多个线程中来加速初始快照,从而改善负载均衡并提高快照的可靠性。

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

当您将 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 信号的 type 设置为 incremental 或 blocking,并按如下表所述提供要包含在快照中的表名:

表 2. 即席 execute-snapshot 信号记录示例

字段 默认值 说明
type incremental 指定你希望运行的快照类型。目前你可以请求 incremental(增量)或 blocking(阻塞)快照。
data-collections 不适用 一个数组,包含与要纳入快照的表的完全限定名相匹配的正则表达式。对于 SQL Server 连接器,使用以下格式指定表的完全限定名:database.schema.table。
additional-conditions 不适用 一个可选数组,用于指定连接器求值的一组附加条件,以确定要包含在快照中的记录子集。每个附加条件都是一个对象,用于指定即席快照所捕获数据的过滤标准。每个附加条件可以设置以下参数:
data-collection:过滤器所应用的表的完全限定名。你可以为每张表应用不同的过滤器。
filter:指定数据库记录中必须存在的列值,快照才会将其包含在内,例如 "color='blue'"。
你赋给 filter 参数的值,与为阻塞快照设置 snapshot.select.statement.overrides 属性时在 SELECT 语句的 WHERE 子句中指定的值类型相同。
surrogate-key 不适用 一个可选字符串,用于指定连接器在快照过程中作为表主键使用的列名。

触发即席增量快照

你可以通过向信号表添加一条 type 为 execute-snapshot 的记录,或向 Kafka 信号主题发送信号消息来发起即席增量快照。连接器处理该消息后,便会开始快照操作。快照过程会读取每个表的首个和末个主键值,并将这些值作为每张表的起始点和结束点。根据表中的记录数量以及配置的区块大小,Debezium 将表划分为若干区块,并依次逐个对每个区块进行快照。

有关更多信息,请参见增量快照。

触发即席阻塞快照

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

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

增量快照

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

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

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

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

增量快照流程

运行增量快照时,Debezium 会按主键对每张表排序,然后依据配置的分块大小将表拆分为多个块。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 Server 排序规则注意事项

每个 SQL Server 服务器或数据库都配置为使用特定的排序规则,该排序规则决定了数据库如何存储、排序、比较和显示字符数据。某些 SQL Server 排序规则集(例如 (SQL_*))的排序规则不符合 Unicode 排序算法。在某些情况下,不兼容的排序规则可能导致连接器执行临时快照时丢失数据。例如,如果 SQL Server 配置为以 Unicode 格式发送字符串(即连接属性 sendStringParametersAsUnicode 设置为 true),连接器在快照期间可能会跳过记录。为防止在临时快照期间丢失数据,请将连接字符串属性 driver.sendStringParametersAsUnicode 的值设置为 false。

其他资源

Debezium 的 SQL Server 连接器不支持在增量快照运行期间进行模式变更。

触发增量快照

要启动增量快照,你可以向源数据库中的信号表发送一个临时快照信号。快照信号以 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\""。

前提条件

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

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

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

例如,

INSERT INTO db1.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 支持以下物理行标识符作为代理键:

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

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

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

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

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

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

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

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

SELECT * FROM <tableName> ....

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

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

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

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

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

  • id(主键)
  • color
  • quantity

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

INSERT INTO db1.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. Execute snapshot 数据字段

字段 默认值

type

incremental

要执行的快照类型。目前 Debezium 支持 incremental 和 blocking 两种类型。更多详情请参阅下一节。

data-collections

N/A

一个由逗号分隔的正则表达式数组,用于匹配要包含在快照中的表的完全限定名称。

请使用与 signal.data.collection 配置选项相同的格式指定名称。数据集合名称区分大小写。

additional-conditions

N/A

一个可选的附加条件数组,用于指定连接器评估的条件标准,以确定要包含在快照中的记录子集。

每个附加条件都是一个对象,用于指定临时快照捕获数据的过滤标准。你可以为每个附加条件设置以下参数:

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 检测到信号表中的变更后,会读取该信号,并在增量快照操作正在进行时将其停止。

其他资源

前提条件

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

  1. 向信号表发送 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 字段必须包含以下字段:

字段默认值说明
typeincremental要执行的快照类型。目前 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"}]}

可能的重复记录

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

读取变更数据表

连接器首次启动时,会获取所捕获表的结构快照,并将该信息持久化到其内部数据库模式历史主题中。随后,连接器会为每个源表确定一个变更表,并完成以下步骤。

  1. 对于每个变更表,连接器读取上次存储的最大 LSN 与当前最大 LSN 之间创建的所有变更。
  2. 连接器按提交 LSN 和变更 LSN 的值升序排列读取到的变更。此排序顺序确保 Debezium 以与数据库中相同顺序回放这些变更。
  3. 连接器将提交 LSN 和变更 LSN 作为偏移量传递给 Kafka Connect。
  4. 连接器存储最大 LSN,并从步骤 1 重新开始该过程。

重启后,连接器从中断处读取的最后一个偏移量(提交 LSN 和变更 LSN)恢复处理。

连接器能够检测所包含的源表是否已启用或禁用 CDC,并据此调整其行为。

数据库中未记录最大 LSN

如果数据库中未记录最大 LSN,Debezium SQL Server 连接器将无法开始流处理。LSN 缺失最常见的原因是 SQL Server 代理未运行。

最大 LSN 可能因以下任一原因而缺失:

  • SQL Server 代理未运行。
  • 变更表中尚无任何变更记录。
  • 数据库活动量较低,CDC 清理作业已定期清除了 CDC 表中的条目。

在这些原因中,只有 SQL Server 代理已停止才是真正的错误状态,因为 SQL Server 代理正在运行是连接器正常工作的前提条件。其他两种情况属于正常情况。为了区分代理已停止与正常情况,连接器会查询 SQL Server 代理的状态。如果 SQL Server 代理未运行,连接器会在日志中写入 ERROR 级别信息:"数据库中未记录最大 LSN;SQL Server 代理未运行"。

为帮助诊断问题,请运行以下查询检查 SQL Server 代理的状态:

"SELECT CASE WHEN dss.[status]=4 THEN 1 ELSE 0 END AS isRunning FROM [#db].sys.dm_server_services dss WHERE dss.[servicename] LIKE N'SQL Server Agent (%';"

运行前面的查询来检查 SQL Server Agent 的运行状态需要 VIEW SERVER STATE 权限。VIEW SERVER STATE 是一个实例级别的权限,授予该权限后,用户可以查看该实例上承载的每个数据库。作为授予此权限的替代方案,你可以配置 database.sqlserver.agent.status.query 属性来运行自定义查询。

例如,database.sqlserver.agent.status.query=SELECT [#db].func_is_sql_server_agent_running(),其中 [#db] 是数据库名称的占位符。

请定义一个函数,当 SQL Server Agent 正在运行时返回 true 或 1,否则返回 false 或 0。通过配置该属性来运行这个函数,你就可以在不向任何用户授予 VIEW SERVER STATE 权限的情况下,安全地使用高级权限。

其他资源

限制

Debezium SQL Server 连接器并不支持对所有 SQL Server 数据库对象类型进行变更捕获。索引视图(也称为物化视图)不受支持,因为 SQL Server CDC 要求捕获实例的基础对象必须是一个表。

主题名称

默认情况下,SQL Server 连接器会将表中发生的所有 INSERT、UPDATE 和 DELETE 操作的事件写入到该表专属的单个 Apache Kafka 主题中。连接器使用以下约定为变更事件主题命名:<topicPrefix>.<schemaName>.<tableName>

以下列表给出了默认名称中各组成部分的定义:

topicPrefix

服务器的逻辑名称,由 topic.prefix 配置属性指定。

schemaName

发生变更事件的数据库架构的名称。

tableName

发生变更事件的数据库表的名称。

例如,如果 fulfillment 是服务器的逻辑名称,dbo 是架构名称,数据库中包含名为 products、products_on_hand、customers 和 orders 的表,那么连接器将把变更事件记录流式传输到以下 Kafka 主题:

  • fulfillment.testDB.dbo.products
  • fulfillment.testDB.dbo.products_on_hand
  • fulfillment.testDB.dbo.customers
  • fulfillment.testDB.dbo.orders

连接器采用类似的命名约定来标记其内部数据库 schema history 主题、schema 变更主题以及事务元数据主题。

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

Schema history 主题

Debezium SQL Server 连接器使用一个内部的 schema history 主题来存储 schema 变更记录,使其能够在任何时间捕获的变更事件上应用正确的 schema。schema history 主题仅供连接器内部使用,不用于供应用程序消费。

由于数据库 schema 可能随时发生变化,连接器必须能够识别每条变更事件被捕获时生效的 schema。

为确保正确处理 schema 变更之后发生的变更事件,Debezium SQL Server 连接器会基于 SQL Server 变更表中的结构存储新 schema 的快照,这些变更表镜像了其对应数据表的结构。连接器使用存储的 schema 表示来生成变更事件,从而正确反映每次插入、更新或删除操作时表的结构。

无论连接器是因崩溃还是正常停止而重启,它都会从上次读取的位置继续读取 SQL Server CDC 表中的条目。根据连接器从数据库 schema history 主题中读取的 schema 信息,连接器会应用其重启位置上存在的表结构。

如果你更新了处于捕获模式的 SQL Server 表的 schema,还必须更新相应变更表的 schema。你必须是拥有提升权限的 SQL Server 数据库管理员才能更新数据库 schema。有关在 Debezium 环境中更新 SQL Server 数据库 schema 的更多信息,请参阅数据库 schema 演进。

此外,连接器还可以选择将 schema 变更事件发送到另一个面向消费者应用程序的主题。

其他资源

Schema 变更主题

对于启用了 CDC 的每个表,Debezium SQL Server 连接器会存储应用于数据库中表的 schema 变更事件的历史记录。连接器将 schema 变更事件写入名为 <topicPrefix> 的 Kafka topic,其中 topicPrefix 是在 topic.prefix 配置属性中指定的逻辑服务器名称。

连接器发送到 schema 变更 topic 的消息包含一个负载,并且可选地包含变更事件消息的 schema。

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

name

schema 变更事件消息的名称。

type

变更事件消息的类型。

version

schema 的版本号。版本号是一个整数,每当 schema 发生更改时该整数就会递增。

fields

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

示例:SQL Server 连接器 schema 变更 topic 的 schema

以下示例展示了一个典型的 JSON 格式 schema。

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

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

databaseName

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

tableChanges

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

当连接器被配置为捕获某张表时,它不仅会将该表的架构变更历史存储到架构变更主题中,还会存储到一个内部数据库架构历史主题中。内部数据库架构历史主题仅供连接器使用,不打算供消费应用直接使用。请确保需要接收架构变更通知的应用仅从架构变更主题中消费该信息。
连接器发送到其架构变更主题的消息格式处于孵化阶段,可能会在不另行通知的情况下发生变化。

当发生以下事件时,Debezium 会向架构变更主题发送消息:

  • 你为某张表启用了 CDC。
  • 你为某张表禁用了 CDC。
  • 你按照架构演进流程修改了已启用 CDC 的表的结构。

示例:发送到 SQL Server 连接器架构变更主题的消息

以下示例展示了架构变更主题中的一条消息。该消息包含表架构的逻辑表示。

{
  "schema": {
  ...
  },
  "payload": {
    "source": {
      "version": "3.6.3.Final",
      "connector": "sqlserver",
      "name": "server1",
      "ts_ms": 0,
      "snapshot": "true",
      "db": "testDB",
      "schema": "dbo",
      "table": "customers",
      "change_lsn": null,
      "commit_lsn": "00000025:00000d98:00a2",
      "event_serial_no": null
    },
    "ts_ms": 1588252618953,
    "databaseName": "testDB",
    "schemaName": "dbo",
    "ddl": null,
    "tableChanges": [
      {
        "type": "CREATE",
        "id": "\"testDB\".\"dbo\".\"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 表示在数据库中进行更改的时间。通过比较 payload.source.ts_ms 的值与 payload.ts_ms 的值,可以确定源数据库更新与 Debezium 之间的延迟。

databaseName、schemaName

标识包含该更改的数据库和架构。

ddl

对于 SQL Server 连接器,该字段始终为 null。对于其他连接器,该字段包含导致架构更改的 DDL。SQL Server 连接器无法获取该 DDL。

tableChanges

一个包含一项或多项的数组,其中包含 DDL 命令生成的架构更改。

type

描述更改的类型,值为以下之一:

CREATE

创建表。

ALTER

修改表。

DROP

删除表。

id

被创建、修改或删除的表的完整标识符。

table

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

primaryKeyColumnNames

构成表主键的列的列表。

columns

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

attributes

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

在连接器发送到架构更改主题的消息中,键为包含该架构更改的数据库名称。在以下示例中,payload 字段包含该键:

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

数据变更事件

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

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

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

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

以下列表说明了上一个变更事件示例中的各个字段:

schema(第一个)

第一个 schema 字段是事件键的一部分。它指定了一个 Kafka Connect 模式,用于描述事件键的 payload 部分中包含的内容。换言之,第一个 schema 字段描述的是被更改表的主键结构;如果该表没有主键,则描述的是唯一键的结构。

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

payload(第一个)

第一个 payload 字段是事件键的一部分。它的结构由前一个 schema 字段描述,其中包含被更改行的键。

schema(第二个)

第二个 schema 字段是事件值的一部分。它指定了一个 Kafka Connect 模式,用于描述事件值的 payload 部分中包含的内容。换言之,第二个 schema 描述的是被更改行的结构。通常,该模式包含嵌套模式。

payload(第二个)

第二个 payload 字段是事件值的一部分。它的结构由前一个 schema 字段描述,其中包含被更改行的实际数据。

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

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

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

变更事件键

变更事件的键包含被更改表的键模式以及被更改行的实际键。模式及其对应的载荷都为被更改表的主键(或唯一键约束)中的每一列包含一个字段,列的内容是连接器创建该事件时的值。

请看以下 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
);

示例变更事件键

捕获 customers 表变更的每个变更事件都具有相同的事件键模式。只要 customers 表保持之前的定义,捕获其变更的每个变更事件都具有如下的键结构,以 JSON 表示如下:

{
    "schema": {
        "type": "struct",
        "fields": [
            {
                "type": "int32",
                "optional": false,
                "field": "id"
            }
        ],
        "optional": false,
        "name": "server1.testDB.dbo.customers.Key"
    },
    "payload": {
        "id": 1004
    }
}

以下列表说明了前面的更改事件键示例中的字段:

schema

键中的 schema 部分指定一个 Kafka Connect schema,用于描述键的 payload 部分中包含的内容。

fields

指定 payload 中预期出现的每个字段,包括每个字段的名称、类型以及是否为必填。在此示例中,有一个名为 id 的必填字段,类型为 int32。

optional

指示事件键的 payload 字段中是否必须包含值。在此示例中,键的 payload 中必须包含值。当表没有主键时,键的 payload 字段中的值是可选的。

server1.dbo.testDB.customers.Key

定义键的 payload 结构的 schema 的名称。此 schema 描述了发生更改的表的主键结构。键 schema 的名称格式为 connector-name.database-schema-name.table-name.Key。在此示例中:

server1

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

dbo

发生更改的表所属的数据库 schema。

customers

被更新的表。

payload

包含生成此更改事件的行的键。在此示例中,键包含一个 id 字段,其值为 1004。

当 Debezium 发出更改事件记录时,它会将每条记录的消息键设置为源表的主键或唯一键列的名称。Debezium 必须能够读取这些列才能正常运行。如果在连接器配置中设置了 column.include.list 或 column.exclude.list 属性,请确保设置允许连接器捕获所需的主键或唯一键列。
如果表没有主键或唯一键,则更改事件的键为 null。这是合理的,因为没有主键或唯一键约束的表中的行无法被唯一标识。

更改事件的值

每个 Debezium SQL Server 变更事件值都包含一个 schema 部分和一个 payload 部分。payload 使用信封(envelope)结构,记录每次 INSERT、UPDATE 或 DELETE 操作前后行的状态。

schema 部分描述 payload 部分的 Envelope 结构,包括其嵌套字段。

请考虑用于演示变更事件键示例的同一个示例表:

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
);

针对该表的变更,变更事件的值部分会按事件类型分别说明。

create 事件

每当有行被插入到被捕获表中时,连接器就会发出 create 事件。该事件的负载(payload)在 after 字段中包含新行的完整状态,而 before 字段为 null。

{
  "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": "server1.dbo.testDB.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": "server1.dbo.testDB.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": "change_lsn"
          },
          {
            "type": "string",
            "optional": true,
            "field": "commit_lsn"
          },
          {
            "type": "int64",
            "optional": true,
            "field": "event_serial_no"
          }
        ],
        "optional": false,
        "name": "io.debezium.connector.sqlserver.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": "server1.dbo.testDB.customers.Envelope"
  },
  "payload": {
    "before": null,
    "after": {
      "id": 1005,
      "first_name": "john",
      "last_name": "doe",
      "email": "[email protected]"
    },
    "source": {
      "version": "3.6.3.Final",
      "connector": "sqlserver",
      "name": "server1",
      "ts_ms": 1559729468470,
      "ts_us": 1559729468470000,
      "ts_ns": 1559729468470000000,
      "snapshot": false,
      "db": "testDB",
      "schema": "dbo",
      "table": "customers",
      "change_lsn": "00000027:00000758:0003",
      "commit_lsn": "00000027:00000758:0005",
      "event_serial_no": "1"
    },
    "op": "c",
    "ts_ms": 1559729471739,
    "ts_ms": 1559729471739876,
    "ts_ms": 1559729471739876149
  }
}

以下列表描述了前一个 create 事件值示例中的部分字段:

schema

值的 schema,用于描述值中 payload 的结构。对于某个特定表,连接器生成的每个变更事件中,其值的 schema 都是相同的。

name(server1.dbo.testDB.customers.Value)

在 schema 部分中,每个 name 字段指定 payload 中对应字段的 schema。

server1.dbo.testDB.customers.Value 是 payload 中 before 和 after 字段的 schema。此 schema 专用于 customers 表。

before 和 after 字段的 schema 名称采用 logicalName.database-schemaName.tableName.Value 形式,以确保 schema 名称在数据库中唯一。这意味着,当使用 Avro 转换器时,每个逻辑源中每个表所生成的 Avro schema 都有各自独立的演进过程和历史记录。

name(io.debezium.connector.sqlserver.Source)

io.debezium.connector.sqlserver.Source 是 payload 中 source 字段的 schema。此 schema 专用于 SQL Server 连接器。连接器在生成的所有事件中都会使用它。

name(io.debezium.connector.sqlserver.Envelope)

server1.dbo.testDB.customers.Envelope 是 payload 整体结构的 schema,其中 server1 是连接器名称,dbo 是数据库 schema 名称,customers 是表名。

payload

值的实际数据,即变更事件所要提供的信息。

事件的 JSON 表示形式看起来可能比它所描述的行要大得多,这是因为 JSON 表示形式必须包含消息的 schema 和 payload 部分。不过,通过使用 Avro 转换器,可以显著减小连接器流式传输到 Kafka 主题的消息大小。

before

一个可选字段,用于指定事件发生前行的状态。当 op 字段表示创建(c)时(如此示例所示),before 字段为 null,因为此变更事件针对的是新内容。

after

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

source

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

  • Debezium 版本
  • 连接器类型和名称
  • 数据库和 schema 名称
  • 数据库中发生变更的时间戳
  • 该事件是否属于快照的一部分
  • 包含新行的表名
  • 服务器日志偏移量

op

必填字符串,描述导致连接器生成事件的操作类型。在此示例中,c 表示该操作创建了一行。有效值包括:

c

创建(Create)。

u

更新(Update)。

d

删除(Delete)。

r

读取(仅适用于快照)。

ts_ms、ts_us、ts_ns

可选字段,分别以毫秒、微秒和纳秒格式显示连接器处理事件的时间。在事件消息封装中,该时间基于运行 Kafka Connect 任务的 JVM 中的系统时钟。

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

update 事件

示例 customers 表中更新操作的变更事件值,与该表的 create 事件具有相同的 schema。同样,事件值的 payload 也具有相同的结构。不过,update 事件的事件值 payload 中包含的是不同的值。以下是连接器为 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": "sqlserver",
      "name": "server1",
      "ts_ms": 1559729995937,
      "ts_us": 1559729995937000,
      "ts_ns": 1559729995937000000,
      "snapshot": false,
      "db": "testDB",
      "schema": "dbo",
      "table": "customers",
      "change_lsn": "00000027:00000ac0:0002",
      "commit_lsn": "00000027:00000ac0:0007",
      "event_serial_no": "2"
    },
    "op": "u",
    "ts_ms": 1559729998706,
    "ts_us": 1559729998706318,
    "ts_ns": 1559729998706318547
  }
}

以下列表说明了前面 update 事件值示例中的部分字段:

before

一个可选字段,指定事件发生前行的状态。在 update 事件值中,before 字段包含表中每一列对应的字段,以及数据库提交之前该列中的值。在本示例中,email 的值为 [email protected]。

after

一个可选字段,指定事件发生后行的状态。通过比较 before 和 after 结构,可以确定对该行所做的更新内容。在该示例中,email 的值现在为 [email protected]。

source

必填字段,描述事件的源元数据。source 字段的结构与 create 事件中的相同,但部分值有所不同,例如,示例 update 事件具有不同的偏移量(offset)。源元数据包括:

  • Debezium 版本

  • 连接器类型和名称

  • 数据库和 schema 名称

  • 在数据库中做出更改的时间戳

  • 该事件是否属于快照的一部分

  • 包含新行的表名

  • 服务器日志偏移量

    event_serial_no 字段用于区分具有相同提交 LSN 和变更 LSN 的事件。该字段的值不为 1 的典型情况如下:

  • update 事件的值被设置为 2,因为一次更新会在 SQL Server 的 CDC 变更表中生成两个事件(详情参见源文档)。第一个事件包含旧值,第二个事件包含新值。连接器使用第一个事件中的值来创建第二个事件,并丢弃第一个事件。

  • 当主键被更新时,SQL Server 会发出两个事件:一个是删除具有旧主键值的记录的 delete 事件,另一个是添加具有新主键值的记录的 create 事件。这两个操作共享相同的提交 LSN 和变更 LSN,它们的事件编号分别为 1 和 2。

op

必填字符串,描述操作的类型。在 update 事件值中,op 字段的值为 u,表示该行因更新而发生了变化。

ts_ms、ts_us、ts_ns

可选字段,以毫秒、微秒和纳秒格式显示连接器处理该事件的时间。在事件消息信封中,该时间基于运行 Kafka Connect 任务的 JVM 的系统时钟。

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

更新某一行的主键/唯一键所涉及的列,会改变该行键的值。当键发生变化时,Debezium 会输出三个事件:一个 delete 事件和一个带有该行旧键的墓碑事件,随后是一个带有该行新键的 create 事件。
delete 事件

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

{
  "schema": { ... },
  },
  "payload": {
    "before": {
      "id": 1005,
      "first_name": "john",
      "last_name": "doe",
      "email": "[email protected]"
    },
    "after": null,
    "source": {
      "version": "3.6.3.Final",
      "connector": "sqlserver",
      "name": "server1",
      "ts_ms": 1559730445243,
      "ts_us": 1559730445243000,
      "ts_ns": 1559730445243000000,
      "snapshot": false,
      "db": "testDB",
      "schema": "dbo",
      "table": "customers",
      "change_lsn": "00000027:00000db0:0005",
      "commit_lsn": "00000027:00000db0:0007",
      "event_serial_no": "1"
    },
    "op": "d",
    "ts_ms": 1559730450205,
    "ts_us": 1559730450205387,
    "ts_ns": 1559730450205387492
  }
}

以下列表说明了前面 delete 事件值示例中的部分字段:

before

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

after

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

source

必填字段,用于描述事件的源元数据。在 delete 事件值中,source 字段的结构与同一表的 create 和 update 事件相同,许多 source 字段值也相同。在 delete 事件值中,ts_ms 和 pos 字段的值以及其他一些值可能已发生变化。但 delete 事件值中的 source 字段提供的元数据保持不变:

  • Debezium 版本
  • 连接器类型和名称
  • 数据库和 schema 名称
  • 在数据库中做出更改的时间戳
  • 该事件是否属于快照的一部分
  • 包含该行的表名
  • 服务器日志偏移量

op

必填字符串,用于描述操作类型。op 字段的值为 d,表示该行已被删除。

ts_ms、ts_us、ts_ns

可选字段,分别以毫秒、微秒和纳秒格式显示连接器处理该事件的时间。在事件消息信封中,该时间基于运行 Kafka Connect 任务的 JVM 的系统时钟。

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

SQL Server 连接器的事件设计为可与 Kafka 日志压缩 协同工作。日志压缩允许移除部分较旧的消息,前提是至少保留每个键的最新消息。这使 Kafka 能够回收存储空间,同时确保主题包含完整的数据集,可用于重新加载基于键的状态。

墓碑事件

当某行被删除时,delete 事件值仍可与日志压缩配合使用,因为 Kafka 可以移除所有具有相同键的更早消息。但是,要让 Kafka 移除所有具有相同键的消息,消息值必须为 null。为了实现这一点,Debezium 的 SQL Server 连接器在发出 delete 事件后,会发出一个特殊的墓碑事件,该事件具有相同的键但值为 null。

事务元数据

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

Debezium 接收事务元数据的时机限制

Debezium 仅会为在部署连接器之后发生的事务注册并接收元数据。在部署连接器之前发生的事务,其元数据无法获取。

数据库事务由包含在 BEGIN 和 END 关键字之间的语句块表示。Debezium 会为每个事务中的 BEGIN 和 END 分隔符生成事务边界事件。事务边界事件包含以下字段:

status

BEGIN 或 END。

id

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

ts_ms

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

event_count(针对 END 事件)

该事务产生的事件总数。

data_collections(针对 END 事件)

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

当 Debezium SQL Server 连接器从 Always On 只读副本捕获变更时,它无法可靠地在事务完成后立即发出 END 事务标记。在其他配置中,连接器可以查询 sys.dm_cdc_log_scan_sessions 来确认最近一次日志扫描会话是否已结束,但该动态管理视图(DMV)在 Always On 只读副本上不可用。当 Debezium 从 Always On 只读副本捕获变更时,连接器只有在变更表中检测到来自后续事务的第一个事件后,才会发出事务 END 事件。在流量较少或不频繁的环境中,这种行为可能导致 Debezium 延迟发送 END 标记。

以下示例展示了一条典型的事务边界消息:

示例:SQL Server 连接器事务边界事件

{
  "status": "BEGIN",
  "id": "00000025:00000d08:0025",
  "ts_ms": 1486500577125,
  "event_count": null,
  "data_collections": null
}

{
  "status": "END",
  "id": "00000025:00000d08:0025",
  "ts_ms": 1486500577691,
  "event_count": 2,
  "data_collections": [
    {
      "data_collection": "testDB.dbo.testDB.tablea",
      "event_count": 1
    },
    {
      "data_collection": "testDB.dbo.testDB.tableb",
      "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": "1580390884335172",
  "ts_ns": "1580390884335172574",
  "transaction": {
    "id": "00000025:00000d08:0025",
    "total_order": "1",
    "data_collection_order": "1"
  }
}

数据类型映射

Debezium SQL Server 连接器通过生成结构与该行所在表结构相同的事件,来表示表行数据的更改。每个事件都包含用于表示该行各列值的字段。事件中针对某个操作表示列值的方式,取决于该列的 SQL 数据类型。在事件中,连接器会将每种 SQL Server 数据类型的字段映射到字面类型和语义类型两种类型。

连接器可以将 SQL Server 数据类型同时映射为字面类型和语义类型。

字面类型

描述该值如何使用 Kafka Connect 模式类型进行字面表示,具体包括 INT8、INT16、INT32、INT64、FLOAT32、FLOAT64、BOOLEAN、STRING、BYTES、ARRAY、MAP 和 STRUCT。

语义类型

描述 Kafka Connect 模式如何通过字段的 Kafka Connect 模式名称来体现该字段的含义。

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

基本类型

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

表 5. SQL Server 连接器使用的数据类型映射

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

BIT

BOOLEAN

不适用

TINYINT

INT16

不适用

SMALLINT

INT16

不适用

INT

INT32

不适用

BIGINT

INT64

不适用

REAL

FLOAT32

不适用

FLOAT[(N)]

FLOAT64

不适用

CHAR[(N)]

STRING

不适用

VARCHAR[(N)]

STRING

不适用

TEXT

STRING

不适用

NCHAR[(N)]

STRING

不适用

NVARCHAR[(N)]

STRING

不适用

NTEXT

STRING

不适用

XML

STRING

io.debezium.data.Xml

包含 XML 文档的字符串表示形式

DATETIMEOFFSET[(P)]

STRING

io.debezium.time.ZonedTimestamp

带时区信息的时间戳的字符串表示形式,其中时区为 GMT

其他数据类型映射将在后续章节中介绍。

如果某列存在默认值,该默认值会被传递到对应字段的 Kafka Connect 模式中。更改消息将包含该字段的默认值(除非已显式给出列值),因此通常很少需要从模式中获取默认值。不过,传递默认值有助于在使用 Avro 作为序列化格式并配合 Confluent 模式注册表时满足兼容性规则。

时间值

除了 SQL Server 的 DATETIMEOFFSET 数据类型(包含时区信息)之外,其他时间类型取决于 time.precision.mode 配置属性的取值。当 time.precision.mode 配置属性设置为 adaptive(默认值)时,连接器将根据列的数据类型定义来确定时间类型的字面类型和语义类型,从而让事件精确地表示数据库中的值:

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

DATE

INT32

io.debezium.time.Date

表示自纪元以来的天数。

TIME(0)、TIME(1)、TIME(2)、TIME(3)

INT32

io.debezium.time.Time

表示自午夜起经过的毫秒数,不包含时区信息。

TIME(4)、TIME(5)、TIME(6)

INT64

io.debezium.time.MicroTime

表示自午夜起经过的微秒数,不包含时区信息。

TIME(7)

INT64

io.debezium.time.NanoTime

表示自午夜起经过的纳秒数,不包含时区信息。

DATETIME

INT64

io.debezium.time.Timestamp

表示自纪元起经过的毫秒数,不包含时区信息。

SMALLDATETIME

INT64

io.debezium.time.Timestamp

表示自纪元起经过的毫秒数,不包含时区信息。

DATETIME2(0)、DATETIME2(1)、DATETIME2(2)、DATETIME2(3)

INT64

io.debezium.time.Timestamp

表示自纪元起经过的毫秒数,不包含时区信息。

DATETIME2(4)、DATETIME2(5)、DATETIME2(6)

INT64

io.debezium.time.MicroTimestamp

表示自纪元起经过的微秒数,不包含时区信息。

DATETIME2(7)

INT64

io.debezium.time.NanoTimestamp

表示自纪元起经过的纳秒数,不包含时区信息。

由于 Unix 纪元以来的纳秒数存储在有符号的 INT64 中,io.debezium.time.NanoTimestamp 的可表示范围大约为 1677-09-21T00:12:43Z 到 2262-04-11T23:47:16Z。超出此范围的值——例如存储在 DATETIME2(7) 列中的 9999-12-31 时间终点哨兵值——将无法被表示,并导致连接器在值转换器中抛出 IllegalArgumentException。该异常会按照标准的 event.processing.failure.handling.mode 配置进行处理,因此运维人员可以根据自身需求选择 fail(默认值)、warn 或 skip。

如果某列可能合法地包含超出可表示范围的值,请将其声明为 DATETIME2(0 - 6),以便 Debezium 改为发出 Timestamp(毫秒)或 MicroTimestamp(微秒),或者设置 time.precision.mode=isostring,以 ISO 8601 字符串的形式捕获该值。

当 time.precision.mode 配置属性设置为 connect 时,连接器将使用 Kafka Connect 预定义的逻辑类型。当消费者只了解 Kafka Connect 内置的逻辑类型、无法处理可变精度的时间值时,这可能很有用。另一方面,由于 SQL Server 支持十万分之一微秒的精度,当数据库列的小数秒精度大于 3 时,使用 connect 时间精度模式的连接器生成的事件会导致精度丢失:

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

DATE

INT32

org.apache.kafka.connect.data.Date

表示自纪元以来的天数。

TIME([P])

INT64

org.apache.kafka.connect.data.Time

表示从午夜开始的毫秒数,不包含时区信息。SQL Server 允许 P 的取值范围为 0-7,可存储最高达十分之一微秒的精度,但当 P > 3 时,该模式会导致精度损失。

DATETIME

INT64

org.apache.kafka.connect.data.Timestamp

表示从纪元(epoch)开始的毫秒数,不包含时区信息。

SMALLDATETIME

INT64

org.apache.kafka.connect.data.Timestamp

表示自纪元(epoch)以来的毫秒数,不包含时区信息。

DATETIME2

INT64

org.apache.kafka.connect.data.Timestamp

表示从纪元(epoch)开始的毫秒数,不包含时区信息。SQL Server 允许 P 的取值范围为 0-7,可存储最高达十分之一微秒的精度,但当 P > 3 时,该模式会导致精度损失。

当 time.precision.mode 配置属性设置为 isostring 时,连接器会使用 io.debezium.time.IsoDate、io.debezium.time.IsoTime 和 io.debezium.time.IsoTimestamp 语义类型,将时间值映射为 UTC 时区下符合 ISO-8601 格式的字符串。由于值以字符串形式表示,因此该模式可避免 connect 模式下小数秒精度大于 3 时可能出现的精度损失:

SQL Server 数据类型 字面量类型(架构类型) 语义类型(架构名称)及说明

DATE

STRING

io.debezium.time.IsoDate

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

TIME([P])

STRING

io.debezium.time.IsoTime

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

DATETIME、SMALLDATETIME、DATETIME2

STRING

io.debezium.time.IsoTimestamp

按照 ISO 8601 标准,以 UTC 格式表示时间戳值,例如 2019-07-09T02:28:57.123456Z。

时间戳值

DATETIME、SMALLDATETIME 和 DATETIME2 类型表示不带时区信息的时间戳。此类列会基于 UTC 转换为等价的 Kafka Connect 值。例如,DATETIME2 的值 "2018-06-20 15:13:16.945104" 由 io.debezium.time.MicroTimestamp 表示,其值为 "1529507596945104"。

请注意,运行 Kafka Connect 和 Debezium 的 JVM 所使用的时区不会影响此转换。

十进制值

Debezium 连接器按照 decimal.handling.mode 连接器配置属性的设置来处理十进制数。

decimal.handling.mode=precise

表 6. decimal.handling.mode=precise 时的映射关系

SQL Server 类型 字面量类型(架构类型) 语义类型(架构名称)

NUMERIC[(P[,S])]

BYTES

org.apache.kafka.connect.data.Decimal

scale 架构参数包含一个整数,表示小数点向左移动的位数。

DECIMAL[(P[,S])]

BYTES

org.apache.kafka.connect.data.Decimal

scale 架构参数包含一个整数,表示小数点移动的位数。

SMALLMONEY

BYTES

org.apache.kafka.connect.data.Decimal

scale 架构参数包含一个整数,表示小数点移动的位数。

MONEY

BYTES

org.apache.kafka.connect.data.Decimal

scale 架构参数包含一个整数,表示小数点移动的位数。

decimal.handling.mode=double

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

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

decimal.handling.mode=string

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

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

设置 SQL Server

要让 Debezium 捕获 SQL Server 表的变更事件,首先必须由具备相应权限的 SQL Server 管理员运行查询,在数据库上启用 CDC。然后,管理员还需要为你希望 Debezium 捕获的每个表启用 CDC。

启用 CDC 后,它会捕获已为启用 CDC 的表提交的所有 INSERT、UPDATE 和 DELETE 操作。随后,Debezium 连接器即可捕获这些事件,并将其发送到 Kafka 主题中。

加密连接

默认情况下,到 Microsoft SQL Server 的 JDBC 连接使用 SSL 加密。你可以使用 JDBC 传递属性配置 SSL 加密设置,以强制执行证书验证,或在 SQL Server 不使用 SSL 的环境中禁用加密。

如果 SQL Server 数据库未启用 SSL,或者你希望不使用 SSL 连接到该数据库,可以将连接器配置中 driver.encrypt 属性的值设置为 false 来禁用 SSL。

操作步骤

  • 要强制使用有效证书进行 SSL 验证,请设置以下 JDBC 传递属性来配置连接:

    {
      "driver.encrypt": true,
      "driver.trustServerCertificate": false,
      "driver.trustStore": "path/to/trust-store",
      "driver.trustStorePassword": "password-for-trust-store"
    }
  • 在非生产环境中,如果 SQL Server 配置为使用你信任的已知自签名证书,可以通过设置以下属性来启用加密并跳过证书路径验证:

    {
      "driver.encrypt": true,
      "driver.trustServerCertificate": true
    }
不要在生产环境中启用对服务器证书的自动信任。如果禁用验证,连接器会跳过证书链和主机名验证,这可能使连接容易受到中间人攻击。仅在测试和开发环境中,为方便起见才使用显式信任配置

在 SQL Server 数据库上启用 CDC

要为某个表启用 CDC,必须先为 SQL Server 数据库启用它。SQL Server 管理员通过运行系统存储过程来启用 CDC。系统存储过程可以使用 SQL Server Management Studio 运行,也可以使用 Transact-SQL 运行。

前提条件

  • 你是该 SQL Server 的 sysadmin 固定服务器角色成员。
  • 你是该数据库的 db_owner。
  • SQL Server Agent 正在运行。
SQL Server 的 CDC 功能仅处理用户创建的表中发生的更改。不能在 SQL Server master 数据库上启用 CDC。

步骤

  1. 在 SQL Server Management Studio 中,从 View(视图)菜单点击 Template Explorer(模板资源管理器)。

  2. 在 Template Browser(模板浏览器)中,展开 SQL Server Templates(SQL Server 模板)。

  3. 展开 Change Data Capture > Configuration,然后点击 Enable Database for CDC(为 CDC 启用数据库)。

  4. 在模板中,将 USE 语句中的数据库名替换为你想要为 CDC 启用的数据库的名称。

  5. 运行存储过程 sys.sp_cdc_enable_db 以将该数据库启用为 CDC 数据库。

    数据库启用 CDC 后,会创建一个名为 cdc 的架构,以及一个 CDC 用户、元数据表和其他系统对象。

    以下示例展示了如何为数据库 MyDB 启用 CDC:

    示例:为 SQL Server 数据库启用 CDC 模板

    USE MyDB
    GO
    EXEC sys.sp_cdc_enable_db
    GO

在 SQL Server 表上启用 CDC

SQL Server 管理员必须在希望 Debezium 捕获的源表上启用变更数据捕获(CDC)。数据库本身必须已经启用 CDC。要在表上启用 CDC,SQL Server 管理员需要对该表运行存储过程 sys.sp_cdc_enable_table。该存储过程可以通过 SQL Server Management Studio 或 Transact-SQL 来执行。对于希望捕获的每张表,都必须启用 SQL Server CDC。

前置条件

  • SQL Server 数据库已启用 CDC。
  • SQL Server Agent 正在运行。
  • 你是该数据库的 db_owner 固定数据库角色成员。

操作步骤

  1. 在 SQL Server Management Studio 中,从 View 菜单点击 Template Explorer。

  2. 在 Template Browser 中,展开 SQL Server Templates。

  3. 展开 Change Data Capture > Configuration,然后点击 Enable Table Specifying Filegroup Option。

  4. 在模板中,将 USE 语句中的数据库名替换为你要捕获的数据库的名称。

  5. 将 source_schema、source_name、role_name 和 filegroup_name 参数分别替换为表的架构名、表名、Debezium 使用的数据库用户或角色名,以及文件组名。你可以按 CTRL+SHIFT+M 打开参数替换对话框。

  6. 运行存储过程 sys.sp_cdc_enable_table。

    以下示例展示了如何为表 MyTable 启用 CDC:

    示例:为 SQL Server 表启用 CDC

    USE MyDB
    GO
    
    EXEC sys.sp_cdc_enable_table
    @source_schema = N'dbo',
    @source_name   = N'MyTable',
    @role_name     = N'MyRole',
    @filegroup_name = N'MyDB_CT',
    @supports_net_changes = 0
    GO

@source_name

指定要捕获的表的名称。

@role_name

指定一个角色 MyRole,你可以将需要获授权限的用户添加到该角色中,以便授予这些用户对源表被捕获列的 SELECT 权限。sysadmin 或 db_owner 角色中的用户同样可以访问指定的变更表。如果 @role_name 的值被显式设置为 NULL,则不使用任何角色来限制对被捕获信息的访问。

@filegroup_name

指定 SQL Server 为被捕获的表存放变更表所在的 filegroup。所指定的 filegroup 必须已存在。最好不要将变更表与源表放置在同一个 filegroup 中。

验证用户对 CDC 表的访问权限

SQL Server 管理员可以运行系统存储过程来查询数据库或表,以获取其 CDC 配置信息。这些存储过程可以通过 SQL Server Management Studio 运行,也可以通过 Transact-SQL 运行。

前提条件

  • 你拥有对该捕获实例所有被捕获列的 SELECT 权限。db_owner 数据库角色的成员可以查看所有已定义捕获实例的信息。
  • 你拥有查询内容所涉及表信息中定义的任何控制角色的成员身份。

操作步骤

  1. 在 SQL Server Management Studio 中,从 View(视图) 菜单点击 Object Explorer(对象资源管理器)。

  2. 在对象资源管理器中,展开 Databases(数据库),然后展开你的数据库对象,例如 MyDB。

  3. 展开 Programmability(可编程性) > Stored Procedures(存储过程) > System Stored Procedures(系统存储过程)。

  4. 运行 sys.sp_cdc_help_change_data_capture 存储过程来查询该表。

    查询不应返回空结果。

    以下示例在数据库 MyDB 上运行存储过程 sys.sp_cdc_help_change_data_capture:

    示例:查询表的 CDC 配置信息

    USE MyDB;
    GO
    EXEC sys.sp_cdc_help_change_data_capture
    GO

该查询返回数据库中启用了 CDC 且包含调用者有权访问的变更数据的每个表的配置信息。如果结果为空,请验证用户是否拥有访问捕获实例和 CDC 表的权限。

Azure 上的 SQL Server

你可以将 Debezium SQL Server 连接器与 Azure 上的 SQL Server 配合使用。在 Azure SQL 数据库上启用 CDC,并配置连接器以连接到 Azure SQL 终结点。

关于如何为 Azure 上的 SQL Server 配置 CDC 并将其与 Debezium 配合使用,请参阅此示例。

SQL Server Always On

SQL Server 连接器可以从 Always On 只读副本中捕获变更。

前提条件

  • 变更数据捕获已在主节点上配置并启用。SQL Server 不支持直接在副本上使用 CDC。

  • 配置选项 driver.applicationIntent 被设置为 ReadOnly。这是 SQL Server 的要求。当 Debezium 检测到该配置选项时,会采取以下操作:

    • 将 snapshot.isolation.mode 设置为 snapshot,这是只读副本唯一支持的事务隔离模式。
    • 在流式查询循环的每次执行中提交(只读)事务,这是获取 CDC 数据最新视图所必需的。

SQL Server 捕获作业代理配置对服务器负载和延迟的影响

SQL Server 捕获作业代理处理来自事务日志的变更事件记录,其轮询频率直接影响数据库的 CPU 负载以及 Debezium 可用变更事件的延迟。调整代理参数有助于你在工作负载中平衡性能与响应能力。

当数据库管理员为源表启用变更数据捕获时,捕获作业代理开始运行。代理从事务日志中读取新的变更事件记录,并将这些事件记录复制到变更数据表中。从变更在源表中提交的时刻,到该变更出现在对应变更表中的时刻之间,始终存在一个很小的延迟间隔。这个延迟间隔代表了变更在源表中发生与变更可被 Debezium 流式传输到 Apache Kafka 之间的时间差。

理想情况下,对于必须快速响应数据变化的应用,你需要在源表和变更表之间保持紧密同步。你可能会认为,让捕获代理以尽可能快的速度持续处理变更事件,能够提高吞吐量并降低延迟——在事件发生后尽快(接近实时地)将新事件记录填充到变更表中。然而,情况未必如此。追求更即时的同步是要付出性能代价的。捕获作业代理每次查询数据库以获取新事件记录时,都会增加数据库主机上的 CPU 负载。服务器上额外的负载可能对数据库的整体性能产生负面影响,并可能降低事务效率,尤其是在数据库使用高峰期。

监控数据库指标非常重要,这样你才能知道数据库是否已经达到了服务器无法再支撑捕获代理当前活动水平的程度。如果你发现性能问题,可以修改 SQL Server 捕获代理的某些设置,以帮助平衡数据库主机上的整体 CPU 负载与可接受程度的延迟。

SQL Server 捕获作业代理配置参数

SQL Server 的 msdb.dbo.cdc_jobs 表定义了控制捕获作业代理行为的参数。你可以将这些参数与 sys.sp_cdc_change_job 存储过程配合使用,调整轮询频率和事务批处理大小,从而平衡数据库服务器的负载与事件延迟。

如果你在运行捕获作业代理时遇到性能问题,可以通过执行 sys.sp_cdc_change_job 存储过程并传入新的值来调整捕获作业设置,从而降低 CPU 负载。

关于如何配置 SQL Server 捕获作业代理参数的具体指导超出了本文档的范围。

以下参数对于配置与 Debezium SQL Server 连接器配合使用的捕获代理行为最为重要:

pollinginterval

  • 指定捕获代理在两次日志扫描循环之间等待的秒数。
  • 值越大,数据库主机的负载越低,但延迟越高。
  • 值为 0 表示两次扫描之间不等待。
  • 默认值为 5。

maxtrans

  • 指定每次日志扫描循环中要处理的最大事务数。捕获作业处理完指定数量的事务后,会暂停由 pollinginterval 指定的时长,然后才开始下一次扫描。
  • 值越小,数据库主机的负载越低,但延迟越高。
  • 默认值为 500。

maxscans

  • 指定捕获作业在捕获完整数据库事务日志内容时可尝试的扫描周期次数上限。如果 continuous 参数设置为 1,作业会按照 pollinginterval 指定的时间长度暂停,然后恢复扫描。
  • 较小的值会降低对数据库主机的负载,但会增加延迟。
  • 默认值为 10。

其他资源

部署

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

前提条件

操作步骤

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

如果你使用的是不可变容器,请参阅适用于 Apache Kafka 和 Kafka Connect 的 Debezium 容器镜像。你可以从 Docker Hub 拉取官方的 Linux 版 Microsoft SQL Server 容器镜像。

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

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

SQL Server 连接器配置示例

通过在 JSON 文件中设置属性来配置 Debezium SQL Server 连接器,该文件指定数据库连接、主题前缀、表捕获列表以及 schema history 主题。你还可以选择性地过滤、掩码或截断列值,以控制连接器捕获和发出哪些数据。

{
    "name": "inventory-connector",
    "config": {
        "connector.class": "io.debezium.connector.sqlserver.SqlServerConnector",
        "database.hostname": "192.168.99.100",
        "database.port": "1433",
        "database.user": "sa",
        "database.password": "Password!",
        "database.names": "testDB1,testDB2",
        "topic.prefix": "fullfillment",
        "table.include.list": "dbo.customers",
        "schema.history.internal.kafka.bootstrap.servers": "kafka:9092",
        "schema.history.internal.kafka.topic": "schemahistory.fullfillment",
        "driver.trustStore": "path/to/trust-store",
        "driver.trustStorePassword": "password-for-trust-store"
    }
}

以下列表说明了前面连接器配置示例中的部分字段:

name

在向 Kafka Connect 服务注册时,我们为连接器指定的名称。

connector.class

此 SQL Server 连接器类的名称。

database.hostname

SQL Server 实例的地址。

database.port

SQL Server 实例的端口号。

database.user

SQL Server 用户名。

database.password

SQL Server 用户的密码。

database.names

要捕获其更改的数据库名称。

topic.prefix

SQL Server 实例/集群的主题前缀,它构成一个命名空间,并用于该连接器写入的所有 Kafka 主题名称、Kafka Connect 模式名称,以及使用 Avro 转换器 时相应 Avro 模式的命名空间中。

table.include.list

Debezium 应捕获其更改的所有表的列表。

schema.history.internal.kafka.bootstrap.servers

此连接器用于将数据库架构历史主题中的 DDL 语句写入和恢复的 Kafka 代理列表。

schema.history.internal.kafka.topic

连接器写入和恢复 DDL 语句的数据库架构历史主题的名称。该主题仅供内部使用,消费者不应使用它。

driver.trustStore

存储服务器签名者证书的 SSL 信任库路径。除非禁用了数据库加密(driver.encrypt=false),否则此属性为必填项。

driver.trustStorePassword

SSL 信任库密码。除非禁用了数据库加密(driver.encrypt=false),否则此属性为必填项。

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

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

  • 连接到 SQL Server 数据库。
  • 读取事务日志。
  • 将更改事件记录到 Kafka 主题。

添加连接器配置

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

前提条件

步骤

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

结果

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

连接器属性

Debezium SQL Server 连接器提供了众多配置属性,可用于为你的应用调校出合适的连接器行为。其中许多属性都有默认值。

属性信息按如下方式组织:

必需的 Debezium SQL Server 连接器配置属性

必需的配置属性用于建立数据库连接,并定义连接器的基本标识与作用范围。在连接器启动之前,你必须设置这些属性。

除非存在默认值,否则下列配置属性均为必需。

属性默认值说明
name无默认值连接器的唯一名称。使用相同的名称再次注册将会失败。(所有 Kafka Connect 连接器都要求此属性。)

connector.class

无默认值

连接器所使用的 Java 类名。对于 SQL Server 连接器,始终使用 io.debezium.connector.sqlserver.SqlServerConnector。

tasks.max

1

指定连接器可用于从数据库实例采集数据的最大任务数。如果 database.names 列表包含多个元素,可以将该属性的值增大,但不能超过列表中元素的数量。

database.hostname

无默认值

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

database.port

1433

SQL Server 数据库服务器的整数端口号。如果同时指定了 database.port 和 database.instance,则 database.instance 会被忽略。更多详情请参阅 SQL Server 的 JDBC 驱动程序文档。

driver.authentication

不适用

指定 Microsoft SQL Server JDBC 驱动程序所使用的身份验证模式。该属性对应于驱动程序的 authentication 连接属性。支持的值包括 ActiveDirectoryManagedIdentity。在连接 Azure SQL Database 时,可使用该属性配置 Microsoft Entra 身份验证。遗留值 ActiveDirectoryMSI 可用于连接 Azure SQL Database。

database.user

无默认值

连接 SQL Server 数据库服务器时使用的用户名。使用 Kerberos 身份验证时可以省略该属性,Kerberos 身份验证可通过透传属性进行配置。

database.password

无默认值

连接 SQL Server 数据库服务器时使用的密码。使用 Microsoft Entra 身份验证时为可选项。

database.instance

无默认值

指定 SQL Server 命名实例的实例名称。如果同时指定了 database.port 和 database.instance,则 database.instance 会被忽略。更多详情请参阅 SQL Server JDBC 驱动程序文档。

database.names

无默认值

以逗号分隔的 SQL Server 数据库名称列表,用于从中捕获变更。

topic.prefix

无默认值

主题前缀,为希望 Debezium 捕获的 SQL Server 数据库服务器提供命名空间。该前缀应在所有其他连接器中保持唯一,因为它是从此连接器接收记录的所有 Kafka 主题名称的前缀。数据库服务器的逻辑名称中只能使用字母、数字、连字符、点和下划线。

不要更改此属性的值。如果更改了名称值,重启后连接器将不再继续向原有主题发送事件,而是将后续事件发送到基于新值命名的主题。连接器也将无法恢复其数据库架构历史主题。

schema.include.list

无默认值

一个可选的、以逗号分隔的正则表达式列表,用于匹配你想要捕获变更的架构名称。任何未包含在 schema.include.list 中的架构名称,其变更都不会被捕获。默认情况下,连接器会捕获所有非系统架构的变更。

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

schema.exclude.list

无默认值

一个可选的、以逗号分隔的正则表达式列表,用于匹配你不希望捕获其变更的 schema 名称。任何名称未包含在 schema.exclude.list 中的 schema,其变更都会被捕获,系统 schema 除外。

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

table.include.list

无默认值

一个可选的、以逗号分隔的正则表达式列表,用于匹配你希望 Debezium 捕获的表的完全限定标识符。默认情况下,连接器会为指定的 schema 捕获所有非系统表。设置此属性后,连接器将仅捕获指定表的变更。每个标识符的形式为 schemaName.tableName。

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

table.exclude.list

无默认值

一个可选的、以逗号分隔的正则表达式列表,用于匹配你希望排除在捕获之外的表的完全限定标识符。Debezium 会捕获所有未包含在 table.exclude.list 中的表。每个标识符的形式为 schemaName.tableName。

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

column.include.list

空字符串

一个可选的、以逗号分隔的正则表达式列表,用于匹配应包含在变更事件消息值中的列的完全限定名称。列的完全限定名称形式为 schemaName.tableName.columnName。

Debezium 为某张表发出的每条更改事件记录都包含一个事件键,该事件键包含该表主键或唯一键中每个列的字段。为确保事件键正确生成,如果设置了此属性,请务必显式列出所有被捕获表的主键列。

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

column.exclude.list

空字符串

一个可选的、以逗号分隔的正则表达式列表,用于匹配应从更改事件消息值中排除的列的完全限定名称。列的完全限定名称的形式为 schemaName.tableName.columnName。请注意,即使主键列被从值中排除,也始终会包含在事件键中。

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

skip.messages.without.change

false

指定当所包含的列没有发生变化时是否跳过发布消息。如果按照 column.include.list 或 column.exclude.list 属性所包含的列没有任何变化,这实际上会过滤掉这些消息。

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

不适用

一个可选的、以逗号分隔的正则表达式列表,用于匹配字符类型列的完全限定名。列的完全限定名格式为 *`<schemaName>*.*<tableName>*.*<columnName>`。为匹配列名,Debezium 会将你指定的正则表达式作为 anchored(锚定) 正则表达式来应用。也就是说,指定的表达式会与列的整个名称字符串进行匹配,而不会匹配列名中可能存在的子字符串。在生成的变更事件记录中,指定列的值会被替换为化名(pseudonym)。

化名由应用所指定的 hashAlgorithm(哈希算法)和 salt(盐值)后得到的哈希值组成。根据所使用的哈希函数,在将列值替换为化名的同时仍可保持引用完整性。受支持的哈希函数见 MessageDigest 章节(Java 密码体系结构标准算法名称文档)。

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

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

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

根据所使用的 hashAlgorithm、所选的 salt 以及实际数据集的不同,生成的数据集可能无法被完全遮蔽。

如果要在不同位置或系统中对同一值进行哈希,应使用哈希策略版本 2,以确保结果一致。

time.precision.mode

adaptive

时间、日期和时间戳可以采用不同类型的精度表示,包括:adaptive(默认值)按照数据库列的类型,使用毫秒、微秒或纳秒精度精确捕获数据库中的时间和时间戳值;connect 则始终使用 Kafka Connect 内置的 Time、Date 和 Timestamp 表示形式来表示时间和时间戳值,无论数据库列的精度如何,均使用毫秒精度。有关更多信息,请参阅时间值。

decimal.handling.mode

precise

指定连接器应如何处理 DECIMAL 和 NUMERIC 列的值。可设置为以下值之一:

precise

(默认值)使用 java.math.BigDecimal 值精确表示列数据,在变更事件中以二进制形式表示。

double

使用 double 值表示列值,可能会损失精度,但更易于使用。

string

将列值编码为格式化的字符串,易于使用,但会丢失原始类型的语义信息。

include.schema.changes

true

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

tombstones.on.delete

true

控制 delete 事件之后是否跟随一个墓碑事件。可设置为以下值之一:

true

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

false

执行删除操作后,连接器仅发出一个 delete 事件。

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

column.truncate.to.length.chars

n/a

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

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

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

column.mask.with.length.chars

n/a

列的完全限定名称的形式为 schemaName.tableName.columnName。

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

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

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

column.propagate.source.type

n/a

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

  • __debezium.source.column.type
  • __debezium.source.column.length
  • __debezium.source.column.scale

这些参数分别传递列的原始类型名称和长度(对于可变宽度类型)。启用连接器发送这些额外数据,有助于在接收端数据库中正确设置特定数值列或字符列的大小。

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

datatype.propagate.source.type

无

一个可选的、以逗号分隔的正则表达式列表,用于指定数据库中为列定义的数据类型的完全限定名。设置此属性后,对于数据类型匹配的列,连接器会发出事件记录,其 schema 中包含以下额外字段:

  • __debezium.source.column.type
  • __debezium.source.column.length
  • __debezium.source.column.scale

这些参数分别传递列的原始类型名称和长度(对于可变宽度类型)。启用连接器发送这些额外数据,有助于在接收端数据库中正确设置特定数值列或字符列的大小。

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

有关 SQL Server 特有的数据类型名称列表,请参阅 SQL Server 数据类型映射。

message.key.columns

无

一个表达式列表,用于指定连接器为发布到指定表对应的 Kafka 主题的变更事件记录构造自定义消息键时所使用的列。

默认情况下,Debezium 使用表的主键列作为它所发出记录的消息键。若要取代默认设置,或为缺少主键的表指定键,你可以基于一个或多个列配置自定义消息键。

要为表建立自定义消息键,请列出该表,后跟用作消息键的列。每个列表项遵循以下格式:

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

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

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

<schemaName>.<tableName>

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

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

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

对于 inventory.customer 表,pk1 和 pk2 两列被指定为消息键;对于任意模式(schema)中的 purchaseorders 表,pk3 和 pk4 两列作为消息键。

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

binary.handling.mode

bytes

指定变更事件中二进制(binary、varbinary)列的表示方式,包括:bytes 将二进制数据表示为字节数组(默认值),base64 将二进制数据表示为 base64 编码的字符串,base64-url-safe 将二进制数据表示为 base64-url-safe 编码的字符串,hex 将二进制数据表示为十六进制(base16)编码的字符串。

schema.name.adjustment.mode

none

指定如何调整模式名称(schema name),以使其与连接器所使用的消息转换器兼容。可选设置如下:

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

field.name.adjustment.mode

none

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

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

有关更多信息,请参阅 Avro 命名。

SQL Server 连接器高级配置属性

以下高级配置属性已具备良好的默认值,在大多数情况下均可直接使用,因此很少需要在连接器配置中指定。

属性 默认值 说明

converters

无默认值

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

isbn

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

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

<converterSymbolicName>.type

例如:

isbn.type: io.debezium.test.IsbnConverter

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

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

snapshot.mode

initial

用于对被捕获表的结构以及(可选的)数据执行初始快照的模式。快照完成后,连接器将继续从数据库的重做日志中读取变更事件。支持以下取值:

always

每次连接器启动时都执行快照。快照完成后,连接器开始流式传输后续数据库变更的事件记录。

initial

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

initial_only

连接器执行数据库快照,然后在流式传输任何变更事件记录之前停止,不允许捕获任何后续的变更事件。

schema_only

已弃用,请参阅 no_data。

no_data

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

recovery

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

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

when_needed

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

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

when_needed_no_data

连接器启动后,其快照行为与 when_needed 选项类似。与 no_data 选项一样,它不会创建 READ 事件来表示连接器启动时的数据集(第 7.b 步)。

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.name

无默认值

如果 snapshot.mode 设置为 custom,请使用此设置来指定自定义实现的名称,该名称由 'io.debezium.spi.snapshot.Snapshotter' 接口中定义的 name() 方法提供。连接器重启后,Debezium 会调用指定的自定义实现,以确定是否执行快照。有关更多信息,请参阅自定义快照器 SPI。

snapshot.locking.mode

exclusive

控制连接器是否持有表锁以及持锁的时长。表锁可防止在连接器执行快照期间发生某些类型的表更改操作。您可以设置以下值:

exclusive

控制连接器在执行架构快照时如何对表加锁,此时 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。

snapshot.query.mode

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。

snapshot.include.collection.list

table.include.list 中指定的所有表

一个可选的、以逗号分隔的正则表达式列表,用于匹配要包含在快照中的表的完全限定名称(<dbName>.<schemaName>.<tableName>)。指定的项必须在连接器的 table.include.list 属性中列出。只有当连接器的 snapshot.mode 属性设置为 no_data 以外的值时,此属性才会生效。此属性不影响增量快照的行为。

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

snapshot.isolation.mode

repeatable_read

用于控制采用哪种事务隔离级别,以及连接器对指定捕获的表加锁时长的模式。支持以下值:

  • read_uncommitted
  • read_committed
  • repeatable_read
  • snapshot
  • exclusive(exclusive 模式使用可重复读隔离级别,但它会对所有待读取的表加排他锁)。

snapshot、read_committed 和 read_uncommitted 模式不会阻止其他事务在初始快照期间更新表中的行。exclusive 和 repeatable_read 模式则会阻止并发更新。

模式的选择也会影响数据一致性。只有 exclusive 和 snapshot 模式能保证完全的一致性,即初始快照与流式日志构成一条线性历史。对于 repeatable_read 和 read_committed 模式,可能会出现例如某条新增记录出现两次的情况——一次在初始快照中,一次在流式阶段中。尽管如此,这种一致性级别对于数据镜像来说通常是足够的。对于 read_uncommitted,则完全没有数据一致性保证(部分数据可能丢失或损坏)。

event.processing.failure.handling.mode

fail

指定连接器在处理事件时应对异常的反应方式。fail 会抛出异常(指示问题事件的偏移量),导致连接器停止。warn 会跳过问题事件,并记录该问题事件的偏移量。skip 会跳过问题事件。

poll.interval.ms

500(0.5 秒)

一个正整数值,指定连接器在检查数据库中是否有新变更事件之前需要等待的毫秒数。

指定的值会影响 heartbeat.interval.ms 的行为。连接器只能在指定的轮询周期内发出心跳消息。

为防止此设置延迟心跳消息的发出,请将其设置为小于或等于 heartbeat.interval.ms 的值。

max.queue.size

8192

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

max.queue.size.in.bytes

0

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

max.batch.size

2048

一个正整数值,指定该连接器每次迭代处理的最大事件批次大小。

heartbeat.interval.ms

0

指定一个以毫秒为单位的间隔,用于确定连接器向 Kafka 心跳主题发送消息的频率,无论数据库中是否发生变更。默认情况下,连接器不发送心跳消息。

设置此属性有助于确认连接器是否仍在接收来自数据库的变更事件。在被捕获的表长时间没有发生变化的数据库中,这一点尤为重要。当数据库频繁出现被捕获的表长时间没有变更的情况时,尽管连接器仍会像往常一样从事务日志中读取数据,但它极少向 Kafka 提交偏移量值。其结果是,在连接器重启后,由于偏移量值已过期,连接器必须发送大量变更事件。

相比之下,当你配置连接器定期发送心跳消息时,它可以更频繁地更新 Kafka 中的偏移量。由于 Kafka 中的偏移量值始终保持最新,连接器重启后需要重新发送的变更事件就会更少。

心跳消息仅在轮询周期内发出。也就是说,在 Debezium 环境中,发送心跳消息的实际间隔由 heartbeat.interval.ms 和 poll.interval.ms 属性的设置共同决定。发送心跳消息的实际频率以这两个值中较小的那个为准。为避免心跳消息发送延迟而降低其有效性,请将此属性设置为大于或等于 poll.interval.ms 的值。例如,如果你将 poll.interval.ms 设置为 100,则应将 heartbeat.interval.ms 设置为 5000。

heartbeat.action.query

无默认值

指定连接器在发送心跳消息时,会在源数据库上执行的查询。

这对于防止从低流量数据库捕获变更时偏移量过期非常有用。在低流量数据库中创建一个心跳表,并将此属性设置为向该表插入记录的语句,例如:

INSERT INTO test_heartbeat_table (text) VALUES ('test_heartbeat')

这样,连接器就能从低流量数据库接收变更并确认其 LSN,从而防止偏移量过期。

snapshot.delay.ms

无默认值

连接器在启动后、执行快照前应等待的时间间隔(以毫秒为单位)。可用于避免在集群中启动多个连接器时快照被打断的情况,这种情况可能导致连接器重新均衡。

statistics.metrics.enabled

true

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

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

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

streaming.delay.ms

0

指定连接器在完成快照后延迟多久再开始流式处理过程(以毫秒为单位)。设置延迟间隔有助于防止在快照完成后、流式处理过程开始之前立即发生故障时,连接器重新启动快照。请将延迟值设置为高于 Kafka Connect 工作程序上配置的 offset.flush.interval.ms 属性的值。

snapshot.fetch.size

2000

指定在执行快照时,每次从每张表一次性读取的最大行数。连接器将按此大小分多批次读取表内容。默认值为 2000。

query.fetch.size

无默认值

指定给定查询在每次数据库往返中获取的行数。默认使用 JDBC 驱动程序的默认获取大小。

snapshot.lock.timeout.ms

10000

一个整数值,指定在执行快照时获取表锁的最长等待时间(以毫秒为单位)。如果在该时间间隔内无法获取表锁,快照将会失败(另请参阅快照)。设置为 0 时,连接器在无法获取锁时会立即失败。值 -1 表示无限等待。

snapshot.select.statement.overrides

无默认值

指定连接器针对哪些表使用自定义 SELECT 语句来决定将哪些行包含在快照中。

此属性仅影响快照。它不适用于连接器在流式处理阶段从日志中读取的事件。

此属性由两部分组成,两者协同工作:

主属性

以 <databaseName>.<schemaName>.<tableName> 格式表示的完全限定表名的逗号分隔列表。该列表标识了你要为其指定自定义快照查询的表。

例如:

"snapshot.select.statement.overrides": "mydb.inventory.products,mydb.customers.orders"

如果完全限定的数据库名、模式名或表名包含特殊字符,例如空格、方括号([ 或 ])或句点字符(.),请将该字符串用双引号括起来,以防止连接器将这些特殊字符解释为分隔符。如果完全限定的表名不包含空格或特殊字符,则表名两侧的双引号是可选的。

次级属性

对于你在主属性中列出的每个表,你都必须定义相应的 snapshot.select.statement.overrides.<databaseName>.<schemaName>.<tableName> 属性,用于指定在快照期间要运行的自定义 SELECT 语句。该 SELECT 语句决定从该表中选取哪些行包含到快照中。

例如,要为 mydb.customer.orders 表指定一条 SELECT 语句,请添加以下属性:

snapshot.select.statement.overrides.mydb.customers.orders

如果某个表列在主属性中,但缺少相应的次属性,连接器会记录警告日志,并对该表使用默认的快照行为。

配置示例

以下示例展示了如何配置 snapshot.select.statement.overrides 属性,以便对 mydb.customers.orders 表执行快照,且仅包含未被软删除的记录;也就是说,软删除字段 delete_flag 的值为 0。

"snapshot.select.statement.overrides": "mydb.customer.orders",
"snapshot.select.statement.overrides.mydb.customer.orders": "SELECT * FROM mydb.customers.orders WHERE delete_flag = 0 ORDER BY id DESC"

source.struct.version

v2

CDC 事件中 source 块的 Schema 版本;为了在所有连接器之间统一暴露的结构,Debezium 0.10 对 source 块的结构做了一些破坏性更改。将此选项设置为 v1 可以生成早期版本中使用的结构。请注意,不建议使用此设置,并且计划在未来的 Debezium 版本中移除。

provide.transaction.metadata

false

设置为 true 时,Debezium 会生成带有事务边界的事件,并在数据事件信封中补充事务元数据。

retriable.restart.connector.wait.ms

10000(10 秒)

发生可重试错误后,重启连接器之前需要等待的毫秒数。

skipped.operations

t

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

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

如果不想让连接器跳过任何操作,请将值设置为 none。由于 Debezium SQL Server 连接器不支持 truncate 变更事件,因此将默认值 t 与将值设置为 none 的效果相同。

signal.data.collection

无默认值

用于向连接器发送信号的数据集合的完全限定名称。集合名称区分大小写。

请使用以下格式指定集合名称:

<databaseName>.<schemaName>.<tableName>

对于多数据库部署,如果 tasks.max > 1,你可以通过提供以逗号分隔的完全限定信号表名列表,为每个数据库配置一个信号表。例如:

signal.data.collection=db1.dbo.debezium_signal,db2.dbo.debezium_signal

为每个数据库指定专用的信号表,可确保每个任务都能正确处理其所分配数据库的信号。如果为多任务连接器仅配置一个信号表,则可能无法正确处理所有数据库的信号。

在 Debezium SQL Server 连接器中使用多个信号表是一项技术预览功能。技术预览功能不受 Red Hat 生产服务级别协议(SLA)支持,功能可能尚不完整。Red Hat 不建议在生产环境中使用这些功能。这些功能让客户可以提前访问即将推出的产品功能,从而在开发过程中测试功能并提供反馈。有关 Red Hat 技术预览功能支持范围的更多信息,请参阅技术预览功能支持范围。

signal.enabled.channels

source

为连接器启用的信号通道名称列表。默认情况下,可以使用以下通道:

notification.enabled.channels

无默认值

为连接器启用的通知通道名称列表。默认情况下,可以使用以下通道:

incremental.snapshot.allow.schema.changes

false

允许在增量快照期间进行模式变更。启用后,连接器会在增量快照期间检测模式变更,并重新选择当前数据块,以避免锁定 DDL。

请注意,主键的变更不受支持,在增量快照期间执行可能导致结果不正确。另一个限制是,如果模式变更仅涉及列的默认值,那么在事务日志流中处理到该 DDL 之前,该变更不会被检测到。这不会影响快照事件的值,但快照事件的模式可能包含过期的默认值。

incremental.snapshot.chunk.size

1024

连接器在增量快照分块期间获取并读入内存的最大行数。增大分块大小可提高效率,因为快照执行的查询次数更少、单次查询的数据量更大。不过,较大的分块大小也需要更多内存来缓冲快照数据。请根据环境调整分块大小,以获得最佳性能。

incremental.snapshot.watermarking.strategy

insert_insert

指定连接器在增量快照期间使用的水印机制,用于对可能先被增量快照捕获、随后在流式传输恢复后又被再次捕获的事件进行去重。你可以指定以下选项之一:

insert_insert

当你发送信号以启动增量快照时,对于 Debezium 在快照期间读取的每个分块,它都会向信号数据集合写入一条记录,以记录打开快照窗口的信号。快照完成后,Debezium 会插入第二条记录,记录关闭窗口的信号。

insert_delete

当你发送信号以启动增量快照时,对于 Debezium 读取的每个分块,它只会向信号数据集合写入一条记录,以记录打开快照窗口的信号。快照完成后,该记录会被删除。不会为关闭快照窗口的信号创建任何记录。设置此选项可防止信号数据集合快速增长。

max.iteration.transactions

500

指定每次迭代的最大事务数,用于在从数据库的多个表流式传输变更时减少内存占用。设置为 0 时,连接器使用当前最大 LSN 作为获取变更的范围。设置为大于零的值时,连接器使用此设置指定的第 n 个 LSN 作为获取变更的范围。默认值为 500。

incremental.snapshot.option.recompile

false

为增量快照期间使用的所有 SELECT 语句添加 OPTION(RECOMPILE) 查询选项。这有助于解决可能出现的参数嗅探问题,但会根据查询执行频率增加源数据库的 CPU 负载。

topic.naming.strategy

io.debezium.schema.SchemaTopicNamingStrategy

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

topic.delimiter

.

指定主题名称中使用的分隔符。

topic.cache.size

10000

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

topic.heartbeat.prefix

__debezium-heartbeat

控制连接器发送心跳消息的主题名称。主题名称遵循以下模式:

topic.heartbeat.prefix.topic.prefix

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

如果设置了 topic.heartbeat.name,则忽略此属性。

topic.heartbeat.name

空

为连接器发送心跳消息的主题指定显式的完整名称,覆盖由 topic.heartbeat.prefix 和 topic.prefix 派生的基于前缀的命名方式。

设置后,无论连接器的 topic.prefix 是什么,所有心跳消息都会路由到这个确切的主题名称。当运行多个连接器并希望将心跳事件集中到一个共享主题中时,此配置非常有用,可避免创建大量单分区的心跳主题。

例如,将其设置为 debezium-heartbeat,会将所有心跳消息路由到名为 debezium-heartbeat 的主题。

如果此属性为空或未设置,连接器将回退到默认行为:topic.heartbeat.prefix.topic.prefix。

topic.transaction

transaction

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

topic.prefix.topic.transaction

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

有关更多信息,请参阅事务元数据。

snapshot.max.threads

1

指定连接器执行初始快照时使用的线程数。该值必须是大于或等于 1 的正整数。当设置为大于 1 的值时,连接器会将每张表按主键范围划分成若干分块,并通过所有可用线程并发处理这些分块,从而执行并行快照。

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

snapshot.max.threads.multiplier

1

一个全局乘数,用于控制并行快照期间每张表创建的分块数量。默认情况下,Debezium 为每个线程创建一个分块。设置更高的值会使连接器在执行快照时为每张表创建更多、更小的分块。更小的分块可以让各线程的负载更加均衡。例如,在 4 个线程且乘数为 2 的情况下,Debezium 会创建 8 个分块,而不是 4 个。

要为特定表覆盖该乘数,请将 snapshot.max.threads.multiplier.<完全限定表名> 设置为所需值。

与提高乘数相比,增加线程数通常对提升快照吞吐量更为有效。

legacy.snapshot.max.threads

false

默认值(false)使连接器能够通过将表划分成若干分块,并使用独立线程并发处理每个分块来加速初始快照。若要恢复到旧版的并行快照行为(每个线程处理一张表),请将该值设置为 true。

启用遗留行为后,完成表快照的线程会保持空闲状态,等待其他线程完成。在强制执行连接超时的环境中,空闲连接会导致连接器在快照结束后无法正常关闭连接。这样即使快照成功捕获了所有数据,也可能引发异常。如果遇到此问题,请将 snapshot.max.threads 设置为 1 并重试快照。
属性 internal.legacy.snapshot.max.threads 是 legacy.snapshot.max.threads 的已弃用别名,不应在新配置中使用。

custom.metric.tags

无默认值

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

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

custom.sanitize.pattern

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

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

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

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

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

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

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

errors.max.retries

-1

指定当某个操作导致可重试错误(例如连接错误)时,连接器的响应方式。可设置为以下选项之一:

-1

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

0

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

> 0

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

data.query.mode

direct

控制连接器查询 CDC 数据的方式。支持以下模式:

  • function:通过调用 cdc.[fn_cdc_get_all_changes_#] 函数来查询数据。
  • direct:使连接器直接查询更改表。这是默认模式。

database.query.timeout.ms

600000(10 分钟)

指定连接器等待查询完成的时间,以毫秒为单位。将该值设置为 0(零)可移除超时限制。

streaming.fetch.size

0

指定在流式传输期间,每次从每张表中一次性读取的最大行数。连接器将按此大小分批读取表内容。默认值为 0,表示无限制。

guardrail.collections.max

0

指定连接器可以捕获的最大表数量。此限制不仅适用于正在被捕获以进行数据更改的表,也适用于其架构被作为架构历史记录一部分进行跟踪的任何表。超出此限制会触发 guardrail.collections.limit.action 中指定的操作。将此属性设置为 0 可防止连接器触发护栏操作。

guardrail.collections.limit.action

warn

指定当连接器捕获的表数量超过你在 guardrail.collections.max 属性中指定的数量时要触发的操作。将该属性设置为以下值之一:

fail

连接器失败并报告异常。

warn

连接器记录一条警告。

extended.headers.enabled

true

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

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

该属性会添加以下标头:

__debezium.context.connectorLogicalName

Debezium 连接器的逻辑名称。

__debezium.context.taskId

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

__debezium.context.connectorName

Debezium 连接器的名称。

Debezium SQL Server 连接器数据库 schema 历史记录配置属性

Debezium 提供了一组 schema.history.internal.* 属性,用于控制连接器与 schema 历史记录主题的交互方式。

下表描述了用于配置 Debezium 连接器的 schema.history.internal 属性。

表 9. 连接器数据库 schema 历史记录配置属性

属性 默认值 说明

schema.history.internal.kafka.topic

无默认值

连接器用于存储数据库 schema 历史记录的 Kafka 主题的完整名称。

schema.history.internal.kafka.bootstrap.servers

无默认值

连接器用于与 Kafka 集群建立初始连接的主机/端口对列表。此连接用于检索连接器先前存储的数据库 schema 历史记录,以及写入从源数据库读取的每条 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。跳过解析应谨慎使用,因为在处理二进制日志时可能导致数据丢失或损坏。

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 的连接器共享同一个全局驻留池。当有大量连接器跟踪结构相似的表时,可实现最大程度的去重。

SQL Server 连接器的透传配置属性

连接器支持透传属性,使 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=test1234

Debezium 会先去除属性名称中的前缀,然后再将该属性传递给 Kafka 客户端。

有关 Kafka 生产者配置属性 和 Kafka 消费者配置属性 的更多信息,请参阅 Apache Kafka 文档。

用于配置 SQL Server 连接器如何与 Kafka 信号主题交互的透传属性

Debezium 提供了一组 signal.* 属性,用于控制连接器如何与 Kafka 信号主题交互。

下表描述了 Kafka signal 属性。

表 10. Kafka 信号配置属性

属性 | 默认值 | 说明

signal.kafka.topic

<topic.prefix>-signal

连接器用于监视临时信号的 Kafka 主题名称。

如果自动创建主题被禁用,你必须手动创建所需的信号主题。为保证信号顺序,必须使用信号主题。信号主题必须只有一个分区。

signal.kafka.groupId

kafka-signal

Kafka 消费者所使用的组 ID 名称。

signal.kafka.bootstrap.servers

无默认值

连接器用于建立与 Kafka 集群初始连接的主机和端口对列表。每个对引用 Debezium Kafka Connect 进程所使用的 Kafka 集群。

signal.kafka.poll.timeout.ms

100

一个整数值,用于指定连接器在轮询信号时等待的最大毫秒数。

用于配置信号通道 Kafka 消费者客户端的透传属性

Debezium 连接器支持对信号 Kafka 消费者进行透传配置。透传信号属性以 signal.consumer.* 为前缀。例如,连接器会将 signal.consumer.security.protocol=SSL 这类属性传递给 Kafka 消费者。

Debezium 会先去除这些属性的前缀,然后再将属性传递给 Kafka 信号消费者。

用于配置 SQL Server 连接器 sink 通知通道的透传属性

下表描述了可用于配置 Debezium sink notification 通道的属性。

属性默认值说明
notification.sink.topic.name无默认值接收来自 Debezium 通知的主题名称。当您将 notification.enabled.channels 属性配置为包含 sink 作为已启用的通知通道之一时,此属性为必填项。

表 11. sink 通知配置属性

Debezium 连接器的数据库驱动透传配置属性

Debezium 连接器支持对数据库驱动进行透传配置。透传的数据库属性以 driver.* 前缀开头。例如,连接器会将 driver.foobar=false 之类的属性传递给 JDBC URL。

Debezium 在将属性传递给数据库驱动之前,会移除这些属性的前缀。

数据库模式演进

当您更改已启用 CDC 的 SQL Server 表的模式时,必须先更新相应的捕获表,连接器才能继续发出结构正确的变更事件。本节中的过程介绍了如何在模式变更后刷新捕获表。

对于启用了变更数据捕获(CDC)的 SQL Server 表,随着更改的发生,事件记录会持久化到服务器上的捕获表中。如果你对源表的结构进行了更改,例如添加新列,该更改不会动态反映到变更表中。只要捕获表继续使用过时的架构,Debezium 连接器就无法为该表正确发出数据变更事件。

由于 SQL Server 中 CDC 的实现方式,你无法使用 Debezium 来更新捕获表。要刷新捕获表,操作者必须是拥有高级权限的 SQL Server 数据库管理员。作为 Debezium 用户,你必须与 SQL Server 数据库管理员协调完成架构刷新,并恢复到 Kafka 主题的流式传输。

在架构变更后,你可以使用以下方法之一来更新捕获表:

每种方法各有利弊。

无论使用在线更新方法还是离线更新方法,你都必须在对同一源表进行后续架构更新之前,完成整个架构更新过程。最佳实践是以单个批次执行所有 DDL,这样该过程只需运行一次。
某些架构变更不支持在启用了 CDC 的源表上进行。例如,如果某个表启用了 CDC,当你重命名其中某一列或更改列类型时,SQL Server 不允许你更改该表的架构。
当你将源表中的某列从 NULL 修改为 NOT NULL,或从 NOT NULL 修改为 NULL 后,SQL Server 连接器在创建新的捕获实例之前无法正确捕获这一变更信息。如果在列属性变更后未创建新的捕获表,连接器发出的变更事件记录将无法正确标识该列是否为可空列。也就是说,之前定义为可空(即 NULL)的列,即使现在已定义为 NOT NULL,仍会被视为可空列;同样,之前定义为必填(NOT NULL)的列,即使现在已定义为 NULL,仍会保留必填的标识。
使用 sp_rename 函数重命名表后,在连接器重启之前,变更仍会以旧的源表名发出。连接器重启后,将以新的源表名发出变更。

离线模式更新

离线模式更新是更新捕获表最安全的方式。但是,对于要求高可用性的应用程序,离线更新可能并不可行。

前置条件

  • 已对启用 CDC 的 SQL Server 表的架构提交了变更。
  • 你是一名拥有提升权限的 SQL Server 数据库操作员。

操作步骤

  1. 暂停更新数据库的应用程序。
  2. 等待 Debezium 连接器将所有尚未流出的变更事件记录全部流出。
  3. 停止 Debezium 连接器。
  4. 对源表架构应用所有变更。
  5. 使用 sys.sp_cdc_enable_table 存储过程为更新后的源表创建新的捕获表,并为参数 @capture_instance 指定一个唯一的值。
  6. 恢复在步骤 1 中暂停的应用程序。
  7. 启动 Debezium 连接器。
  8. 当 Debezium 连接器开始从新的捕获表流出数据后,运行存储过程 sys.sp_cdc_disable_table 并将参数 @capture_instance 设置为旧的捕获实例名称,以删除旧的捕获表。连接器在读取完旧捕获实例后会发出一条通知。有关更多信息,参见通知。

在线架构更新

完成在线架构更新的过程比执行离线架构更新的过程更简单,而且你可以在不要求应用程序和数据处理停机的情况下完成它。但是,对于在线架构更新来说,在你更新源数据库中的架构之后、创建新的捕获实例之前,可能会出现潜在的处理间隙。在此期间,变更事件仍由旧的变更表实例捕获,保存到旧表中的变更数据保留的仍是先前架构的结构。因此,例如,如果你向源表添加了一个新列,那么在新的捕获表就绪之前产生的变更事件不会包含该新列对应的字段。如果你的应用程序无法容忍这样的过渡期,最好使用离线架构更新过程。

前提条件

  • 已经对启用了 CDC 的 SQL Server 表的架构提交了一次更新。
  • 你是拥有更高权限的 SQL Server 数据库操作员。

操作步骤

  1. 对源表架构应用所有变更。
  2. 运行 sys.sp_cdc_enable_table 存储过程并为参数 @capture_instance 指定一个唯一的值,为更新后的源表创建新的捕获表。
  3. 当 Debezium 开始从新的捕获表流出数据后,你可以运行 sys.sp_cdc_disable_table 存储过程并将参数 @capture_instance 设置为旧的捕获实例名称,以删除旧的捕获表。连接器在读取完旧捕获实例后会发出一条通知。有关更多信息,参见通知。

示例:在数据库架构变更后执行在线架构更新

我们部署基于 SQL Server 的 Debezium 教程来演示在线架构更新。

在下面的示例中,向 customers 表添加了一个 phone_number 列。

  1. 输入以下命令以启动数据库 shell:
docker-compose -f docker-compose-sqlserver.yaml exec sqlserver bash -c '/opt/mssql-tools/bin/sqlcmd -U sa -P $SA_PASSWORD -d testDB'
  1. 修改 customers 源表的结构,运行以下查询以添加 phone_number 字段:

    ALTER TABLE customers ADD phone_number VARCHAR(32);
  2. 通过运行 sys.sp_cdc_enable_table 存储过程创建新的捕获实例。

    EXEC sys.sp_cdc_enable_table @source_schema = 'dbo', @source_name = 'customers', @role_name = NULL, @supports_net_changes = 0, @capture_instance = 'dbo_customers_v2';
    GO
  3. 运行以下查询,向 customers 表中插入新数据:

    INSERT INTO customers(first_name,last_name,email,phone_number) VALUES ('John','Doe','[email protected]', '+1-555-123456');
    GO

Kafka Connect 的日志会通过类似下面的消息条目来报告配置更新:

connect_1    | 2019-01-17 10:11:14,924 INFO   ||  Multiple capture instances present for the same table: Capture instance "dbo_customers" [sourceTableId=testDB.dbo.customers, changeTableId=testDB.cdc.dbo_customers_CT, startLsn=00000024:00000d98:0036, changeTableObjectId=1525580473, stopLsn=00000025:00000ef8:0048] and Capture instance "dbo_customers_v2" [sourceTableId=testDB.dbo.customers, changeTableId=testDB.cdc.dbo_customers_v2_CT, startLsn=00000025:00000ef8:0048, changeTableObjectId=1749581271, stopLsn=NULL]   [io.debezium.connector.sqlserver.SqlServerStreamingChangeEventSource]
connect_1    | 2019-01-17 10:11:14,924 INFO   ||  Schema will be changed for ChangeTable [captureInstance=dbo_customers_v2, sourceTableId=testDB.dbo.customers, changeTableId=testDB.cdc.dbo_customers_v2_CT, startLsn=00000025:00000ef8:0048, changeTableObjectId=1749581271, stopLsn=NULL]   [io.debezium.connector.sqlserver.SqlServerStreamingChangeEventSource]
...
connect_1    | 2019-01-17 10:11:33,719 INFO   ||  Migrating schema to ChangeTable [captureInstance=dbo_customers_v2, sourceTableId=testDB.dbo.customers, changeTableId=testDB.cdc.dbo_customers_v2_CT, startLsn=00000025:00000ef8:0048, changeTableObjectId=1749581271, stopLsn=NULL]   [io.debezium.connector.sqlserver.SqlServerStreamingChangeEventSource]

最终,phone_number 字段会被添加到模式中,其值也会出现在写入 Kafka 主题的消息里。

...
     {
        "type": "string",
        "optional": true,
        "field": "phone_number"
     }
...
    "after": {
      "id": 1005,
      "first_name": "John",
      "last_name": "Doe",
      "email": "[email protected]",
      "phone_number": "+1-555-123456"
    },
  1. 运行 sys.sp_cdc_disable_table 存储过程来删除旧的捕获实例。

    EXEC sys.sp_cdc_disable_table @source_schema = 'dbo', @source_name = 'dbo_customers', @capture_instance = 'dbo_customers';
    GO

通知

Debezium SQL Server 连接器可以发出通知事件,用来报告连接器的重要状态变化,例如架构演进变更的完成。你可以配置通知来监控连接器的活动,并触发下游工作流。

如果为连接器启用了通知功能,那么在进行架构更改后,连接器会发出一条通知来报告该更改的状态。有关通知及其启用方式的更多信息,请参阅 Debezium 通知文档。

以下示例展示了 Debezium 在架构演进变更之后可能发出的状态通知。

{
      "id":"ff81ba59-15ea-42ae-b5d0-4d74f1f4038f",
      "aggregate_type":"Capture Instance",
      "type":"COMPLETED",
      "additional_data":{
         "connector_name":"my-connector",
         "capture_instance":"dbo_customers",
         "server":"server1",
         "database":"testDB",
         "start_lsn":"00000027:00000758:0003",
         "stop_lsn":"00000028:00000759:0004",
         "commit_lsn":"00000029:00000760:0005"
      },
      "timestamp": "1695817046353"
}

监控

除了 Kafka 和 Kafka Connect 提供的内置 JMX 指标支持之外,Debezium SQL Server 连接器还提供了三种类型的指标。该连接器提供以下指标:

有关如何通过 JMX 暴露上述指标的信息,请参阅 Debezium 监控文档。

自定义 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 名称

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

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

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

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

快照指标

MBean 为 debezium.sql_server:type=connector-metrics,server=<topic.prefix>,task=<task.id>,context=snapshot。

下表列出了可用于监控 Debezium 快照操作的 JMX 指标,包括行计数、表进度、持续时间和队列容量。除非快照操作正在进行,或者自上次连接器启动以来发生过快照,否则不会暴露快照指标。

属性类型描述
LastEventstring连接器读取的最后一个快照事件。
MilliSecondsSinceLastEventlong自连接器读取并处理最近一个事件以来经过的毫秒数。
NumberOfErroneousEventslong记录连接器在快照操作期间识别为错误的变更事件数量。每当连接器在初始快照、增量快照或临时快照过程中遇到无法处理的事件时,该指标都会递增。事件处理失败的原因可能包括:事件格式错误、与模式(schema)不兼容,或在转换过程中发生故障。该指标值在连接器任务的整个生命周期内持续保留。如果快照被中断且连接器任务重新启动,该指标计数将重置为 0。
TotalNumberOfEventsSeenlong自上次启动或重置以来,该连接器看到的事件总数。
NumberOfEventsFilteredlong被连接器上配置的包含/排除列表过滤规则过滤掉的事件数量。
CapturedTablesstring[]连接器捕获的表的列表。
QueueTotalCapacityint用于在快照程序和主 Kafka Connect 循环之间传递事件的队列长度。
QueueRemainingCapacityint用于在快照程序和主 Kafka Connect 循环之间传递事件的队列的剩余容量。
TotalTableCountint被纳入快照的表的总数。
RemainingTableCountint快照尚未复制的表的数量。
SnapshotRunningboolean快照是否已启动。
SnapshotPausedboolean快照是否已暂停。
SnapshotAbortedboolean快照是否已中止。
SnapshotCompletedboolean快照是否已完成。
SnapshotSkippedboolean快照是否已被跳过。
SnapshotDurationInSecondslong到目前为止快照所花费的总秒数,即使尚未完成。其中也包含快照被暂停的时间。
SnapshotPausedDurationInSecondslong快照被暂停的总秒数。如果快照被暂停多次,暂停时间将累加。
RowsScannedMap<String, Long>一个映射,包含快照中每张表已扫描的行数。在处理过程中,表会被逐步添加到该映射中。每扫描 10,000 行以及完成一张表时都会更新。
TableChunkCountsMap<String, Long>一个映射,包含使用基于分块的多线程快照时,快照中每张表的分块数量。
TableChunksCompletedCountsMap<String, Long>一个映射,包含使用基于分块的多线程快照时,快照中每张表已完成的分块数量。
MaxQueueSizeInByteslong队列的最大缓冲区大小(以字节为单位)。当 max.queue.size.in.bytes 被设置为正的 long 值时,该指标可用。
CurrentQueueSizeInByteslong队列中记录的当前体积(以字节为单位)。

下表列出了连接器执行增量快照时可用的其他 JMX 指标,其中包括可用于跟踪快照进度的分块和表边界标识符。

属性类型描述
ChunkIdstring当前快照分块的标识符。
ChunkFromstring定义当前分块的主键集合的下界。
ChunkTostring定义当前分块的主键集合的上界。
TableFromstring当前正在执行快照的表的主键集合的下界。
TableTostring当前正在执行快照的表的主键集合的上界。

流处理指标

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

下表列出了可用于监控 Debezium 流处理操作的 JMX 指标,其中包括按类型统计的事件数量、相对源端的延迟、队列容量以及连接状态。

属性类型描述
LastEventstring连接器读取到的最后一个流式事件。
MilliSecondsSinceLastEventlong自连接器读取并处理最近一个事件以来经过的毫秒数。
NumberOfErroneousEventslong记录连接器在流式传输过程中识别为错误的变更事件数量。在流式会话的生命周期内,每当连接器遇到无法处理的事件时,该指标就会递增。事件处理失败的原因可能包括:事件格式错误、与 Schema 不兼容,或在转换过程中发生失败。该指标的值在连接器任务的整个生命周期内保持不变。连接器重启后,该指标计数会重置为 0。
TotalNumberOfEventsSeenlong自上次启动连接器或重置指标以来,源数据库上报的数据变更事件总数。代表 Debezium 需要处理的数据变更工作负载。
TotalNumberOfCreateEventsSeenlong自上次启动或重置指标以来,连接器处理的创建(create)事件总数。
TotalNumberOfUpdateEventsSeenlong自上次启动或重置指标以来,连接器处理的更新(update)事件总数。
TotalNumberOfDeleteEventsSeenlong自上次启动或重置指标以来,连接器处理的删除(delete)事件总数。
NumberOfEventsFilteredlong被连接器上配置的包含/排除列表过滤规则过滤掉的事件数量。
NumberOfUnchangedEventsSkippedlong自上次启动连接器或重置指标以来,因被监控列未发生变化而跳过的更新事件数量。如果 skip.messages.without.change 为 false,则默认值为 -1。
CapturedTablesstring[]连接器所捕获的表列表。
QueueTotalCapacityint用于在流式传输器和 Kafka Connect 主循环之间传递事件的队列长度。
QueueRemainingCapacityint用于在流式传输器和 Kafka Connect 主循环之间传递事件的队列的剩余可用容量。
Connectedboolean表示连接器当前是否已连接到数据库服务器的标志。
MilliSecondsBehindSourcelong最后一个变更事件的时间戳与连接器处理该事件之间相差的毫秒数。该值会包含数据库服务器和连接器所在机器之间时钟的差异。
MilliSecondsBehindSourceMinValuelong连接器运行期间观测到的相对于源的最小滞后毫秒数。
MilliSecondsBehindSourceMaxValuelong连接器运行期间观测到的相对于源的最大滞后毫秒数。
MilliSecondsBehindSourceAverageValuedouble连接器运行期间所有观测值的相对于源的平均滞后毫秒数。
MilliSecondsBehindSourceP50double相对于源的滞后毫秒数的第 50 百分位数(中位数)。与平均值相比,该指标对典型滞后的衡量更为稳健,因为它受异常值的影响较小。仅当 statistics.metrics.enabled 设置为 true(默认值)时可用。
MilliSecondsBehindSourceP95double相对于源的滞后毫秒数的第 95 百分位数。该指标表示 95% 的滞后测量值低于此值,有助于识别尾部延迟并设置 SLA 阈值。仅当 statistics.metrics.enabled 设置为 true(默认值)时可用。
MilliSecondsBehindSourceP99double相对于源的滞后毫秒数的第 99 百分位数。该指标表示 99% 的滞后测量值低于此值,有助于了解最坏情况下的性能表现。仅当 statistics.metrics.enabled 设置为 true(默认值)时可用。
NumberOfCommittedTransactionslong已提交的已处理事务数量。
SourceEventPositionMap<String, String>最后接收到的事件的位置信息。

| LastTransactionId | string | 最后处理的事务

架构历史指标

MBean 为 debezium.sql_server:type=connector-metrics,context=schema-history,server=<topic.prefix>,task=<task.id>。

下表列出了可用于监控连接器架构历史过程的 JMX 指标,包括恢复状态、已应用的架构变更数量以及最近变更的时间戳。

属性类型描述
Statusstring描述数据库架构历史状态的取值之一:STOPPED、RECOVERING(正在从存储中恢复历史记录)或 RUNNING。
RecoveryStartTimelong恢复开始的时间,以纪元秒(epoch seconds)表示。
ChangesRecoveredlong恢复阶段期间读取的更改数量。
ChangesAppliedlong恢复和运行期间已应用的架构更改总数。
MilliSecondsSinceLast​RecoveredChangelong自从从历史存储中恢复上一次更改以来经过的毫秒数。
MilliSecondsSinceLast​AppliedChangelong自上一次更改被应用以来经过的毫秒数。
LastRecoveredChangestring从历史存储中恢复的最近一次更改的字符串表示。
LastAppliedChangestring最近一次已应用更改的字符串表示。

评论

登录后参与评论

正在加载评论…