源连接器

PostgreSQL

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

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

Debezium PostgreSQL 连接器

Debezium PostgreSQL 连接器捕获 PostgreSQL 数据库模式中的行级变更。有关与该连接器兼容的 PostgreSQL 版本信息,请参阅 Debezium 发布版本概述。

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

概述

PostgreSQL 的逻辑解码(logical decoding)功能是在 9.4 版本中引入的。它是一种机制,能够提取已提交到事务日志中的变更,并借助输出插件以对用户友好的方式处理这些变更。输出插件使客户端能够消费这些变更。

PostgreSQL 连接器包含两个协同工作的主要部分,用于读取和处理数据库变更:

逻辑解码输出插件

要使用该连接器,必须将 PostgreSQL 配置为使用逻辑解码,并选择一个输出插件。你可能需要安装所选的输出插件。默认情况下,只要 Debezium 以拥有所需权限的用户身份连接数据库,连接器启动后会自动创建已配置的复制槽(如果该复制槽尚不存在)。在生产环境中,如果你想对安装过程施加更强的控制,可以手动创建复制槽。你可以将连接器配置为使用以下任一插件:

decoderbufs

基于 Protobuf 并由 Debezium 社区维护的解码插件。

pgoutput

PostgreSQL 10+ 中的标准逻辑解码输出插件。它由 PostgreSQL 社区维护,并被 PostgreSQL 自身用于逻辑复制。该插件始终存在,因此无需安装额外的库。Debezium 连接器直接将原始复制事件流解释为变更事件。

Java 代码

用于读取所选逻辑解码输出插件所产生的变更的 Kafka Connect 连接器。该连接器通过 PostgreSQL JDBC 驱动程序,使用 PostgreSQL 的流式复制协议。

对于捕获到的每一行级插入、更新和删除操作,连接器都会产生一个变更事件,并为每张表的变更事件记录发送到各自独立的 Kafka 主题中。客户端应用程序读取与感兴趣数据库表相对应的 Kafka 主题,并可对接收自这些主题的每一个行级事件作出响应。

PostgreSQL 通常会在一段时间后清除预写日志(WAL)段。这意味着连接器并不具备对数据库所做全部变更的完整历史记录。因此,当 PostgreSQL 连接器首次连接到特定的 PostgreSQL 数据库时,它会先对数据库的每个架构执行一致性快照。连接器完成快照后,会从快照完成的准确位置继续流式传输变更。这样一来,连接器一开始便拥有一致的全量数据视图,且不会遗漏快照期间所发生的任何变更。

该连接器具备故障容错能力。连接器在读取变更并产生事件时,会为每个事件记录对应的 WAL 位置。如果连接器因任何原因停止运行(包括通信故障、网络问题或崩溃),重启后连接器会从上次中断的位置继续读取 WAL。快照也遵循此机制:如果连接器在快照期间停止,重启后将开始新的快照。

该连接器依赖并反映 PostgreSQL 的逻辑解码功能,而该功能存在以下限制:

  • 逻辑解码不支持 DDL 变更。这意味着连接器无法将 DDL 变更事件报告给消费者。
  • 由于逻辑解码复制槽在提交时(而非提交后)发布变更,因此可能产生不良副作用。客户端可能观察到不一致状态的两种主要场景如下:第一,在复制完成前主库宕机,导致未提交的变更被发布;第二,由于变更正在复制过程中,暂时无法读取(即读后写一致性)的变更被发布。例如,DebeziumEngine 消费者收到某行已创建的通知,却无法通过事务读取该行。

此外,pgoutput 逻辑解码输出插件不会捕获生成列的值,导致连接器输出中缺少这些列的数据。

出现问题时的行为描述了连接器在遇到问题时的响应方式。

Debezium 目前仅支持 UTF-8 字符编码的数据库。对于单字节字符编码,无法正确处理包含扩展 ASCII 字符的字符串。

连接器的工作原理

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

安全性

要使用 Debezium 连接器从 PostgreSQL 数据库流式传输变更,该连接器必须以数据库中的特定权限运行。虽然授予所需权限的一种方式是为用户提供 superuser 权限,但这样做可能会使您的 PostgreSQL 数据面临未授权访问的风险。与其给 Debezium 用户授予过多权限,最佳做法是创建一个专用的 Debezium 复制用户,并仅向其授予所需的特定权限。

有关为 Debezium PostgreSQL 用户配置权限的更多信息,请参阅设置权限。有关 PostgreSQL 逻辑复制安全性的更多信息,请参阅 PostgreSQL 文档。

快照

大多数 PostgreSQL 服务器的配置不会在 WAL 段中保留数据库的完整历史记录。这意味着 PostgreSQL 连接器仅通过读取 WAL 将无法看到数据库的全部历史记录。因此,连接器首次启动时会对数据库执行一次初始一致性快照。

初始快照的默认工作流行为

以下步骤描述了连接器在初始快照期间执行的默认操作。您可以通过将 snapshot.mode 连接器配置属性设置为 initial 以外的其他值来更改此行为。

  1. 启动一个事务,其隔离级别由 snapshot.isolation.mode 属性指定。指定的模式决定该事务中后续的读取操作是针对单一一致的数据版本。根据所选模式,其他客户端随后执行 INSERT、UPDATE 和 DELETE 操作所造成的数据变更,可能对该事务可见。
  2. 读取服务器事务日志中的当前位置。
  3. 扫描数据库表和模式,为每一行生成一个 READ 事件,并将该事件写入相应表专属的 Kafka 主题。
  4. 提交事务。
  5. 在连接器偏移量中记录快照已成功完成。

如果连接器在第 1 步开始之后、第 5 步完成之前失败、被重新平衡或停止,重启时连接器会开始新的快照。在连接器完成初始快照之后,PostgreSQL 连接器会从第 2 步中读取的位置继续流式传输。这可确保连接器不会遗漏任何更新。如果连接器因任何原因再次停止,重启后它会从之前中断的位置继续流式传输变更。

表 1. snapshot.mode 连接器配置属性的可选项

选项 说明

always

连接器在启动时始终执行快照。快照完成后,连接器会从上述序列中的第 3 步继续流式传输变更。此模式适用于以下情形:

  • 已知某些 WAL 段已被删除且不再可用。
  • 集群故障后提升出新的主节点时。always 快照模式可确保连接器不会遗漏在新主节点提升之后、但在连接器于新主节点上重启之前所做的任何变更。

initial(默认)

当不存在 Kafka 偏移量主题时,连接器执行数据库快照。数据库快照完成后,会写入 Kafka 偏移量主题。如果 Kafka 偏移量主题中存在之前存储的 LSN,连接器将从该位置继续流式传输变更。

initial_only

连接器执行数据库快照,并在流式传输任何变更事件记录之前停止。如果连接器在停止之前已启动快照但未完成,连接器会重新开始快照流程,并在快照完成时停止。

no_data

连接器从不执行快照。当连接器以此方式配置时,启动后的行为如下:

如果 Kafka 偏移量主题中存在之前存储的 LSN,连接器将从该位置继续流式传输变更。如果没有存储 LSN,连接器会从服务器上创建 PostgreSQL 逻辑复制槽的位置开始流式传输变更。仅当你确定所有相关数据仍反映在 WAL 中时,才使用此快照模式。

when_needed

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

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

configuration_based

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

custom

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

基于分块的并行快照

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

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

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

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

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

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

特性 基于分块(默认) 传统的每线程一表

每个线程的工作单元

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

每个线程处理整张表。

负载分配

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

在整个快照过程中,每个线程始终分配给同一张表。处理完较小表的线程会处于空闲状态,而其他线程则继续处理较大的表。

空闲连接风险

低。

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

高。

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

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

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

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

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

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

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

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

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

你可以通过发起所谓的临时快照,对之前已捕获过快照的表重新运行快照。临时快照需要使用信号表。你通过向 Debezium 信号表发送信号请求来发起临时快照。

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

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

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

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

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

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

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

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

触发临时增量快照

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

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

触发临时阻塞快照

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

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

增量快照

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

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

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

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

增量快照的处理过程

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

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

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

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

快照窗口

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

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

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

连接器会针对每个快照块重复这一过程。

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

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

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

Debezium 的 PostgreSQL 连接器不支持在增量快照运行期间进行模式变更。如果在增量快照开始之前、但发送信号之后执行了模式变更,则需将直通配置选项 database.autosave 设置为 conservative,以便正确处理该模式变更。

触发增量快照

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

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

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

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

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

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

如果要包含在快照中的表名包含点号(.)、空格或其他非字母数字字符,则必须用双引号对表名进行转义。例如,要包含位于 public 模式中、名称为 My.Table 的表,请使用以下格式:"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 myschema.debezium_signal (id, type, data)
values ('ad-hoc-1',
    'execute-snapshot',
    '{"data-collections": ["schema1.table1", "schema1.table2"],
    "type":"incremental",
    "additional-conditions":[{"data-collection": "schema1.table1" ,"filter":"color=\'blue\'"}]}');

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

INSERT INTO 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 运行即时增量快照。

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

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

物理行标识符在以下场景中特别有用:

具有复合主键的表

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

索引利用效率低下

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

通过使用物理行标识符作为代理键,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> ....

下面的示例展示了如何向信号表发送一个临时增量快照请求,并附加一个额外条件:

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 myschema.debezium_signal (id, type, data) VALUES('ad-hoc-1', 'execute-snapshot', '{"data-collections": ["schema1.products"],"type":"incremental", "additional-conditions":[{"data-collection": "schema1.products", "filter": "color=blue"}]}');

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

INSERT INTO myschema.debezium_signal (id, type, data) VALUES('ad-hoc-1', 'execute-snapshot', '{"data-collections": ["schema1.products"],"type":"incremental", "additional-conditions":[{"data-collection": "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 消息的 key 必须与连接器配置选项 topic.prefix 的值相匹配。

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

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

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

字段 默认值

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": ["schema1.products"], "type": "INCREMENTAL", "additional-conditions": [{"data-collection": "schema1.products" ,"filter":"color='blue'"}]}}`

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

Key = `test_connector`

Value = `{"type":"execute-snapshot","data": {"data-collections": ["schema1.products"], "type": "INCREMENTAL", "additional-conditions": [{"data-collection": "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 myschema.debezium_signal (id, type, data)
values ('ad-hoc-1',
    'stop-snapshot',
    '{"data-collections": ["schema1.table1", "schema1.table2"],
    "type":"incremental"}');

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

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

schema.debezium_signal

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

ad-hoc-1

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

stop-snapshot

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

data-collections

信号 data 字段的可选组成部分,用于指定一个表名数组或用于匹配表名的正则表达式,以从快照中移除这些表。该数组列出的正则表达式按照 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不适用一个可选的字符串数组,包含若干用逗号分隔的正则表达式,用于匹配需要从快照中移除的表的全限定名(表名数组或用于匹配表名的正则表达式)。请使用 schema.table 格式指定表名。

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

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

Key = `test_connector`

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

只读增量快照

你可以配置一个对数据库具有只读连接的 PostgreSQL 连接器,使其无需信号数据收集表即可运行增量快照。为了在只读访问下运行增量快照,连接器会将当前正在进行的事务用作高水位标记和低水位标记。通过将预写日志事件或心跳事件的事务 ID 与低水位标记和高水位标记进行比较,来更新数据块窗口的状态。

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

前提条件

  • PostgreSQL 版本 13 或更高。

当 PostgreSQL 连接器使用只读连接时,你可以在不使用信号数据收集表的情况下触发增量快照。相反,你可以使用任何可用的信号通道。

自定义快照器 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"}]}

可能出现重复记录

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

流式传输变更

PostgreSQL 连接器在绝大多数时间里,都在将其所连接的 PostgreSQL 服务器的变更以流式方式传输过来。这一机制依赖于 PostgreSQL 复制协议。该协议使客户端能够在特定位置接收服务器事务日志中已提交的变更,这些位置称为日志序列号(LSN)。

每当服务器提交一个事务时,一个独立的服务器进程会从逻辑解码插件中调用一个回调函数。该函数处理该事务中的变更,将它们转换为特定格式(对于 Debezium 插件而言是 Protobuf 或 JSON),并将其写入输出流,随后客户端即可消费这些数据。

Debezium PostgreSQL 连接器充当 PostgreSQL 客户端的角色。当连接器接收到变更时,它会将这些事件转换为包含事件 LSN 的 Debezium 创建(create)、更新(update) 或 删除(delete) 事件。PostgreSQL 连接器以记录的形式将这些变更事件转发给在同一进程中运行的 Kafka Connect 框架。Kafka Connect 进程以异步方式,按照变更事件生成的相同顺序,将这些变更事件记录写入相应的 Kafka 主题。

Kafka Connect 会定期在另一个 Kafka 主题中记录最新的 偏移量(offset)。偏移量表示 Debezium 随每个事件一起包含的、与源特定的位置信息。对于 PostgreSQL 连接器而言,每个变更事件中记录的 LSN 就是偏移量。

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

PostgreSQL 连接器会将模式(schema)信息作为逻辑解码插件所发送事件的一部分来获取。但是,连接器不会获取由哪些列组成主键的信息。连接器通过 JDBC 元数据(旁路通道)获取此信息。如果某张表的主键定义发生变化(添加、删除或重命名主键列),则在一段极短的时间内,来自 JDBC 的主键信息与逻辑解码插件生成的更改事件不同步。在这段极短的时间内,可能会生成键结构不一致的消息。为避免这种不一致,请按以下步骤更新主键结构:

  1. 将数据库或应用程序置于只读模式。
  2. 让 Debezium 处理完所有剩余事件。
  3. 停止 Debezium。
  4. 更新相关表的主键定义。
  5. 将数据库或应用程序置于读写模式。
  6. 重新启动 Debezium。

PostgreSQL 10+ 逻辑解码支持(pgoutput)

从 PostgreSQL 10+ 开始,PostgreSQL 原生支持一种名为 pgoutput 的逻辑复制流模式。这意味着 Debezium PostgreSQL 连接器无需安装额外的插件即可消费该复制流。对于不支持或不允许安装插件的环境而言,这一点尤其有价值。

有关更多信息,请参阅设置 PostgreSQL。

主题名称

默认情况下,PostgreSQL 连接器会将表中发生的所有 INSERT、UPDATE 和 DELETE 操作的更改事件写入与该表对应的单个 Apache Kafka 主题中。连接器使用以下约定为更改事件主题命名:

topicPrefix.schemaName.tableName

以下列表说明了默认名称中各个组成部分的含义:

topicPrefix

由 topic.prefix 配置属性指定的主题前缀。

schemaName

发生更改事件的数据库模式的名称。

tableName

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

例如,假设 fulfillment 是某个连接器配置中的逻辑服务器名称,该连接器正在捕获一个 PostgreSQL 实例的更改。该实例包含一个 postgres 数据库,其中有一个 inventory 模式,该模式包含四张表:products、products_on_hand、customers 和 orders。连接器将把记录流式传输到以下四个 Kafka 主题:

  • fulfillment.inventory.products
  • fulfillment.inventory.products_on_hand
  • fulfillment.inventory.customers
  • fulfillment.inventory.orders

现在假设这些表不属于某个特定模式,而是在 PostgreSQL 默认的 public 模式中创建的。那么 Kafka 主题的名称将是:

  • fulfillment.public.products
  • fulfillment.public.products_on_hand
  • fulfillment.public.customers
  • fulfillment.public.orders

连接器采用类似的命名约定来标记其事务元数据主题。

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

事务元数据

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

Debezium 接收事务元数据的限制

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

对于每个事务的 BEGIN 和 END,Debezium 会生成一个包含以下字段的事件:

status

BEGIN 或 END。

id

唯一事务标识符的字符串表示,由 Postgres 事务 ID 本身与给定操作的 LSN 组成,以冒号分隔,即格式为 txID:LSN。

ts_ms

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

event_count(用于 END 事件)

该事务产生的事件总数。

data_collections(用于 END 事件)

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

示例

{
  "status": "BEGIN",
  "id": "571:53195829",
  "ts_ms": 1486500577125,
  "event_count": null,
  "data_collections": null
}

{
  "status": "END",
  "id": "571:53195832",
  "ts_ms": 1486500577691,
  "event_count": 2,
  "data_collections": [
    {
      "data_collection": "s1.a",
      "event_count": 1
    },
    {
      "data_collection": "s2.a",
      "event_count": 1
    }
  ]
}

除非通过 topic.transaction 选项覆盖,否则事务事件将写入名为 <topic.prefix>.transaction 的主题。

变更数据事件的丰富处理

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

id

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

total_order

该事件在该事务生成的所有事件中的绝对位置。

data_collection_order

该事件在该事务发出的所有事件中,针对单个数据集合的相对位置。

以下是消息示例:

{
  "before": null,
  "after": {
    "pk": "2",
    "aa": "1"
  },
  "source": {
   ...
  },
  "op": "c",
  "ts_ms": "1580390884335",
  "ts_us": "1580390884335451",
  "ts_ns": "1580390884335451325",
  "transaction": {
    "id": "571:53195832",
    "total_order": "1",
    "data_collection_order": "1"
  }
}

数据变更事件

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

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

下面的 JSON 骨架展示了一个变更事件的四个基本部分。不过,你在应用中选择使用的 Kafka Connect 转换器(converter)如何配置,决定了这四个部分在变更事件中的表现形式。只有当你配置转换器生成 schema 字段时,变更事件中才会包含 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 字段所描述的结构,并且包含发生变更的行的实际数据。

默认情况下,连接器会将变更事件记录流式传输到主题名与其事件来源表同名的主题中。

从 Kafka 0.10 开始,Kafka 可以选择性地将事件键和值与消息创建时(由生产者记录)或由 Kafka 写入日志时的时间戳 一并记录。

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

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

变更事件键

对于给定的表,变更事件的键(key)具有这样的结构:在事件创建时,键中会为该表主键中的每个列包含一个字段。或者,如果该表的 REPLICA IDENTITY 设置为 FULL 或 USING INDEX,则会为每个唯一键约束包含一个字段。

考虑在 public 数据库模式中定义的 customers 表,以及该表的变更事件键示例。

示例表

列名 类型
id numeric
first_name varchar(50)
last_name varchar(50)
email varchar(725)
CREATE TABLE customers (
  id SERIAL,
  first_name VARCHAR(255) NOT NULL,
  last_name VARCHAR(255) NOT NULL,
  email VARCHAR(255) NOT NULL,
  PRIMARY KEY(id)
);

变更事件键示例

如果 topic.prefix 连接器配置属性的值为 PostgreSQL_server,那么只要 customers 表保持当前的定义,其所有的变更事件都具有相同的键结构,用 JSON 表示如下:

{
  "schema": {
    "type": "struct",
    "name": "PostgreSQL_server.public.customers.Key",
    "optional": false,
    "fields": [
          {
              "name": "id",
              "index": "0",
              "schema": {
                  "type": "INT32",
                  "optional": "false"
              }
          }
      ]
  },
  "payload": {
      "id": "1"
  },
}

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

schema

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

schema.name

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

  • PostgreSQL_server 是生成此事件的连接器的名称。
  • inventory 是包含被更改表的数据库。
  • customers 是被更新的表。

optional

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

fields

指定 payload 中预期的每个字段,包括每个字段的名称、索引和 schema。

payload

包含生成此变更事件的行的键。在本示例中,该键包含一个值为 1 的 id 字段。

尽管 column.exclude.list 和 column.include.list 连接器配置属性允许你仅捕获表列的子集,但主键或唯一键中的所有列始终会包含在事件的键中。
如果表没有主键或唯一键,则变更事件的键为 null。没有主键或唯一键约束的表中的行无法被唯一标识。

变更事件值

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

考虑与展示变更事件键示例时所用的同一示例表:

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

该表发生变更时,变更事件中 value 部分的内容取决于 REPLICA IDENTITY 设置以及事件所对应的操作类型。

副本身份

REPLICA IDENTITY 是 PostgreSQL 特有的表级设置,用于确定逻辑解码插件在处理 UPDATE 和 DELETE 事件时能够获得多少信息。更具体地说,每当发生 UPDATE 或 DELETE 事件时,REPLICA IDENTITY 的设置决定了涉及的表列的旧值(如果有的话)是否可用。

REPLICA IDENTITY 有 4 个可选值:

  • DEFAULT - 默认行为是:如果表有主键,则 UPDATE 和 DELETE 事件包含该表主键列的旧值。对于 UPDATE 事件,只包含值发生变化的主键列。

    如果表没有主键,连接器就不会为该表发出 UPDATE 或 DELETE 事件。对于没有主键的表,连接器只会发出 create(创建)事件。通常,没有主键的表用于在表末尾追加消息,这意味着 UPDATE 和 DELETE 事件并无用处。

  • NOTHING - 发出的 UPDATE 和 DELETE 操作事件不包含任何表列的旧值信息。

  • FULL - 发出的 UPDATE 和 DELETE 操作事件包含表中所有列的旧值。

  • INDEX index-name - 发出的 UPDATE 和 DELETE 操作事件包含指定索引中所含各列的旧值。UPDATE 事件还包含更新后的索引列值。

create(创建)事件

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

{
    "schema": {
        "type": "struct",
        "fields": [
            {
                "type": "struct",
                "fields": [
                    {
                        "type": "int32",
                        "optional": false,
                        "field": "id"
                    },
                    {
                        "type": "string",
                        "optional": false,
                        "field": "first_name"
                    },
                    {
                        "type": "string",
                        "optional": false,
                        "field": "last_name"
                    },
                    {
                        "type": "string",
                        "optional": false,
                        "field": "email"
                    }
                ],
                "optional": true,
                "name": "PostgreSQL_server.inventory.customers.Value",
                "field": "before"
            },
            {
                "type": "struct",
                "fields": [
                    {
                        "type": "int32",
                        "optional": false,
                        "field": "id"
                    },
                    {
                        "type": "string",
                        "optional": false,
                        "field": "first_name"
                    },
                    {
                        "type": "string",
                        "optional": false,
                        "field": "last_name"
                    },
                    {
                        "type": "string",
                        "optional": false,
                        "field": "email"
                    }
                ],
                "optional": true,
                "name": "PostgreSQL_server.inventory.customers.Value",
                "field": "after"
            },
            {
                "type": "struct",
                "fields": [
                    {
                        "type": "string",
                        "optional": false,
                        "field": "version"
                    },
                    {
                        "type": "string",
                        "optional": false,
                        "field": "connector"
                    },
                    {
                        "type": "string",
                        "optional": false,
                        "field": "name"
                    },
                    {
                        "type": "int64",
                        "optional": false,
                        "field": "ts_ms"
                    },
                    {
                        "type": "int64",
                        "optional": false,
                        "field": "ts_us"
                    },
                    {
                        "type": "int64",
                        "optional": false,
                        "field": "ts_ns"
                    },
                    {
                        "type": "boolean",
                        "optional": true,
                        "default": false,
                        "field": "snapshot"
                    },
                    {
                        "type": "string",
                        "optional": false,
                        "field": "db"
                    },
                    {
                        "type": "string",
                        "optional": false,
                        "field": "schema"
                    },
                    {
                        "type": "string",
                        "optional": false,
                        "field": "table"
                    },
                    {
                        "type": "int64",
                        "optional": true,
                        "field": "txId"
                    },
                    {
                        "type": "int64",
                        "optional": true,
                        "field": "lsn"
                    },
                    {
                        "type": "int64",
                        "optional": true,
                        "field": "xmin"
                    }
                ],
                "optional": false,
                "name": "io.debezium.connector.postgresql.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": "PostgreSQL_server.inventory.customers.Envelope"
    },
    "payload": {
        "before": null,
        "after": {
            "id": 1,
            "first_name": "Anne",
            "last_name": "Kretchmar",
            "email": "[email protected]"
        },
        "source": {
            "version": "3.6.3.Final",
            "connector": "postgresql",
            "name": "PostgreSQL_server",
            "ts_ms": 1559033904863,
            "ts_us": 1559033904863123,
            "ts_ns": 1559033904863123000,
            "snapshot": true,
            "db": "postgres",
            "sequence": "[\"24023119\",\"24023128\"]",
            "schema": "public",
            "table": "customers",
            "txId": 555,
            "lsn": 24023128,
            "xmin": null
        },
        "op": "c",
        "ts_ms": 1559033904863,
        "ts_us": 1559033904863841,
        "ts_ns": 1559033904863841257
    }
}

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

schema

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

name(PostgreSQL_server.inventory.customers.Value)

在 schema 部分中,每个 name 字段指定值负载中某个字段的 schema。

PostgreSQL_server.inventory.customers.Value 是负载中 before 和 after 字段的 schema。该 schema 是 customers 表特有的。

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

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

io.debezium.connector.postgresql.Source 是负载中 source 字段的 schema。该 schema 是 PostgreSQL 连接器特有的,连接器对其生成的所有事件都使用该 schema。

name(PostgreSQL_server.inventory.customers.Envelope)

PostgreSQL_server.inventory.customers.Envelope 是负载整体结构的 schema,其中 PostgreSQL_server 是连接器名称,inventory 是数据库,customers 是表。

payload

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

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

before

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

该字段是否可用取决于每张表的 REPLICA IDENTITY 设置。

after

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

source

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

  • Debezium 版本
  • 连接器类型和名称
  • 包含新行的数据库和表
  • 附加偏移信息的 JSON 字符串数组。第一个值始终是最后提交的 LSN,第二个值始终是当前 LSN。这两个值都可能为 null。
  • Schema 名称
  • 该事件是否属于快照
  • 执行该操作所在事务的 ID
  • 该操作在数据库日志中的偏移量
  • 该变更在数据库中发生的时间戳

op

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

c

创建

u

更新

d

删除

r

读取(仅适用于快照)

t

截断

m

消息

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": 1
        },
        "after": {
            "id": 1,
            "first_name": "Anne Marie",
            "last_name": "Kretchmar",
            "email": "[email protected]"
        },
        "source": {
            "version": "3.6.3.Final",
            "connector": "postgresql",
            "name": "PostgreSQL_server",
            "ts_ms": 1559033904863,
            "ts_us": 1559033904863769,
            "ts_ns": 1559033904863769000,
            "snapshot": false,
            "db": "postgres",
            "schema": "public",
            "table": "customers",
            "txId": 556,
            "lsn": 24023128,
            "xmin": null
        },
        "op": "u",
        "ts_ms": 1465584025523,
        "ts_us": 1465584025523514,
        "ts_ns": 1465584025523514964,
    }
}

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

before

一个可选字段,包含数据库提交前该行中的值。在本示例中,只存在主键列 id,因为该表的 REPLICA IDENTITY 设置默认为 DEFAULT。

如果希望 update 事件包含该行所有列的先前值,就必须通过执行 ALTER TABLE customers REPLICA IDENTITY FULL 来修改 customers 表。

after

一个可选字段,指定事件发生后该行的状态。在本示例中,first_name 的值现在是 Anne Marie。

source

必填字段,描述事件的来源元数据。source 字段的结构与 create 事件中的结构相同,但部分值有所不同。来源元数据包括:

  • Debezium 版本
  • 连接器类型与名称
  • 包含新行的数据库和表
  • 模式名称
  • 该事件是否属于快照的一部分(对于 update 事件始终为 false)
  • 执行该操作所属事务的 ID
  • 该操作在数据库日志中的偏移量
  • 该变更在数据库中产生的时间戳

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 事件和一个带有该行旧键的墓碑事件,随后是一个带有该行新键的事件。有关更多信息,请参阅主键更新。

主键更新

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

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

delete 事件

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

{
    "schema": { ... },
    "payload": {
        "before": {
            "id": 1
        },
        "after": null,
        "source": {
            "version": "3.6.3.Final",
            "connector": "postgresql",
            "name": "PostgreSQL_server",
            "ts_ms": 1559033904863,
            "ts_us": 1559033904863852,
            "ts_ns": 1559033904863852000,
            "snapshot": false,
            "db": "postgres",
            "schema": "public",
            "table": "customers",
            "txId": 556,
            "lsn": 46523128,
            "xmin": null
        },
        "op": "d",
        "ts_ms": 1465581902461,
        "ts_us": 1465581902461496,
        "ts_ns": 1465581902461496187,
    }
}

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

before

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

在此示例中,before 字段只包含主键列,因为该表的 REPLICA IDENTITY 设置为 DEFAULT。

after

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

source

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

  • Debezium 版本
  • 连接器类型和名称
  • 包含被删除行的数据库和表
  • 模式名称
  • 该事件是否属于快照的一部分(对于 delete 事件始终为 false)
  • 执行该操作所属事务的 ID
  • 该操作在数据库日志中的偏移量
  • 该变更在数据库中发生的时间戳

op

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

ts_ms、ts_us、ts_ns

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

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

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

为了使消费者能够处理为没有主键的表生成的 delete 事件,请将该表的 REPLICA IDENTITY 设置为 FULL。当表没有主键且该表的 REPLICA IDENTITY 设置为 DEFAULT 或 NOTHING 时,delete 事件将不包含 before 字段。

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

Tombstone 事件

当删除一行时,删除事件的值仍可用于日志压缩,因为 Kafka 可以删除所有具有相同键的更早消息。不过,要让 Kafka 删除所有具有相同键的消息,消息的值必须为 null。为实现这一点,PostgreSQL 连接器会在 删除事件之后发送一个特殊的 tombstone 事件,该事件具有相同的键,但值为 null。

truncate 事件

truncate 变更事件表示某个表已被清空。此时消息键为 null,消息值的形式如下:

{
    "schema": { ... },
    "payload": {
        "source": {
            "version": "3.6.3.Final",
            "connector": "postgresql",
            "name": "PostgreSQL_server",
            "ts_ms": 1559033904863,
            "ts_us": 1559033904863112,
            "ts_ns": 1559033904863112000,
            "snapshot": false,
            "db": "postgres",
            "schema": "public",
            "table": "customers",
            "txId": 556,
            "lsn": 46523128,
            "xmin": null
        },
        "op": "t",
        "ts_ms": 1559033904961,
        "ts_us": 1559033904961654,
        "ts_ns": 1559033904961654789
    }
}

以下列表描述了前面的 truncate(截断)事件值示例中的部分字段:

source

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

  • Debezium 版本
  • 连接器类型和名称
  • 包含新行的数据库和表
  • 架构(schema)名称
  • 该事件是否属于快照的一部分(对于 delete 事件始终为 false)
  • 执行该操作所在的事务 ID
  • 该操作在数据库日志中的偏移量(offset)
  • 变更在数据库中发生的时间戳

op

描述操作类型的必填字符串。op 字段的值为 t,表示该表已被截断。

ts_ms、ts_us、ts_ns

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

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

如果单次 TRUNCATE 操作影响多张表,连接器会为每张被截断的表发出一条 truncate 变更事件记录。

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

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

message 事件

message 事件用于捕获通过 pg_logical_emit_message 函数直接写入 WAL 的通用逻辑解码消息。以下示例展示了事务性消息和非事务性消息的事件结构,并说明了各自的字段。

消息键是一个仅包含单个字段 prefix 的 Struct,该字段承载插入消息时指定的前缀。下面的 JSON 展示了一个事务性消息的示例消息值:

此事件类型仅在 Postgres 14+ 上通过 pgoutput 插件支持(PostgreSQL 文档)
{
    "schema": { ... },
    "payload": {
        "source": {
            "version": "3.6.3.Final",
            "connector": "postgresql",
            "name": "PostgreSQL_server",
            "ts_ms": 1559033904863,
            "ts_us": 1559033904863879,
            "ts_ns": 1559033904863879000,
            "snapshot": false,
            "db": "postgres",
            "schema": "",
            "table": "",
            "txId": 556,
            "lsn": 46523128,
            "xmin": null
        },
        "op": "m",
        "ts_ms": 1559033904961,
        "ts_us": 1559033904961621,
        "ts_ns": 1559033904961621379,
        "message": {
            "prefix": "example",
            "content": "ZXhhbXBsZS1tZXNzYWdl"
        }
    }
}

与其他事件类型不同,非事务性消息不会包含任何关联的 BEGIN 或 END 事务事件。非事务性消息的消息值如下所示:

{
    "schema": { ... },
    "payload": {
        "source": {
            "version": "3.6.3.Final",
            "connector": "postgresql",
            "name": "PostgreSQL_server",
            "ts_ms": 1559033904863,
            "ts_us": 1559033904863762,
            "ts_ns": 1559033904863762000,
            "snapshot": false,
            "db": "postgres",
            "schema": "",
            "table": "",
            "lsn": 46523128,
            "xmin": null
        },
        "op": "m",
        "ts_ms": 1559033904961,
        "ts_us": 1559033904961741,
        "ts_ns": 1559033904961741698,
        "message": {
            "prefix": "example",
            "content": "ZXhhbXBsZS1tZXNzYWdl"
    }
}

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

source

描述事件源元数据的必填字段。在 message 事件值中,source 字段结构不包含任何 message 事件的 table 或 schema 信息,且只有在 message 事件属于事务时才包含 txId。

  • Debezium 版本
  • 连接器类型和名称
  • 数据库名称
  • Schema 名称(message 事件始终为 "")
  • 表名(message 事件始终为 "")
  • 该事件是否属于快照的一部分(message 事件始终为 false)
  • 执行该操作所属事务的 ID(非事务性 message 事件为 null)
  • 该操作在数据库日志中的偏移量
  • 事务消息:消息写入 WAL 的时间戳
  • 非事务消息:连接器遇到该消息的时间戳

op

描述操作类型的必填字符串。op 字段的值为 m,表示这是一个 message 事件。

ts_ms、ts_us、ts_ns

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

对于事务性 message 事件,source 对象的 ts_ms 属性表示在数据库中进行变更的时间。通过比较 payload.source.ts_ms 与 payload.ts_ms 的值,可以确定源数据库更新与 Debezium 之间的延迟。

对于非事务性 message 事件,source 对象的 ts_ms 表示连接器遇到 message 事件的时间,而 payload.ts_ms 表示连接器处理该事件的时间。之所以存在这一差异,是因为 Postgres 的通用逻辑消息格式中不包含提交时间戳,且非事务性逻辑消息之前没有带时间信息的 BEGIN 事件。

message

包含消息元数据的字段。

数据类型映射

PostgreSQL 连接器使用与行所在表结构相同的事件来表示行的变更。事件中为每个列值包含一个字段。该值在事件中的表示方式取决于该列的 PostgreSQL 数据类型。以下各节介绍连接器如何将 PostgreSQL 数据类型映射为事件字段中的字面类型和语义类型。

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

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

基本类型

下表描述了连接器如何映射基本类型。

表 5. PostgreSQL 基本数据类型的映射

PostgreSQL 数据类型字面类型(模式类型)语义类型(模式名称)及说明
BOOLEANBOOLEAN无
BIT(1)BOOLEAN无
BIT( > 1)BYTESio.debezium.data.Bits。length 模式参数包含一个整数,表示位数。生成的 byte[] 以小端序形式包含这些位,其大小足以容纳指定的位数。例如,numBytes = n/8 + (n % 8 == 0 ? 0 : 1),其中 n 是位数。
BIT VARYING[(M)]BYTESio.debezium.data.Bits。length 模式参数包含一个整数,表示位数(如果列未指定长度,则为 2^31 - 1)。生成的 byte[] 以小端序形式包含这些位,其大小根据内容确定。指定的大小 (M) 存储在 io.debezium.data.Bits 类型的 length 参数中。
SMALLINT, SMALLSERIALINT16无
INTEGER, SERIALINT32无
BIGINT, BIGSERIAL, OIDINT64无
REALFLOAT32无
DOUBLE PRECISIONFLOAT64无
CHAR[(M)]STRING无
VARCHAR[(M)]STRING无
CHARACTER[(M)]STRING无
BPCHAR[(M)]STRING无
CHARACTER VARYING[(M)]STRING无
TIMESTAMPTZ, TIMESTAMP WITH TIME ZONESTRINGio.debezium.time.ZonedTimestamp。带时区信息的时间戳的字符串表示,时区为 GMT。
TIMETZ, TIME WITH TIME ZONESTRINGio.debezium.time.ZonedTime。带时区信息的时间值的字符串表示,时区为 GMT。
INTERVAL [P]INT64io.debezium.time.MicroDuration(默认)。时间间隔的近似微秒数,使用 365.25 / 12.0 公式计算每月平均天数。
INTERVAL [P]STRINGio.debezium.time.Interval(当 interval.handling.mode 设置为 string 时)。遵循 P<年>Y<月>M<日>DT<小时>H<分钟>M<秒>S 模式的间隔值字符串表示,例如 P1Y2M3DT4H5M6.78S。
BYTEABYTES 或 STRING无。根据连接器的二进制处理模式设置,取原始字节(默认)、base64 编码字符串、base64 URL 安全编码字符串或十六进制编码字符串。

Debezium 仅支持 Postgres 的 bytea_output 配置为 hex 值。有关 PostgreSQL 二进制数据类型的更多信息,请参阅 PostgreSQL 文档。

JSON、JSONB

STRING

io.debezium.data.Json

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

XML

STRING

io.debezium.data.Xml

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

UUID

STRING

io.debezium.data.Uuid

包含 PostgreSQL UUID 值的字符串表示形式。

POINT

STRUCT

io.debezium.data.geometry.Point

包含一个具有两个 FLOAT64 字段 (x,y) 的结构。每个字段表示几何点的坐标。

LTREE

STRING

io.debezium.data.Ltree

包含 PostgreSQL LTREE 值的字符串表示形式。

CITEXT

STRING

不适用

INET

STRING

不适用

INT4RANGE

STRING

不适用

整数的范围。

INT8RANGE

STRING

不适用

bigint 的范围。

NUMRANGE

STRING

不适用

numeric 的范围。

TSRANGE

STRING

不适用

包含不含时区的时间戳范围的字符串表示形式。

TSTZRANGE

STRING

不适用

包含带本地系统时区的时间戳范围的字符串表示形式。

DATERANGE

STRING

不适用

包含日期范围的字符串表示形式。其上界始终为开区间(不含)。

ENUM

STRING

io.debezium.data.Enum

包含 PostgreSQL ENUM 值的字符串表示形式。allowed 架构参数包含以逗号分隔的允许值列表。这些值按照 PostgreSQL 中定义的逻辑排序顺序(enumsortorder)排列,而不是按字母顺序排列。这样,即使后续向枚举中插入了新值,架构也能保持确定性。

时间类型

除了包含时区信息的 PostgreSQL TIMESTAMPTZ 和 TIMETZ 数据类型之外,时间类型的映射方式取决于 time.precision.mode 连接器配置属性的值。以下各节描述了这些映射方式:

time.precision.mode=adaptive

当 time.precision.mode 属性设置为默认值 adaptive 时,连接器会根据列的数据类型定义来确定字面量类型和语义类型。这样可以确保事件精确地表示数据库中的值。

表 6. time.precision.mode 为 adaptive 时的映射

PostgreSQL 数据类型 字面量类型(schema 类型) 语义类型(schema 名称)及说明
DATE INT32 io.debezium.time.Date 表示自纪元以来的天数。
TIME(1)、TIME(2)、TIME(3) INT32 io.debezium.time.Time 表示自午夜起的毫秒数,不包含时区信息。
TIME(4)、TIME(5)、TIME(6) INT64 io.debezium.time.MicroTime 表示自午夜起的微秒数,不包含时区信息。
TIMESTAMP(1)、TIMESTAMP(2)、TIMESTAMP(3) INT64 io.debezium.time.Timestamp 表示自纪元起的毫秒数,不包含时区信息。
TIMESTAMP(4)、TIMESTAMP(5)、TIMESTAMP(6)、TIMESTAMP INT64 io.debezium.time.MicroTimestamp 表示自纪元起的微秒数,不包含时区信息。

time.precision.mode=adaptive_time_microseconds

当 time.precision.mode 配置属性设置为 adaptive_time_microseconds 时,连接器会根据列的数据类型定义来确定时间类型的字面量类型和语义类型。这样可以确保事件精确地表示数据库中的值,区别仅在于所有 TIME 字段都以微秒为单位采集。

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

PostgreSQL 数据类型 字面量类型(schema 类型) 语义类型(schema 名称)及说明
DATE INT32 io.debezium.time.Date 表示自纪元以来的天数。
TIME([P]) INT64 io.debezium.time.MicroTime 以微秒表示时间值,不包含时区信息。PostgreSQL 允许精度 P 的取值范围为 0-6,最多可存储到微秒精度。
TIMESTAMP(1)、TIMESTAMP(2)、TIMESTAMP(3) INT64 io.debezium.time.Timestamp 表示自纪元起的毫秒数,不包含时区信息。
TIMESTAMP(4)、TIMESTAMP(5)、TIMESTAMP(6)、TIMESTAMP INT64 io.debezium.time.MicroTimestamp 表示自纪元起的微秒数,不包含时区信息。

time.precision.mode=connect

当 time.precision.mode 配置属性设置为 connect 时,连接器使用 Kafka Connect 逻辑类型。当消费者只能处理 Kafka Connect 内置的逻辑类型、而无法处理可变精度的时间值时,这种方式可能很有用。但是,由于 PostgreSQL 支持微秒精度,当数据库列的小数秒精度值大于 3 时,使用 connect 时间精度模式的连接器所产生的事件会导致精度丢失。

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

PostgreSQL 数据类型 字面量类型(schema 类型) 语义类型(schema 名称)及说明
DATE INT32 org.apache.kafka.connect.data.Date 表示自纪元(epoch)以来的天数。
TIME([P]) INT64 org.apache.kafka.connect.data.Time 表示自午夜以来的毫秒数,不包含时区信息。PostgreSQL 允许 P 的取值范围为 0-6,可存储高达微秒级的精度,但当 P 大于 3 时,此模式会导致精度损失。
TIMESTAMP([P]) INT64 org.apache.kafka.connect.data.Timestamp 表示自纪元(epoch)以来的毫秒数,不包含时区信息。PostgreSQL 允许 P 的取值范围为 0-6,可存储高达微秒级的精度,但当 P 大于 3 时,此模式会导致精度损失。

time.precision.mode=isostring

将 time.precision.mode 属性设置为 isostring,可配置连接器以 UTC 时区的 ISO-8601 格式字符串映射时间值。应用此设置后,连接器使用语义类型 io.debezium.time.IsoTimestamp、io.debezium.time.IsoTime 和 io.debezium.time.IsoDate 来映射时间戳、日期时间、日期和时间值。

表 9. time.precision.mode 为 isostring 时的映射关系

PostgreSQL 数据类型 字面量类型(schema 类型) 语义类型(schema 名称)及说明
DATE STRING io.debezium.time.IsoDate 以 UTC 格式表示日期值,遵循 ISO 8601 标准,例如 2017-09-15Z。
TIME([P]) STRING io.debezium.time.IsoTime 以 UTC 格式表示时间值,遵循 ISO 8601 标准,例如 04:05:11.789Z。
TIMESTAMP([P]) STRING io.debezium.time.IsoTimestamp 以 UTC 格式表示时间戳值,遵循 ISO 8601 标准,例如 2019-07-09T02:28:57.123456Z。

time.precision.mode=microseconds

将 time.precision.mode 属性设置为 microseconds,可配置连接器以微秒精度表示时间值。应用此设置后,连接器使用语义类型 io.debezium.time.MicroTime 和 io.debezium.time.MicroTimestamp 来映射时间戳、日期时间和时间值。

表 10. time.precision.mode 为 microseconds 时的映射关系

PostgreSQL 数据类型 字面量类型(schema 类型) 语义类型(schema 名称)及说明
DATE INT32 io.debezium.time.Date 表示自纪元(epoch)以来的天数。
TIME([P]) INT64 io.debezium.time.MicroTime 以微秒表示时间值,不包含时区信息。在 PostgreSQL 中,精度参数 p 指定时间值中秒部分的小数位数。精度取值范围为 0(无小数秒)到 6(微秒精度)。
TIMESTAMP([P]) INT64 io.debezium.time.MicroTimestamp

表示自纪元以来的微秒数,不包含时区信息。在 PostgreSQL 中,精度参数 p 指定时间值中秒部分的小数位数。精度范围可从 0(无小数秒)到 6(微秒精度)。

time.precision.mode=nanoseconds

将 time.precision.mode 属性设置为 nanoseconds,可配置连接器以纳秒精度表示时间值。应用此设置后,连接器会使用语义类型 io.debezium.time.NanoTime 和 io.debezium.time.NanoTimestamp(它们以纳秒精度存储值)来映射时间戳、日期时间和时间值。

表 11. time.precision.mode 为 nanoseconds 时的映射

PostgreSQL 数据类型 字面量类型(schema 类型) 语义类型(schema 名称)及说明

DATE

INT32

io.debezium.time.Date

表示自纪元以来的天数。

TIME([P])

INT64

io.debezium.time.NanoTime

以纳秒表示时间值,不包含时区信息。在 PostgreSQL 中,精度参数 p 指定时间值中秒部分的小数位数。精度范围可从 0(无小数秒)到 6(微秒精度)。

TIMESTAMP([P])

INT64

io.debezium.time.NanoTimestamp

表示自纪元以来的纳秒数,不包含时区信息。在 PostgreSQL 中,精度参数 p 指定时间值中秒部分的小数位数。精度范围可从 0(无小数秒)到 6(微秒精度)。

PostgreSQL 支持在 TIMESTAMP 列中使用 +/-infinite 值。当 time.precision.mode 设置为 nanoseconds 时,这些特殊值会被转换为 Long.MAX_VALUE(9223372036854775807)表示正无穷大,转换为 Long.MIN_VALUE(-9223372036854775808)表示负无穷大,以防止在以纳秒表示时间戳时发生溢出。

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

如果源模式中可能包含超出可表示范围的值,请改用 time.precision.mode=microseconds(大约覆盖 ±290000 年)或 time.precision.mode=isostring。

TIMESTAMP 类型

TIMESTAMP 类型表示不带时区信息的时间戳。这类列会被转换为基于 UTC 的等效 Kafka Connect 值。例如,当 time.precision.mode 未设置为 connect 时,TIMESTAMP 值 "2018-06-20 15:13:16.945104" 会由值为 "1529507596945104" 的 io.debezium.time.MicroTimestamp 来表示。

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

PostgreSQL 支持在 TIMESTAMP 列中使用 +/-infinite 值。这些特殊值会被转换为时间戳:正无穷对应值 9223372036825200000,负无穷对应值 -9223372036832400000。此行为与 PostgreSQL JDBC 驱动的标准行为一致。供参考,请参见 org.postgresql.PGStatement 接口。

十进制类型

PostgreSQL 连接器配置属性 decimal.handling.mode 的设置决定了连接器如何映射十进制类型。

当 decimal.handling.mode 属性设置为 precise 时,连接器对所有 DECIMAL、NUMERIC 和 MONEY 列使用 Kafka Connect 的 org.apache.kafka.connect.data.Decimal 逻辑类型。这是默认模式。

表 12. decimal.handling.mode 为 precise 时的映射

PostgreSQL 数据类型 字面类型(schema 类型)语义类型(schema 名称)及说明

NUMERIC[(M[,D])]

BYTES

org.apache.kafka.connect.data.Decimal

scale schema 参数包含一个整数,表示小数点向右移动了多少位。

DECIMAL[(M[,D])]

BYTES

org.apache.kafka.connect.data.Decimal

scale schema 参数包含一个整数,表示小数点向右移动了多少位。

MONEY[(M[,D])]

BYTES

org.apache.kafka.connect.data.Decimal

scale schema 参数包含一个整数,表示小数点向右移动了多少位。scale schema 参数由连接器配置属性 money.fraction.digits 确定。

此规则有一个例外。当 NUMERIC 或 DECIMAL 类型未使用精度(scale)约束时,来自数据库的每个值具有不同的(可变的)精度。在这种情况下,连接器使用 io.debezium.data.VariableScaleDecimal,它同时包含所传输值的值和精度。

表 13. 无精度约束时 DECIMAL 和 NUMERIC 类型的映射

PostgreSQL 数据类型 字面类型(schema 类型)语义类型(schema 名称)及说明

NUMERIC

STRUCT

io.debezium.data.VariableScaleDecimal

包含一个具有两个字段的结构:类型为 INT32 的 scale,其中包含所传输值的精度;以及类型为 BYTES 的 value,其中包含未缩放形式的原始值。

DECIMAL

STRUCT

io.debezium.data.VariableScaleDecimal

该类型包含一个结构,其中有两个字段:INT32 类型的 scale,用于存放所传输值的标度;以及 BYTES 类型的 value,用于存放未按标度缩放的原始值。

当 decimal.handling.mode 属性设置为 double 时,连接器会将所有 DECIMAL、NUMERIC 和 MONEY 值表示为 Java double 值,并按下表所示进行编码。

PostgreSQL 数据类型字面量类型(schema 类型)语义类型(schema 名称)
NUMERIC[(M[,D])]FLOAT64
DECIMAL[(M[,D])]FLOAT64
MONEY[(M[,D])]FLOAT64

表 14. decimal.handling.mode 为 double 时的映射关系

decimal.handling.mode 配置属性的最后一个可选设置是 string。在这种情况下,连接器会将 DECIMAL、NUMERIC 和 MONEY 值表示为其格式化后的字符串形式,并按下表所示进行编码。

PostgreSQL 数据类型字面量类型(schema 类型)语义类型(schema 名称)
NUMERIC[(M[,D])]STRING
DECIMAL[(M[,D])]STRING
MONEY[(M[,D])]STRING

表 15. decimal.handling.mode 为 string 时的映射关系

当 decimal.handling.mode 设置为 string 或 double 时,PostgreSQL 支持将 NaN(非数字)作为特殊值存储在 DECIMAL/NUMERIC 值中。在这种情况下,连接器会将 NaN 编码为 Double.NaN 或字符串常量 NAN。

HSTORE 类型

PostgreSQL 连接器配置属性 hstore.handling.mode 的设置决定了连接器如何映射 HSTORE 值。

当 hstore.handling.mode 属性设置为 json(默认值)时,连接器会将 HSTORE 值表示为 JSON 值的字符串形式,并按下表所示进行编码。当 hstore.handling.mode 属性设置为 map 时,连接器会对 HSTORE 值使用 MAP schema 类型。

表 16. HSTORE 数据类型的映射关系

PostgreSQL 数据类型 字面量类型(schema 类型) 语义类型(schema 名称)及备注

HSTORE

STRING

io.debezium.data.Json

示例:使用 JSON 转换器时的输出表示为 {"key" : "val"}

HSTORE

MAP

不适用

示例:使用 JSON 转换器时的输出表示为 {"key" : "val"}

域类型

PostgreSQL 支持基于其他底层类型构建的用户自定义类型。当使用此类列类型时,Debezium 会根据完整的类型层次结构来呈现该列的表示形式。

捕获使用 PostgreSQL 域类型的列的变更需要特别注意。当某列被定义为包含一个扩展了数据库默认类型之一、且该域类型定义了自定义长度或精度的域类型时,生成的架构会继承所定义的长度或精度。

当某列被定义为包含一个扩展了另一个已定义自定义长度或精度的域类型的域类型时,生成的架构则不会继承所定义的长度或精度,因为 PostgreSQL 驱动程序的列元数据中并不包含该信息。

网络地址类型

PostgreSQL 提供了可用于存储 IPv4、IPv6 和 MAC 地址的数据类型。存储网络地址时,使用这些类型而不是纯文本类型更为合适。网络地址类型提供了输入错误检查以及专门的运算符和函数。

表 17. 网络地址类型的映射

PostgreSQL 数据类型 字面量类型(schema 类型) 语义类型(schema 名称)及说明

INET

STRING

不适用

IPv4 和 IPv6 网络

CIDR

STRING

不适用

IPv4 和 IPv6 主机及网络

MACADDR

STRING

不适用

MAC 地址

MACADDR8

STRING

不适用

采用 EUI-64 格式的 MAC 地址

PostGIS 类型

PostgreSQL 连接器支持所有 PostGIS 数据类型。

表 18. PostGIS 数据类型的映射

PostGIS 数据类型 字面量类型(schema 类型) 语义类型(schema 名称)及说明

GEOMETRY(平面)

STRUCT

io.debezium.data.geometry.Geometry

包含一个具有以下两个字段的结构:

  • srid (INT32) - 空间参考系标识符,用于定义结构中存储的几何对象的类型。
  • wkb (BYTES) - 采用 Well-Known-Binary 格式编码的几何对象的二进制表示。

有关格式详情,请参阅 Open Geospatial Consortium Simple Features Access 规范。

GEOGRAPHY(球面)

STRUCT

io.debezium.data.geometry.Geography

包含一个具有以下两个字段的结构:

  • srid (INT32) - 空间参考系标识符,用于定义结构中存储的地理对象的类型。
  • wkb (BYTES) - 采用 Well-Known-Binary 格式编码的几何对象的二进制表示。

有关格式详情,请参阅 Open Geospatial Consortium Simple Features Access 规范。

pgvector 类型

PostgreSQL 连接器支持所有 pgvector 扩展数据类型。

表 19. PostgreSQL pgvector 数据类型的映射

pgvector 数据类型 字面量类型(schema 类型) 语义类型(schema 名称)及说明

VECTOR

ARRAY (FLOAT64)

io.debezium.data.DoubleVector

HALFVEC

ARRAY (FLOAT32)

io.debezium.data.FloatVector

SPARSEVEC

STRUCT

io.debezium.data.SparseVector

包含一个包含以下字段的结构:

dimensions (INT16)

稀疏向量的总长度。

vector (MAP (INT16, FLOAT64))

表示稀疏向量的映射。每个映射值包含以下元素:

  • 向量元素的索引号(从 1 开始)。
  • 向量元素的值。

TSVECTOR 类型

PostgreSQL 连接器无需任何扩展即可支持原生的 TSVECTOR 数据类型。

表 20. PostgreSQL TSVECTOR 数据类型的映射

PostgreSQL 数据类型 字面量类型(schema 类型)语义类型(schema 名称)及说明

TSVECTOR

STRING

io.debezium.data.Tsvector

包含以规范化词素格式表示的 TSVECTOR 值的字符串。

示例:'direct':6 'insert':8 'test':4

TSVECTOR 类型是 PostgreSQL 的原生数据类型(OID 3614),用于全文搜索操作。连接器会将 TSVECTOR 值转换为其字符串表示形式,并保留词素及其位置信息。

TOAST 值

PostgreSQL 对页面大小有严格的限制。这意味着大于约 8 KB 的值需要使用 TOAST 存储来保存。这会影响来自数据库的复制消息。通过 TOAST 机制存储且未被修改的值不会包含在消息中,除非它们属于表的副本标识(replica identity)。Debezium 无法安全地以带外方式直接从数据库中读取缺失的值,因为这可能导致竞态条件。因此,Debezium 遵循以下规则来处理 TOAST 值:

  • 表的 REPLICA IDENTITY FULL —— TOAST 列的值与其他列一样,包含在变更事件的 before 和 after 字段中。
  • 表的 REPLICA IDENTITY DEFAULT —— 当从数据库接收到 UPDATE 事件时,任何不属于副本标识且未发生变化的 TOAST 列值都不包含在事件中。同样,当接收到 DELETE 事件时,before 字段中不包含任何 TOAST 列(如果存在的话)。由于 Debezium 在这种情况下无法安全地提供列值,连接器会返回由连接器配置属性 unavailable.value.placeholder 定义的占位值。

默认值

如果数据库架构中为某列指定了默认值,PostgreSQL 连接器将尽可能尝试将该值传播到 Kafka 架构中。最常见的数据类型都受支持,包括以下类型:

  • BOOLEAN
  • 数值类型(INT、FLOAT、NUMERIC 等)
  • 文本类型(CHAR、VARCHAR、TEXT 等)
  • 时间类型(DATE、TIME、INTERVAL、TIMESTAMP、TIMESTAMPTZ)
  • JSON、JSONB、XML
  • UUID

请注意,对于时间类型,默认值的解析由 PostgreSQL 库提供;因此,PostgreSQL 通常支持的任何字符串表示形式,连接器也应当支持。

如果默认值由函数生成而非直接内联指定,连接器会导出该数据类型对应的等效值。这些值包括:

  • BOOLEAN 对应 FALSE
  • 数值类型对应具有适当精度的 0
  • 文本/XML 类型对应空字符串
  • JSON 类型对应 {}
  • DATE、TIMESTAMP、TIMESTAMPTZ 类型对应 1970-01-01
  • TIME 对应 00:00
  • INTERVAL 对应 EPOCH
  • UUID 对应 00000000-0000-0000-0000-000000000000

目前这种支持仅限于显式使用函数的情况。例如,CURRENT_TIMESTAMP(6) 带有括号时受支持,而 CURRENT_TIMESTAMP 则不受支持。

支持传播默认值主要是为了在使用 PostgreSQL 连接器并配合会强制校验模式版本兼容性的模式注册中心时,能够安全地进行模式演进。由于这一主要考虑因素,再加上不同插件的刷新行为,Kafka 模式中出现的默认值并不保证始终与数据库模式中的默认值保持一致。

  • 默认值可能在 Kafka 模式中“延后”出现,具体取决于插件何时以及如何触发内存中模式的刷新。如果在两次刷新之间默认值发生了多次变化,那么该值可能永远不会出现在 Kafka 模式中,或被跳过。
  • 如果在连接器仍有待处理的记录时触发了模式刷新,默认值可能在 Kafka 模式中“提前”出现。这是因为在刷新时列元数据是从数据库中读取的,而不是包含在复制消息中。如果连接器处理滞后时发生了刷新,或者连接器被停止一段时间(在此期间源数据库仍在持续写入更新)后重新启动时,都可能出现这种情况。

这种行为可能出人意料,但仍然是安全的。只有模式定义会受到影响,而消息中实际存在的值仍与写入源数据库的值保持一致。

自定义转换器

默认情况下,Debezium 不会复制具有自定义数据类型的列中的数据,例如使用 SQL CREATE TYPE 语句创建的复合类型。要复制具有自定义数据类型的列,请按照创建自定义转换器中的说明进行操作,但需注意以下几点重要事项:

  • 在连接器配置中将 include.unknown.datatypes 属性设置为 true。默认的 false 设置会导致自定义转换器始终返回 null 值。

  • 传递给转换器的值的类型取决于为复制槽配置的逻辑解码输出插件。

    • decoderbufs 传递列数据的字节数组(byte[])表示形式。
    • pgoutput 传递列数据的字符串表示形式。

设置 PostgreSQL

在运行 Debezium PostgreSQL 连接器之前,需要先配置 PostgreSQL 以支持逻辑复制,设置所需的用户权限,并允许连接器主机的网络访问。

首先确定你打算使用的逻辑解码插件。如果你计划不使用原生的 pgoutput 逻辑复制流支持,则必须将逻辑解码插件安装到 PostgreSQL 服务器中。之后,配置一个具有足够权限来执行复制的用户。连接器启动时,如果复制槽尚不存在,会自动创建所配置的复制槽。

如果你的数据库托管在 Heroku Postgres 这类服务上,你可能无法安装该插件。如果是这样,并且你使用的是 PostgreSQL 10 或更高版本,则可以使用 pgoutput 解码器支持来捕获数据库中的更改。如果这也不是一个选项,你就无法将 Debezium 与你的数据库配合使用。

配置复制槽

PostgreSQL 逻辑解码使用复制槽。要让 PostgreSQL 能够承载 Debezium 使用的复制槽,请在 postgresql.conf 文件中指定以下内容:

wal_level=logical
max_wal_senders=1
max_replication_slots=1

这些设置对 PostgreSQL 服务器作出如下指示:

  • wal_level - 使用预写日志(WAL)进行逻辑解码。
  • max_wal_senders - 最多使用一个独立进程来处理 WAL 变更。
  • max_replication_slots - 最多允许创建一个用于流式传输 WAL 变更的复制槽。

复制槽可保证保留 Debezium 所需的全部 WAL 条目,即使在 Debezium 停机期间也是如此。因此,密切监控复制槽非常重要,以避免出现以下问题:

  • 磁盘占用过多
  • 复制槽长期未被使用而引发的任何状况,例如系统目录膨胀

有关更多信息,请参阅PostgreSQL 复制槽文档。

熟悉 PostgreSQL 预写日志的机制与配置 有助于使用 Debezium PostgreSQL 连接器。

自动创建复制槽

默认情况下,当 Debezium PostgreSQL 连接器启动时,如果所配置的复制槽不存在,连接器会自动创建它。要使自动创建槽成功,必须满足以下先决条件:

  • 为连接器配置的数据库用户必须具有 REPLICATION 权限。
  • 使用 pgoutput 插件时,数据库用户必须在该数据库上具有 CREATE 权限(创建发布所必需)。
  • PostgreSQL 必须配置有足够的复制槽容量。例如,max_replication_slots 必须设置为至少 1。
  • 长时间运行的事务不会阻塞槽的创建。

如果你希望对复制槽的设置有更多的控制权,尤其是在生产环境中,可以在启动连接器之前手动创建复制槽。在以下运维场景中,也必须手动创建槽:

  • 受控的数据库升级流程,你需要在升级过程中保留变更数据。
  • 故障切换恢复场景,你需要在恢复数据库写入之前重新创建槽。
  • 连接器用户不具备 REPLICATION 或 CREATE 权限的环境。

手动创建复制槽

你可以手动创建 PostgreSQL 复制槽。如果由于权限不足而无法自动创建槽,或者你希望更好地控制槽的创建,可以使用此选项。

操作步骤

  1. 输入以下 SQL 命令以手动创建复制槽:

    SELECT pg_create_logical_replication_slot('slot_name', 'plugin_name');

将 plugin_name 替换为你所使用的逻辑解码插件的名称(pgoutput 或 decoderbufs)。

将 slot_name 替换为连接器 slot.name 属性的值。

如果你部署了多个 Debezium PostgreSQL 连接器,则必须为每个连接器分配一个唯一的槽名称。请为每个连接器实例配置具有不同值的 slot.name 属性。

云端的 PostgreSQL

当你在 Amazon RDS、Azure Database for PostgreSQL、Google Cloud SQL 或 CrunchyBridge 等托管云平台上运行 PostgreSQL 时,需要执行额外的步骤才能为 Debezium 启用逻辑复制。

Amazon RDS 上的 PostgreSQL

可以在运行于 Amazon RDS 上的 PostgreSQL 数据库中捕获更改。具体做法如下:

  • 将实例参数 rds.logical_replication 设置为 1。
  • 以数据库 RDS 主用户身份运行查询 SHOW wal_level,验证 wal_level 参数是否设置为 logical。在多可用区复制配置中,该参数可能并非如此。你无法手动设置此选项,当 rds.logical_replication 参数设置为 1 时,它会被自动更改。如果在完成上述更改后 wal_level 仍不是 logical,很可能是因为参数组更改后需要重启实例。重启会在你的维护窗口期间进行,你也可以手动发起重启。
  • 将 Debezium 的 plugin.name 参数设置为 pgoutput。
  • 从具有 rds_replication 角色的 AWS 账户发起逻辑复制。该角色授予管理逻辑复制槽以及使用逻辑复制槽流式传输数据的权限。默认情况下,在 Amazon RDS 上只有 AWS 主用户账户拥有 rds_replication 角色。若要让主用户账户之外的用户账户发起逻辑复制,你必须为该账户授予 rds_replication 角色,例如 grant rds_replication to <my_user>。授予 rds_replication 角色需要具有 superuser 权限。若要让主用户账户之外的账户创建初始快照,你必须为这些账户授予所要捕获表的 SELECT 权限。有关 PostgreSQL 逻辑复制安全性的更多信息,请参阅 PostgreSQL 文档。
  • 若要使用 AWS IAM 认证,请将 database.connection.factory.class 设置为 io.debezium.connector.postgresql.connection.PostgresAwsIamConnectionFactory,并将 AWS Advanced JDBC Wrapper 添加到你的 classpath 中。更多文档请见此处。

Azure 上的 PostgreSQL

Debezium 可以配合 Azure Database for PostgreSQL 使用,该服务支持 Debezium 所支持的 pgoutput 逻辑解码插件。

将 Azure 的复制支持设置为 logical。你可以使用 Azure CLI 或 Azure 门户进行配置。例如,若使用 Azure CLI,你需要执行以下 az postgres server 命令:

az postgres server configuration set --resource-group mygroup --server-name myserver --name azure.replication_support --value logical

az postgres server restart --resource-group mygroup --name myserver

Cloud SQL for PostgreSQL

要将 Debezium 与 Cloud SQL for PostgreSQL 配合使用,必须将数据库配置为使用 pgoutput 逻辑解码插件。

以下各节概述了在准备 Cloud SQL for PostgreSQL 数据库以配合 Debezium 使用时必须完成的任务。

设置 cloudsql.logical_decoding 标志

在 Cloud SQL 中,可以通过将 cloudsql.logical_decoding 标志设置为 on 来启用逻辑解码。设置该标志后,系统会自动将 wal_level 配置参数调整为 logical。

你可以使用 Google Cloud 控制台或 gcloud 命令行工具来修改 cloudsql.logical_decoding 标志。

有关如何更改 Cloud SQL 中标志值的详细说明,请参阅 Google Cloud SQL 文档。

要验证该设置的值是否反映了你的更改,请运行以下查询:

SHOW wal_level;

CrunchyBridge 上的 PostgreSQL

Debezium 可以与 CrunchyBridge 搭配使用;逻辑复制已启用,pgoutput 插件也可用。你需要创建一个复制用户并授予相应的权限。

使用 pgoutput 插件时,建议将 publication.autocreate.mode 配置为 filtered。如果使用 all_tables(这是 publication.autocreate.mode 的默认值),且未找到发布,连接器会尝试执行 CREATE PUBLICATION <publication_name> FOR ALL TABLES; 来创建发布,但由于权限不足,该操作会失败。

安装逻辑解码输出插件

在 PostgreSQL 服务器上安装逻辑解码输出插件,Debezium 连接器才能接收变更事件。pgoutput 插件内置于 PostgreSQL 中,无需安装;其他插件(如 decoderbufs)则需要单独编译和安装。

有关设置和测试逻辑解码插件的更详细说明,请参阅 PostgreSQL 的逻辑解码输出插件安装。

从 PostgreSQL 9.4 开始,读取预写日志(WAL)变更的唯一方式就是安装逻辑解码输出插件。插件使用 C 语言编写,需编译后安装到运行 PostgreSQL 服务器的机器上。插件使用了大量 PostgreSQL 特有的 API,具体说明请参见 PostgreSQL 文档。

PostgreSQL 连接器使用 Debezium 支持的逻辑解码插件之一,以 Protobuf 格式或 pgoutput 格式从数据库接收变更事件。pgoutput 插件随 PostgreSQL 数据库一同提供。有关通过 decoderbufs 插件使用 Protobuf 的更多信息,请参阅该插件的文档,其中讨论了它的要求、限制以及编译方法。

为简化操作,Debezium 还提供了一个基于上游 PostgreSQL 服务器镜像的容器镜像,并在其之上编译和安装这些插件。你可以使用此镜像作为安装所需详细步骤的示例。

Debezium 的逻辑解码插件仅在 Linux 机器上安装和测试过。对于 Windows 和其他操作系统,可能需要不同的安装步骤。

插件差异

不同插件在所有情况下的行为并不完全相同。已识别出以下差异:

  • 虽然所有插件都会在流式传输期间检测到结构变更时从数据库刷新结构元数据,但 pgoutput 插件在触发此类刷新方面更为“积极”。例如,列默认值的更改会触发 pgoutput 的刷新,而其他插件在另一次变更触发刷新之前(例如新增列)不会感知到这一更改。这源于 pgoutput 的行为,而非 Debezium 本身的行为。

所有最新的差异都记录在一个测试套件 Java 类中。

指定 Debezium 的 plugin.name

设置 plugin.name 连接器属性,以指定连接器使用哪个逻辑解码插件从 PostgreSQL 复制流中读取变更事件。

在需要连接器使用 pgoutput 逻辑解码插件的环境中(例如 Google Cloud SQL for PostgreSQL 或 Amazon RDS 上的 PostgreSQL),你必须将连接器的 plugin.name 属性的值显式设置为 pgoutput。

操作步骤

  • 在 Debezium 连接器配置中,将 plugin.name 属性的值设置为 pgoutput,如以下示例所示:

    {
          ..
          "plugin.name": "pgoutput",
          ..
          ..
    }

创建复制用户

要使用逻辑解码功能,必须创建一个具有 REPLICATION 属性的 PostgreSQL 用户,或者将该属性授予现有用户。

操作步骤

  • 完成以下步骤之一:

    创建具有 REPLICATION 属性的用户

    以 postgres 用户身份登录,或以 cloudsqlsuperuser 用户组成员的身份登录,然后运行以下命令:

    CREATE USER replication_user WITH REPLICATION IN ROLE cloudsqlsuperuser LOGIN PASSWORD 'secret';

为现有用户设置 REPLICATION 属性

以 postgres 用户或 cloudsqlsuperuser 用户组成员的身份登录,然后运行以下命令:

ALTER USER existing_user WITH REPLICATION;

配置 PostgreSQL 服务器

配置 PostgreSQL 服务器以支持逻辑解码,并提供 Debezium 连接器可靠捕获变更事件所需的复制参数。

如果你使用的是 pgoutput 以外的逻辑解码插件,则在安装插件后,按下述方式配置 PostgreSQL 服务器:

  1. 若要在启动时加载该插件,请在 postgresql.conf 文件中添加以下内容:

    # MODULES
    shared_preload_libraries = 'decoderbufs'

上一条语句指示 PostgreSQL 服务器在启动时加载 decoderbufs 逻辑解码插件(插件名称在 Protobuf 的 make 文件中设置)。
2. 无论使用哪种解码器,要配置复制槽,请在 postgresql.conf 文件中指定以下内容:

# REPLICATION
wal_level = logical

上述语句指示 PostgreSQL 服务器使用预写日志(WAL)进行逻辑解码。

复制参数与性能注意事项

根据你的需求,可以设置其他 PostgreSQL 流复制参数,以便在 Debezium 环境中自定义性能。例如,可以通过设置 max_wal_senders 和 max_replication_slots 参数来增加可并发访问发送服务器的连接器数量。也可以设置 wal_keep_size 来限制复制槽的最大 WAL 大小。

Debezium 使用 PostgreSQL 的逻辑解码,而逻辑解码依赖复制槽。逻辑解码从 PostgreSQL 预写日志(WAL)中提取数据变更,并将其转换为 Debezium 可以消费的格式。即使在 Debezium 停机期间,复制槽也能保证保留 Debezium 所需的全部 WAL 段。因此,密切监控复制槽非常重要,以避免磁盘过度消耗,以及因复制槽长时间未被使用而导致的系统目录膨胀等问题。

当 PostgreSQL 的 synchronous_commit 参数设置为 on 以外的值时,为最小化变更事件的延迟,可将集群的 wal_writer_delay 设置为 10 毫秒之类的值。如果该属性未指定明确的值,Debezium 将采用 200 毫秒的默认值,这会增加延迟。

其他资源

为 PostgreSQL 连接器启用故障转移槽

当 Debezium 连接到运行 PostgreSQL 17 或更高版本的 PostgreSQL 主服务器时,你可以配置连接器创建支持故障转移的复制槽。启用故障转移的复制槽使 Debezium 能够在故障转移事件之后从新的主服务器继续读取变更,而不会丢失事件。

前提条件

  • Debezium 连接到运行 PostgreSQL 17 或更高版本的主服务器。
  • 你拥有配置 PostgreSQL 服务器参数的管理权限。
  • 你已配置 Debezium 连接器属性。

操作步骤

  1. 在 Debezium 连接器配置中,将 slot.failover 属性设置为 true。

    此设置使连接器能够创建支持故障转移的复制槽。

  2. 在 PostgreSQL 主服务器上,配置 synchronized_standby_slots 参数,使其包含连接器所使用的复制槽名称。

    主服务器使用该参数将复制槽状态与备用服务器进行同步。

    例如,如果你的复制槽名为 debezium,请在 postgresql.conf 文件中添加以下配置:

    synchronized_standby_slots = 'debezium'
  3. 重启 PostgreSQL 服务器以应用配置更改。

  4. 验证故障转移槽已创建并完成同步。

    你可以查询 pg_replication_slots 视图,确认该槽存在且 failover 列的值为 true:

    SELECT slot_name, failover, synced
    FROM pg_replication_slots
    WHERE slot_name = 'debezium';

结果

启用 slot.failover 属性并将插槽列入 synchronized_standby_slots 后,如果主服务器发生故障,备用副本被提升为新的主服务器,Debezium 会继续从新主服务器上指定的故障转移插槽读取变更。

其他资源

设置权限

要配置 PostgreSQL 服务器以运行 Debezium 连接器,需要一个能够执行复制操作的数据库用户。只有具备相应权限的数据库用户才能执行复制,并且只能针对已配置的主机进行复制。

尽管超级用户默认拥有必要的 REPLICATION 和 LOGIN 角色(如安全中所述),但最好不要为 Debezium 复制用户授予过高的权限。相反,应创建一个仅具有所需最低权限的 Debezium 用户。

先决条件

  • PostgreSQL 管理权限。

操作步骤

  1. 为用户提供复制权限,需要定义一个 PostgreSQL 角色,该角色至少具有 REPLICATION 和 LOGIN 权限,然后将该角色授予用户。例如:

    CREATE ROLE <name> REPLICATION LOGIN;

设置权限以启用 Debezium 在使用 pgoutput 时创建 PostgreSQL 发布

当使用 pgoutput 逻辑解码插件时,Debezium 会从 PostgreSQL 发布中流式传输变更事件。请为 Debezium 数据库用户授予所需的权限,并配置共享的表所有权,以便连接器能够为其捕获的表创建和管理发布。

如果你将 pgoutput 用作逻辑解码插件,Debezium 必须以具有特定权限的用户身份在数据库中运行。

发布包含由一个或多个表生成的经过筛选的变更事件集合。每个发布中的数据都依据发布规范进行筛选。该规范可以由 PostgreSQL 数据库管理员创建,也可以由 Debezium 连接器创建。

关于如何创建发布,有多种可选方式。通常,最好在设置连接器之前,手动为要捕获的表创建发布。不过,你也可以配置环境,使 Debezium 能够自动创建发布,并指定添加到发布中的数据。

Debezium 使用包含列表和排除列表属性来指定如何将数据插入发布中。有关启用 Debezium 创建发布的各项选项的更多信息,请参阅 publication.autocreate.mode。

Debezium 要创建 PostgreSQL 发布,必须以具有以下权限的用户身份运行:

  • 数据库中的复制权限,以便将表添加到发布中。
  • 数据库上的 CREATE 权限,以便创建发布。
  • 表上的 SELECT 权限,以便复制初始表数据。表所有者自动拥有该表的 SELECT 权限。

要将表添加到发布中,用户必须是该表的所有者。但由于源表已经存在,你需要一种机制来与原所有者共享所有权。要实现共享所有权,你可以创建一个 PostgreSQL 复制组,然后将现有的表所有者和复制用户添加到该组中。

操作步骤

  1. 创建一个复制组。

    CREATE ROLE <replication_group>;
  2. 将表的原所有者添加到该组。

    GRANT REPLICATION_GROUP TO <original_owner>;
  3. 将 Debezium 复制用户添加到该组。

    GRANT REPLICATION_GROUP TO <replication_user>;
  4. 将表的所有权转移给 <replication_group>。

    ALTER TABLE <table_name> OWNER TO REPLICATION_GROUP;

要让 Debezium 指定捕获配置,publication.autocreate.mode 的值必须设置为 filtered。

配置 PostgreSQL 以允许与 Debezium 连接器主机进行复制

要使 Debezium 能够复制 PostgreSQL 数据,必须配置数据库,允许其与运行 PostgreSQL 连接器的主机进行复制。要指定允许与数据库进行复制的客户端,请在 PostgreSQL 基于主机的身份验证文件 pg_hba.conf 中添加条目。有关 pg_hba.conf 文件的更多信息,请参阅 PostgreSQL 文档。

操作步骤

  • 在 pg_hba.conf 文件中添加条目,指定可以与数据库主机进行复制的 Debezium 连接器主机。例如:

    pg_hba.conf 文件示例:

    local   replication     <youruser>                          trust
    host    replication     <youruser>  127.0.0.1/32            trust
    host    replication     <youruser>  ::1/128                 trust

以下列表说明了前面 pg_hba.conf 示例中若干行的用途:

local replication <youruser>

指示服务器在本地(即在服务器机器上)为 <youruser> 允许复制。

host replication <youruser> 127.0.0.1/32

指示服务器允许 localhost 上的 <youruser> 使用 IPV4 接收复制变更。

host replication <youruser> ::1/128

指示服务器允许 localhost 上的 <youruser> 使用 IPV6 接收复制变更。

有关网络掩码(netmask)的更多信息,请参阅网络地址类型(PostgreSQL 文档)。

支持的 PostgreSQL 拓扑

PostgreSQL 连接器可以与独立的 PostgreSQL 服务器配合使用,也可以与 PostgreSQL 服务器集群配合使用。

PostgreSQL 15 及更早版本的集群

在运行 PostgreSQL 15 及更早版本的环境中部署 Debezium 时,只能在集群中的主服务器上配置逻辑复制槽。无法在集群中的副本服务器上配置逻辑复制。

因此,Debezium PostgreSQL 连接器只能连接并与主服务器通信。如果主服务器发生故障,连接器就会停止。要从故障中恢复,你必须修复集群,然后将原主服务器提升为 primary,或将另一台 PostgreSQL 服务器提升为 primary。有关更多信息,请参阅故障后从新的主服务器采集数据。

PostgreSQL 16 及更高版本的集群

使用 PostgreSQL 16 及更高版本的集群部署 Debezium 时,可以在副本服务器上设置逻辑复制槽。此功能使 Debezium 能够从主服务器以外的服务器采集变更事件。但请注意,Debezium 连接到副本服务器的延迟通常高于连接到主服务器的延迟。

另外,请注意,PostgreSQL 副本服务器上的复制槽不会与主服务器上对应的复制槽自动同步。为了便于 PostgreSQL 16 集群在故障后恢复,你应该定期手动执行同步,将备用服务器上复制槽的位置推进到与主服务器上的位置一致。

Debezium 与 PostgreSQL 17 及更高版本的集群

当你使用 PostgreSQL 17 或更高版本部署 Debezium 时,可以在主服务器上设置逻辑复制槽,并使这些槽可用于故障转移。PostgreSQL 可以自动将故障转移槽的状态传播到一个或多个副本服务器。在启用了自动复制的环境中,如果发生故障,可用的副本会被自动提升为主服务器。Debezium 可以继续从新的主服务器摄取变更,而无需任何配置更改,从而有助于确保连接器不会遗漏任何事件。

WAL 磁盘空间占用

在某些情况下,WAL 文件占用的 PostgreSQL 磁盘空间可能会飙升或超出正常比例地增长。造成这种情况的原因可能有几种:

  • 连接器已接收数据所对应的 LSN 记录在服务器 pg_replication_slots 视图的 confirmed_flush_lsn 列中。比该 LSN 更早的数据已不再可用,由数据库负责回收这些磁盘空间。

    同样在 pg_replication_slots 视图中,restart_lsn 列包含连接器可能需要的最旧 WAL 的 LSN。如果 confirmed_flush_lsn 的值在定期增长而 restart_lsn 的值滞后,则数据库需要回收空间。

    数据库通常以批处理块(batch block)的方式回收磁盘空间。这是预期行为,用户无需采取任何操作。

  • 被跟踪的数据库中存在大量更新,但只有极少量更新与连接器正在捕获变更的表和 schema 相关。这种情况可以通过定期的 heartbeat 事件轻松解决。设置 heartbeat.interval.ms 连接器配置属性。

    为了让连接器能够检测并处理来自心跳表的事件,你必须将该表添加到 publication.name 属性所指定的 PostgreSQL 发布(publication)中。如果该发布早于你的 Debezium 部署,则连接器按原样使用该发布。如果该发布尚未配置为自动复制数据库中 FOR ALL TABLES 的变更,则你必须显式地将心跳表添加到该发布中,例如:

    ALTER PUBLICATION <publicationName> ADD TABLE <heartbeatTableName>;

  • PostgreSQL 实例包含多个数据库,其中一个数据库流量很高。Debezium 捕获的是另一个与之相比流量较低的数据库的变更。由于复制槽是按数据库工作的,而 Debezium 没有被调用,因此 Debezium 无法确认 LSN。由于 WAL 由所有数据库共享,其使用量往往会持续增长,直到 Debezium 正在捕获变更的那个数据库发出事件为止。为了解决这个问题,有必要:

  • 通过 heartbeat.interval.ms 连接器配置属性启用周期性心跳记录的生成。

  • 定期从 Debezium 正在捕获更改的数据库中发出变更事件。

随后,一个独立的进程会定期更新该表,方式可以是插入新行,也可以是反复更新同一行。PostgreSQL 会据此调用 Debezium,Debezium 确认最新的 LSN,从而允许数据库回收 WAL 空间。此任务可以通过 heartbeat.action.query 连接器配置属性自动完成。

对于使用 AWS RDS 上 PostgreSQL 的用户,一种类似于高流量/低流量场景的情况可能会在空闲环境中出现。AWS RDS 会频繁地(每 5 分钟)使其自身系统表上的写入对客户端不可见。同样,定期发出事件即可解决该问题。

为同一数据库服务器设置多个连接器

当您部署多个 Debezium 连接器来捕获同一个 PostgreSQL 数据库服务器的更改时,每个连接器都需要一个唯一的复制槽和发布。此配置可防止数据丢失,并确保每个连接器独立跟踪其在数据库事务日志中的位置。

Debezium 使用复制槽从数据库流式传输更改。复制槽使用日志序列号(LSN)作为持久化指针,来跟踪连接器已处理的预写日志(WAL)中的最后位置。这种跟踪机制使 PostgreSQL 能够保留 WAL 段,直到 Debezium 处理完它们为止。

由于每个复制槽只维护一个 LSN 指针,它只能跟踪一个消费者在 WAL 中的位置。当多个连接器尝试共享同一个复制槽时,它们会争用同一个位置标记,导致行为不可预测并可能造成数据丢失。

如果允许多个连接器从同一个复制槽捕获数据,就会面临数据丢失的风险,因为一个复制槽只能发送一次每项更改。当多个连接器共用一个槽时,PostgreSQL 只会将每个更改事件发送给其中一个竞争的消费者,其他连接器永远收不到这些事件,从而导致数据捕获不完整。发生这种错误配置时,PostgreSQL 不会给出任何警告,因此数据丢失是静默的,难以察觉。

除了复制槽之外,当 Debezium 连接器使用 pgoutput 插件时,它会使用发布(publication)来流式传输事件。发布存在于数据库级别,用于定义连接器从哪些表中捕获更改。与复制槽一样,每个连接器应各自拥有独立的发布,以确保正确的隔离和独立的配置。

虽然从技术上讲,多个连接器可以共享同一个发布,前提是它们从同一组表中捕获更改,但这种配置仍然要求每个连接器使用各自的复制槽。只有在确实需要多个连接器监控完全相同的表集合时(例如将处理负载分散到多条管道中),共享发布才是合适的。

其他资源

升级 PostgreSQL

当你升级 Debezium 所使用的 PostgreSQL 数据库时,必须采取特定的步骤来防止数据丢失并确保 Debezium 持续运行。总体而言,Debezium 对网络故障及其他中断造成的中断具有较强的恢复能力。例如,当被连接器监控的数据库服务器停止或崩溃后,在连接器重新与 PostgreSQL 服务器建立通信时,它会从日志序列号(LSN)偏移量记录的最后位置继续读取。连接器从 Kafka Connect 偏移量主题中检索最后记录的偏移量信息,并在已配置的 PostgreSQL 复制槽中查询具有相同值的 LSN。

要让连接器从先前记录的偏移量处恢复流式传输,相应的复制槽状态必须仍然可用。然而,如果在 PostgreSQL 升级过程中复制槽被删除,升级完成后原始的复制槽不会被恢复。因此,尽管连接器可以自动创建新的复制槽,但 PostgreSQL 无法使用该新槽返回与最后一个已知偏移量相对应的历史位置。

你可以创建新的复制槽,但要防止数据丢失,仅创建新槽是不够的。新的复制槽只能提供在你创建该槽之后发生的变更的 LSN,无法提供升级之前发生的事件的偏移量。当连接器重启时,它会首先从 Kafka 偏移量主题中请求最后一个已知偏移量,然后向复制槽发送请求,以返回从偏移量主题中检索到的偏移量对应的信息。但新的复制槽无法提供连接器从预期位置恢复流式传输所需的信息。于是连接器会跳过日志中所有已有的变更事件,仅从日志中的最新位置恢复流式传输。这可能导致静默的数据丢失:连接器不会为被跳过的事件发出任何记录,也不会提供任何信息来表明事件已被跳过。

有关如何执行 PostgreSQL 数据库升级、使 Debezium 能够在将数据丢失风险降至最低的同时继续捕获事件的指导,请参阅以下过程。

操作步骤

  1. 暂时停止向数据库写入的应用程序,或将它们置为只读模式。
  2. 备份数据库。
  3. 暂时禁用对数据库的写访问。
  4. 确认在你阻止写操作之前数据库中发生的所有更改都已保存到预写日志(WAL)中,并且 WAL 的 LSN 已反映在复制槽上。
  5. 为连接器留出足够的时间,以便捕获写入复制槽的所有事件记录。此步骤可确保停机之前发生的所有更改事件都已被记录,并且已保存到 Kafka 中。
  6. 通过检查已刷新 LSN 的值,确认连接器已完成从复制槽中消费条目。
  7. 通过停止 Kafka Connect 来正常关闭连接器。Kafka Connect 会停止连接器、将所有事件记录刷新到 Kafka,并记录从每个连接器接收到的最后一个偏移量。
作为停止整个 Kafka Connect 集群的替代方案,你可以通过删除连接器来停止它。请不要移除偏移量主题,因为它可能被其他 Kafka 连接器共享。之后,当你恢复了数据库的写入权限并准备重新启动连接器时,必须重新创建该连接器。
  1. 以 PostgreSQL 管理员身份,在主数据库服务器上删除复制槽。请不要使用 slot.drop.on.stop 属性来删除复制槽,该属性仅用于测试。
  2. 停止数据库。
  3. 使用认可的 PostgreSQL 升级流程执行升级,例如 pg_upgrade,或 pg_dump 与 pg_restore。
  4. (可选)使用标准 Kafka 工具从偏移量存储主题中移除连接器偏移量。有关如何移除连接器偏移量的示例,请参阅 Debezium 社区常见问题中的如何移除连接器已提交的偏移量。
  5. 重新启动数据库。
  6. 以 PostgreSQL 管理员身份,在数据库上创建 Debezium 逻辑复制槽。在正常运行期间,当连接器启动时,如果复制槽尚不存在,它会自动创建该槽。但在升级期间,你必须在启用数据库写入之前手动创建该槽,以便升级之后写入的更改能够被捕获保留。否则,Debezium 将无法捕获这些更改,从而导致数据丢失。
  7. 验证定义 Debezium 要捕获的表的发布(publication)在升级后仍然存在。如果该发布不可用,请以 PostgreSQL 管理员身份连接到数据库以创建新的发布。
  8. 如果在上一步中需要创建新的发布,请更新 Debezium 连接器配置,将新发布的名称添加到 publication.name 属性中。
  9. 在连接器配置中,重命名连接器。
  10. 在连接器配置中,将 slot.name 设置为 Debezium 复制槽的名称。
  11. 验证新的复制槽可用。
  12. 恢复数据库的写入权限,并重新启动所有向数据库写入的应用程序。
  13. 在连接器配置中,将 snapshot.mode 属性设置为 no_data,然后重新启动连接器。
如果你在第 6 步中无法确认 Debezium 已读取完所有数据库变更,可以将 snapshot.mode=initial 设置为 initial,让连接器执行一次新的快照。如有必要,你可以检查在升级之前立即创建的数据库备份的内容,以确认连接器是否从复制槽中读取了所有变更。

其他资源

部署

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

前置条件

操作步骤

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

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

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

连接器配置示例

以下是一个 PostgreSQL 连接器的配置示例,该连接器连接到位于 192.168.99.100、端口为 5432 的 PostgreSQL 服务器,其逻辑名称为 fulfillment。通常,你通过在 JSON 文件中设置连接器可用的配置属性来配置 Debezium PostgreSQL 连接器。

你可以选择只为数据库中部分模式和表生成事件。另外,你还可以忽略、屏蔽或截断包含敏感数据的列、超过指定大小的列,或你不需要的列。

{
  "name": "fulfillment-connector",
  "config": {
    "connector.class": "io.debezium.connector.postgresql.PostgresConnector",
    "database.hostname": "192.168.99.100",
    "database.port": "5432",
    "database.user": "postgres",
    "database.password": "postgres",
    "database.dbname" : "postgres",
    "topic.prefix": "fulfillment",
    "table.include.list": "public.inventory"

  }
}

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

name

连接器在向 Kafka Connect 服务注册时使用的名称。

connector.class

此 PostgreSQL 连接器类的名称。

database.hostname

PostgreSQL 服务器的地址。

database.port

PostgreSQL 服务器的端口号。

database.user

具有所需权限的 PostgreSQL 用户名。

database.password

具有所需权限的 PostgreSQL 用户的密码。

database.dbname

要连接的 PostgreSQL 数据库的名称。

topic.prefix

PostgreSQL 服务器/集群的主题前缀,它构成一个命名空间,用于该连接器写入的所有 Kafka 主题名称、Kafka Connect 架构名称,以及使用 Avro 转换器时相应 Avro 架构的命名空间。

table.include.list

此连接器要监控的、由该服务器承载的所有表的列表。这是可选的,还有其他属性可用于列出要包含或排除在监控之外的模式和表。

有关可以在这些配置中指定的完整 PostgreSQL 连接器属性列表,请参阅 PostgreSQL 连接器属性。

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

  • 连接到 PostgreSQL 数据库。
  • 读取事务日志。
  • 将变更事件记录流式传输到 Kafka 主题。

添加连接器配置

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

前提条件

操作步骤

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

结果

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

连接器属性

Debezium PostgreSQL 连接器有许多配置属性,可用于实现适合你应用的连接器行为。许多属性都有默认值。属性相关信息按以下方式组织:

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

除非提供了默认值,否则以下配置属性是必需的。

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

属性 默认值 说明
name 无默认值 连接器的唯一名称。尝试使用相同名称再次注册将会失败。所有 Kafka Connect 连接器都需要此属性。
connector.class 无默认值 连接器的 Java 类名。PostgreSQL 连接器始终使用 io.debezium.connector.postgresql.PostgresConnector 作为该值。
tasks.max 1 应为此连接器创建的最大任务数。PostgreSQL 连接器始终使用单个任务,因此不会使用此值,默认值始终可接受。
plugin.name decoderbufs 安装在 PostgreSQL 服务器上的 逻辑解码插件 的名称。支持的值为 decoderbufs 和 pgoutput。
slot.name debezium

用于连接器从特定数据库或架构中的特定插件流式传输变更的 PostgreSQL 逻辑解码槽的名称。如果该槽尚不存在,连接器会在启动时自动创建它,前提是 Debezium 数据库用户拥有所需的权限。要自动创建槽,用户必须具有 REPLICATION 权限。使用 pgoutput 插件时,用户还必须具有该数据库上的 CREATE 权限。部署多个连接器实例时,每个实例必须使用唯一的槽名称。服务器使用该槽将事件流式传输到你正在配置的 Debezium 连接器。

槽名称必须符合 PostgreSQL 复制槽命名规则,该规则规定名称可以包含小写字母、数字和下划线字符。

slot.drop.on.stop

false

当连接器以优雅、预期的方式停止时,是否删除逻辑复制槽。默认行为是连接器停止时为其保留复制槽的配置。连接器重新启动时,相同的复制槽可使连接器从中断处继续处理。

仅在测试或开发环境中设置为 true。删除槽可让数据库丢弃 WAL 段。连接器重新启动时会执行新的快照,或者可以从 Kafka Connect 偏移量主题中持久化的偏移量继续处理。

slot.failover

false

指定连接器是否创建故障转移槽。如果省略此设置,或者主服务器运行的是 PostgreSQL 16 或更早版本,连接器将不创建故障转移槽。

PostgreSQL 使用 synchronized_standby_slots 参数来配置主服务器与备用服务器之间的复制槽同步。在主服务器上设置此参数,可指定它与备用服务器同步的物理复制槽。

offset.mismatch.strategy

no_validation

指定连接器启动时,如何处理存储的偏移量 LSN 与复制槽的已确认刷新 LSN 之间的不一致情况。

offset.mismatch.strategy 属性是一项技术预览功能。技术预览功能不受 Red Hat 生产服务水平协议(SLA)支持,可能在功能上并不完整。Red Hat 不建议在生产环境中使用它们。这些功能可让客户提前访问即将推出的产品功能,从而在开发过程中测试功能并提供反馈。有关 Red Hat 技术预览功能支持范围的更多信息,请参阅技术预览功能支持范围。

此属性取代了已废弃的 internal.slot.seek.to.known.offset.on.start 布尔属性。

在以下几种场景中可能会出现不一致:

  • 使用 pg_replication_slot_advance() 手动推进了复制槽。
  • 当使用 lsn.flush.mode=connector_and_driver 时,JDBC 驱动程序的保活机制为未受监控的 WAL 活动刷新了 LSN,从而使复制槽的位置超前于存储的偏移量。
  • 复制槽在被删除后重新创建。
  • Kafka Connect 偏移量存储已损坏或被重置。

请设置以下选项之一:

no_validation

连接器尝试从存储的偏移量开始流式传输,而不校验复制槽状态。如果复制槽超前于偏移量,PostgreSQL 会返回错误,因为请求的 LSN 已不在 WAL 中(Postgres 14 及更早版本),或者会从 confirmed_flush_lsn 开始(Postgres 15 及更高版本)。

trust_offset

连接器会校验存储的偏移量是否不落后于复制槽的已确认刷新 LSN。如果偏移量落后于复制槽,连接器将以指示可能发生数据丢失的错误失败退出。如果偏移量超前于或等于复制槽,连接器会在可能的情况下将复制槽推进到偏移量所在位置。当您希望检测并收到可能表明数据丢失的意外复制槽状态变化的警报时,请使用此策略。

trust_slot

连接器将 PostgreSQL 复制槽视为权威的事实来源。如果存储的偏移量落后于该槽的已确认刷新 LSN,连接器会自动推进偏移量以匹配槽的位置。

此策略会跳过存储偏移量与槽位置之间的事件重放。当你确信复制槽可靠地代表了需要捕获的所有逻辑复制事件时,请使用此策略。当使用 lsn.flush.mode=connector_and_driver(该模式要求信任槽位置)时,这是合适的。请务必确认你的环境能够在主从切换期间保证复制槽的持久性。

trust_greater_lsn

连接器会同步到存储偏移量与槽的已确认刷新 LSN 之间较大的那个 LSN。如果偏移量落后于槽,连接器会将偏移量推进到槽的位置;如果偏移量超前于槽,连接器会在可能的情况下将槽推进到偏移量的位置。这提供了双向同步,在使用 lsn.flush.mode=connector_and_driver 时非常有用。与 trust_slot 相同的槽可靠性要求同样适用。

当槽领先于偏移量时,此策略会跳过存储偏移量与槽位置之间的事件重放。

为保持向后兼容,已弃用的 internal.slot.seek.to.known.offset.on.start 属性会被自动映射:

  • internal.slot.seek.to.known.offset.on.start=false 映射为 offset.mismatch.strategy=no_validation
  • internal.slot.seek.to.known.offset.on.start=true 映射为 offset.mismatch.strategy=trust_offset

publication.name

dbz_publication

使用 pgoutput 进行变更流式传输时所创建的 PostgreSQL 发布(publication)的名称。

如果该发布(publication)尚不存在,则会在启动时创建,并且包含所有表。随后,若配置了相应的包含/排除列表过滤,Debezium 会应用该过滤,将发布限制为仅包含所关注特定表的变更事件。连接器用户必须拥有超级用户权限才能创建此发布,因此通常更推荐在首次启动连接器之前先创建好该发布。

如果该发布已经存在——无论是包含所有表,还是配置为仅包含部分表——Debezium 都将按其现有定义使用该发布。

database.hostname

无默认值

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

database.port

5432

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

database.user

无默认值

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

database.password

无默认值

连接 PostgreSQL 数据库服务器时使用的密码。

database.dbname

无默认值

要从中流式传输变更的 PostgreSQL 数据库的名称。

database.connection.factory.class

无默认值

ConnectionFactory 的名称。该属性的一种用途是:当您使用 AWS IAM 身份验证建立到 Amazon RDS 上 PostgreSQL 数据库的 JDBC 连接时,可用它来指定封装凭据和其他连接详情的接口。有关更多信息,参见 Amazon RDS 上的 PostgreSQL)。

topic.prefix

无默认值

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

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

schema.include.list

无默认值

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

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

schema.exclude.list

无默认值

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

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

table.include.list

无默认值

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

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

table.exclude.list

无默认值

一个可选的、以逗号分隔的正则表达式列表,用于匹配你不想捕获其变更的表的完全限定表标识符。每个标识符的形式为 schemaName.tableName。设置此属性后,连接器会捕获除你指定的表之外的所有表的变更。

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

column.include.list

无默认值

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

为了匹配列名,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 属性所包含的列没有变化,这实质上会过滤掉这些消息。

该属性仅在表的 REPLICA IDENTITY 设置为 FULL 时生效。

time.precision.mode

adaptive

时间、日期和时间戳可以使用不同种类的精度来表示:

adaptive 根据数据库列的类型,使用毫秒、微秒或纳秒精度值,精确捕获与数据库中一致的时间和时间戳值。

adaptive_time_microseconds 根据数据库列的类型,使用毫秒、微秒或纳秒精度值,精确捕获与数据库中一致的日期、日期时间和时间戳值。唯一的例外是 TIME 类型字段,它始终以微秒进行捕获。

connect 始终使用 Kafka Connect 内置的 Time、Date 和 Timestamp 表示形式来表示时间和时间戳值,无论数据库列的精度如何,这些表示形式都使用毫秒精度。有关更多信息,请参阅时间值。

decimal.handling.mode

precise

指定连接器应如何处理 DECIMAL 和 NUMERIC 列的值:

precise 使用 java.math.BigDecimal 以二进制形式在变更事件中表示值。

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

string 将值编码为格式化字符串,这些字符串易于使用,但会丢失有关实际类型的语义信息。有关更多信息,请参阅十进制类型。

hstore.handling.mode

json

指定连接器应如何处理 hstore 列的值:

map 使用 MAP 表示值。

json 使用 json string 表示值。此设置会将值编码为格式化字符串,例如 {"key" : "val"}。有关更多信息,请参阅PostgreSQL HSTORE 类型。

interval.handling.mode

numeric

指定连接器应如何处理 interval 列的值:

numeric 使用近似的微秒数表示时间间隔。

string 使用字符串模式表示法 P<years>Y<months>M<days>DT<hours>H<minutes>M<seconds>S 来精确表示时间间隔。例如:P1Y2M3DT4H5M6.78S。有关更多信息,请参阅 PostgreSQL 基本类型。

database.sslmode

prefer

是否使用与 PostgreSQL 服务器的加密连接。可选值包括:

disable 使用非加密连接。

allow 先尝试使用非加密连接,若失败,则使用安全(加密)连接。

prefer 先尝试使用安全(加密)连接,若失败,则使用非加密连接。

require 使用安全(加密)连接,若无法建立则失败。

verify-ca 行为与 require 类似,但还会根据已配置的证书颁发机构(CA)证书验证服务器的 TLS 证书;若未找到有效的匹配 CA 证书,则失败。

verify-full 行为与 verify-ca 类似,但还会验证服务器证书是否与连接器尝试连接的主机相匹配。有关更多信息,请参阅 PostgreSQL 文档。

database.sslcert

无默认值

包含客户端 SSL 证书的文件路径。有关更多信息,请参阅 PostgreSQL 文档。

database.sslkey

无默认值

包含客户端 SSL 私钥的文件路径。有关更多信息,请参阅 PostgreSQL 文档。

database.sslpassword

无默认值

用于从 database.sslkey 指定的文件中访问客户端私钥的密码。有关更多信息,请参阅 PostgreSQL 文档。

database.sslrootcert

无默认值

包含用于验证服务器的根证书的文件路径。有关更多信息,请参阅 PostgreSQL 文档。

database.sslfactory

无默认值

创建 SSL Socket 的类名。在开发环境中,可使用 org.postgresql.ssl.NonValidatingFactory 来禁用 SSL 验证。

database.tcpKeepAlive

true

启用 TCP keep-alive 探测,以验证数据库连接是否仍然存活。有关更多信息,请参阅 PostgreSQL 文档。

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

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

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

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

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

n/a

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

化名由应用所指定的 hashAlgorithm 和 salt 后得到的哈希值构成。根据所使用的哈希函数,可以在列值被替换为化名的同时保持引用完整性。受支持的哈希函数详见 Java 密码学架构标准算法名称文档中的 MessageDigest 章节。

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

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

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

根据所使用的 hashAlgorithm(哈希算法)、所选的 salt(盐值)以及实际数据集,得到的结果数据集可能无法被完全脱敏。

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

column.propagate.source.type

n/a

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

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

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

列的完全限定名遵循以下格式之一:databaseName.tableName.columnName,或 databaseName.schemaName.tableName.columnName。为匹配列名,Debezium 会将您指定的正则表达式作为锚定正则表达式使用,即该表达式会与列的整个名称字符串进行匹配,而不会匹配列名中可能出现的子字符串。

datatype.propagate.source.type

n/a

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

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

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

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

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

message.key.columns

空字符串

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

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

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

<完全限定的表名>:<keyColumn>,<keyColumn>

若要基于多个列名构建表键,请在各列名之间使用逗号分隔。

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

<schemaName>.<tableName>

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

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

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

在该示例中,列 pk1 和 pk2 被指定为 inventory.customer 表的消息键。对于任意 schema 中的 purchaseorders 表,列 pk3 和 pk4 作为消息键。

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

如果为此属性指定的表达式匹配的列不属于该表的主键,请将该表的 REPLICA IDENTITY 设置为 FULL。如果将 REPLICA IDENTITY 设置为其他值(例如 DEFAULT),那么在删除操作之后,连接器将无法生成带有预期 null 值的墓碑事件。

publication.autocreate.mode

all_tables

指定连接器是否以及如何创建 publication。此设置仅在连接器使用 pgoutput 插件流式传输变更时适用。

要创建 publication,连接器必须通过具有特定权限的数据库账户访问 PostgreSQL。有关更多信息,请参阅设置权限以允许 Debezium 创建 PostgreSQL publication。

可指定以下值之一:

all_tables

如果已存在 publication,连接器将使用它。如果不存在 publication,连接器会为连接器捕获变更的数据库中的所有表创建一个 publication。连接器运行以下 SQL 命令来创建 publication:

CREATE PUBLICATION <publication_name> FOR ALL TABLES;

如果配置了 skipped.operations 属性,连接器会追加 WITH (publish = '…​') 子句,以将 publication 限制为未被跳过的操作类型:

CREATE PUBLICATION <publication_name> FOR ALL TABLES WITH (publish = '<allowed_operations>');

disabled

连接器不会尝试创建 publication。在运行连接器之前,数据库管理员或配置为执行复制的用户必须已经创建该 publication。如果连接器找不到该 publication,将抛出异常并停止运行。

filtered

如果不存在 publication,连接器会运行以下格式的 SQL 命令来创建一个:

CREATE PUBLICATION <publication_name> FOR TABLE <tbl1, tbl2, tbl3>

生成的 publication 包含与当前过滤配置匹配的表,如连接器配置属性 schema.include.list、schema.exclude.list、table.include.list 和 table.exclude.list 所指定。

如果已存在 publication,连接器会运行以下格式的 SQL 命令,为与当前过滤配置匹配的表更新该 publication:

ALTER PUBLICATION <publication_name> SET TABLE <tbl1, tbl2, tbl3>

在这两种情况下,如果配置了 skipped.operations 属性,连接器会追加 WITH (publish = '…​') 子句,将发布范围限制为未被跳过的操作类型。

no_tables

如果已存在发布(publication),连接器会直接使用它。如果不存在发布,连接器会执行以下格式的 SQL 命令来创建一个不指定任何表的发布:

CREATE PUBLICATION <publication_name>;

如果你希望连接器只捕获逻辑解码消息,而不捕获其他变更事件(例如由任何表上的 INSERT、UPDATE 和 DELETE 操作引起的事件),请设置 no_tables 选项。

如果选择此选项,为防止连接器发出并处理 READ 事件,你可以指定不希望捕获其变更的架构或表的名称,例如使用 "table.exclude.list": "public.*" 或 "schema.exclude.list": "public"。

replica.identity.autoset.values

空字符串

设置此属性可根据表名,将特定的副本标识设置应用于连接器捕获的一部分表。该属性设置的副本标识值会覆盖数据库中已设置的副本标识值。

该属性接受一个以逗号分隔的键值对列表。每个键是与全限定表名匹配的正则表达式;相应的值指定一种副本标识类型。例如:

<fqTableNameA>:<replicaIdentity1>,<fqTableNameB>:<replicaIdentity2>,<fqTableNameC>:<replicaIdentity3>

使用以下格式指定全限定表名:SchemaName.TableName

将副本标识设置为以下值之一:

DEFAULT

记录变更事件发生前为主键列设置的值(如果存在)。这是非系统表的默认设置。

INDEX indexName

记录变更事件发生前为指定索引定义的所有列所设置的值。该索引必须是唯一的、非部分索引、不可延迟的,并且仅包含标记为 NOT NULL 的列。如果指定的索引被删除,则最终行为与将该值设置为 NOTHING 相同。

FULL

记录变更事件发生前行中所有列所设置的值。

NOTHING

不记录变更事件发生前行状态的任何信息。这是系统表的默认值。

示例:

schema1.*:FULL,schema2.table2:NOTHING,schema2.table3:INDEX idx_name

replica.identity.autoset.values 属性仅适用于连接器捕获的表。其他表即使匹配指定的表达式也会被忽略。请使用以下连接器属性来指定要捕获的表:

binary.handling.mode

bytes

指定二进制(bytea)列在更改事件中应如何表示。可指定以下值之一:

bytes

将二进制数据表示为字节数组。

base64

将二进制数据表示为 base64 编码的字符串。

base64-url-safe

将二进制数据表示为 base64-url-safe 编码的字符串。

hex

将二进制数据表示为十六进制(base16)编码的字符串。

schema.name.adjustment.mode

none

指定为兼容连接器所使用的消息转换器,架构名称应如何进行调整。可设置以下值之一:

none

不进行任何调整。

avro

将不能用于 Avro 类型名称的字符替换为下划线。

avro_unicode

将下划线或不能用于 Avro 类型名称的字符替换为相应的 Unicode 字符,例如 _uxxxx。

在上面的示例中,下划线字符(_)表示转义序列,等同于 Java 中的反斜杠。

field.name.adjustment.mode

none

指定为兼容连接器所使用的消息转换器,字段名称应如何进行调整。可指定以下值之一:

none

不进行任何调整。

avro

将不能用于 Avro 类型名称的字符替换为下划线。

avro_unicode

将下划线或不能用于 Avro 类型名称的字符替换为相应的 Unicode 字符,例如 _uxxxx。

在上面的示例中,下划线字符(_)表示转义序列,等同于 Java 中的反斜杠。

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

money.fraction.digits

2

指定在将 Postgres 的 money 类型转换为 java.math.BigDecimal(用于表示变更事件中的值)时应使用多少位小数。仅当 decimal.handling.mode 设置为 precise 时适用。

message.prefix.include.list

无默认值

一个可选的、以逗号分隔的正则表达式列表,用于匹配你希望连接器捕获的逻辑解码消息前缀的名称。默认情况下,连接器会捕获所有逻辑解码消息。设置此属性后,连接器将仅捕获具有该属性所指定前缀的逻辑解码消息,其他所有逻辑解码消息都会被排除。

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

如果在配置中包含了此属性,请不要再设置 message.prefix.exclude.list 属性。

有关 message 事件的结构及其排序语义的信息,请参阅 message 事件。

message.prefix.exclude.list

无默认值

一个可选的、以逗号分隔的正则表达式列表,用于匹配你不希望连接器捕获的逻辑解码消息前缀的名称。设置此属性后,连接器不会捕获使用指定前缀的逻辑解码消息,其他所有消息都会被捕获。要排除所有逻辑解码消息,请将此属性的值设置为 .*。

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

如果在配置中包含了此属性,请不要再设置 message.prefix.include.list 属性。

有关 message 事件的结构及其排序语义的信息,请参阅 message 事件。

Debezium PostgreSQL 连接器高级配置属性

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

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

属性 默认值 描述

converters

无默认值

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

isbn

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

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

<converterSymbolicName>.type

例如,

isbn.type: io.debezium.test.IsbnConverter

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

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

snapshot.isolation.mode

serializable

指定连接器在初始快照或临时阻塞快照期间读取数据时所使用的事务隔离级别,以及(如有)所应用的锁定类型。

每种隔离级别都在以下两方面取得不同的平衡:一方面优化并发性和性能,另一方面最大化数据的一致性与准确性。使用更严格隔离级别的快照会产生质量更高、更一致的数据,但其代价是由于锁定时间更长、可并发的事务更少而导致性能下降。限制较少的隔离级别可以提升效率,但代价是数据可能不一致。有关 PostgreSQL 中事务隔离级别的更多信息,请参阅 PostgreSQL 文档。

请指定以下隔离级别之一:

serializable

默认级别,也是限制最严格的隔离级别。该选项可防止序列化异常,并提供最高程度的数据完整性。

为确保被捕获表的数据一致性,快照在使用可重复读(repeatable read)隔离级别的事务中运行,阻止对这些表的并发 DDL 变更,并对数据库创建索引的操作加锁。设置此选项后,在快照结束之前,用户或管理员无法执行某些操作(例如创建表索引)。整个表键范围在快照完成之前始终保持锁定状态。此选项与该属性引入之前连接器所具有的快照行为一致。

repeatable_read

阻止其他事务在快照期间更新表行。快照捕获的新记录可能出现两次:首先作为初始快照的一部分出现,随后又在流式处理阶段出现。不过,对于数据库镜像而言,这种程度的一致性是可以接受的。它确保被扫描表之间的一致性,阻止对所选表的 DDL 操作,以及阻止在整个数据库范围内并发创建索引。允许出现序列化异常。

read_committed

在 PostgreSQL 中,读未提交(Read Uncommitted)与读已提交(Read Committed)隔离模式的行为没有区别。因此,对于此属性,read_committed 选项实际上提供了限制最少的隔离级别。设置该选项会使初始快照和临时阻塞快照牺牲部分一致性,但在快照期间为其他用户带来更好的数据库性能。

总体而言,这种事务一致性级别适用于数据镜像。快照期间其他事务无法更新表行。不过,当某条记录在初始快照期间被添加,而连接器在流式处理阶段开始后再次捕获该记录时,可能会出现轻微的数据不一致。

read_uncommitted

名义上,这个选项提供限制最宽松的隔离级别。但是,正如 read-committed 选项的说明中所述,对于 Debezium PostgreSQL 连接器而言,该选项提供的隔离级别与 read_committed 选项相同。

snapshot.mode

initial

指定连接器启动时执行快照的条件:

always

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

initial

仅当逻辑服务器名称没有记录任何偏移量时,连接器才执行快照。

initial_only

连接器执行一次初始快照后即停止,不处理任何后续更改。

no_data

连接器从不执行快照。当连接器按此方式配置时,启动后的行为如下:

如果 Kafka 偏移量主题中存在先前存储的 LSN,连接器将从该位置继续流式传输更改。如果没有存储 LSN,连接器将从服务器上创建 PostgreSQL 逻辑复制槽的时间点开始流式传输更改。只有在您确定所有关注的数据仍反映在 WAL 中时,才使用此快照模式。

when_needed

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

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

configuration_based

通过此选项,您可以使用以 'snapshot.mode.configuration.based' 为前缀的一组连接器属性来控制快照行为。

custom

连接器按照 snapshot.mode.custom.name 属性指定的实现执行快照,该属性定义了 io.debezium.spi.snapshot.Snapshotter 接口的自定义实现。

有关更多信息,请参阅 snapshot.mode 选项表。

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,请设置此属性以指定当 Schema 历史主题不可用时,连接器是否在快照中包含表结构。

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() 方法所提供的自定义实现的名称。连接器重启后会调用该实现,以确定是否执行快照。有关更多信息,请参阅自定义快照器 SPI。

snapshot.locking.mode

none

指定连接器在执行结构快照时如何对表加锁。请设置以下选项之一:

shared

在快照的初始阶段(即读取数据库结构和其他元数据的阶段),连接器持有表锁以阻止其他会话独占访问该表。初始阶段结束后,快照不再需要表锁。

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 中指定的所有表

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

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

snapshot.lock.timeout.ms

10000

正整数值,用于指定执行快照时等待获取表锁的最长时间(以毫秒为单位)。如果连接器无法在此时间间隔内获取表锁,快照将失败。有关详情,请参阅连接器如何执行快照。

snapshot.select.statement.overrides

无默认值

指定连接器对哪些表使用自定义 SELECT 语句,以确定哪些行包含在快照中。

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

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

主属性

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

例如,

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

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

次属性

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

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

snapshot.select.statement.overrides.customers.orders

如果主属性中列出了某个表,但缺少其对应的次级属性,连接器会记录一条警告,并对该表使用默认的快照行为。

配置示例

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

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

event.processing.failure.handling.mode

fail

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

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

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

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

max.batch.size

2048

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

max.queue.size

8192

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

max.queue.size.in.bytes

0

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

poll.interval.ms

500

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

include.unknown.datatypes

false

指定连接器在遇到数据类型未知的字段时的行为。默认情况下,连接器会从变更事件中省略该字段并记录一条警告。

如果希望变更事件包含该字段的不透明二进制表示,请将此属性设置为 true。这样消费者即可自行解码该字段。你可以通过设置 binary handling mode 属性来控制具体的表示形式。

当 include.unknown.datatypes 设置为 true 时,消费者可能面临向后兼容性问题。特定于数据库的二进制表示不仅可能在不同版本之间发生变化,而且如果 Debezium 最终支持了该数据类型,下游收到的将是逻辑类型形式的数据类型,这就需要消费者进行相应的调整。通常情况下,遇到不受支持的数据类型时,应提交功能请求,以便后续添加支持。

database.initial.statements

无默认值

由分号分隔的 SQL 语句列表,连接器在建立到数据库的 JDBC 连接时执行这些语句。如果需要将分号用作字符而非分隔符,请使用两个连续的分号 ;;。

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

连接器在创建用于读取事务日志的连接时,不会执行这些语句。

status.update.interval.ms

10000

向服务器发送复制连接状态更新的频率,以毫秒为单位。

该属性还控制检查数据库状态的频率,以便在数据库关闭的情况下检测到连接失效。

heartbeat.interval.ms

0

控制连接器向 Kafka 主题发送心跳消息的频率。默认情况下,连接器不发送心跳消息。

心跳消息用于监控连接器是否正在接收来自数据库的变更事件。心跳消息有助于减少连接器重启时需要重新发送的变更事件数量。要发送心跳消息,请将此属性设置为一个正整数,该整数表示相邻两条心跳消息之间的毫秒数。

当被跟踪的数据库中存在大量更新,但其中只有极少数更新与连接器正在捕获变更的表和模式相关时,就需要心跳消息。在这种情况下,连接器会照常从数据库事务日志中读取数据,但很少向 Kafka 发送变更记录。这意味着没有偏移量更新被提交到 Kafka,连接器也没有机会将最新检索到的 LSN 发送给数据库。数据库会保留包含已被连接器处理过的事件的 WAL 文件。发送心跳消息使连接器能够将最新检索到的 LSN 发送给数据库,从而让数据库回收不再需要的 WAL 文件所占用的磁盘空间。

heartbeat.action.query

无默认值

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

这对于解决 WAL 磁盘空间消耗 中描述的情况非常有用:当从与高流量数据库位于同一主机上的低流量数据库捕获变更时,会阻止 Debezium 处理 WAL 记录,从而无法与数据库确认 WAL 位置。为了解决此情况,请在低流量数据库中创建一个心跳表,并将此属性设置为向该表插入记录的语句,例如:

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

这使连接器能够接收来自低流量数据库的变更并确认其 LSN,从而防止数据库主机上的 WAL 无限增长。

schema.refresh.mode

columns_diff

指定触发表的内存中模式刷新的条件。

columns_diff 是最安全的模式。它确保内存中模式始终与数据库表的模式保持同步。

columns_diff_exclude_unchanged_toast 指示连接器,如果内存中模式与传入消息派生的模式存在差异,则刷新内存中的模式缓存,除非未更改的 TOAST 数据完全解释了该差异。

此设置可以显著提升连接器性能,前提是存在频繁更新的表,且其 TOAST 数据很少参与更新。不过,如果可 TOAST 列从表中被删除,内存中的模式可能会变得过时。

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

10240

在快照期间,连接器按批次读取表内容。此属性指定每批的最大行数。

slot.stream.params

无默认值

以分号分隔的可选键值对列表,用于表示连接器在启动复制流时传递给已配置的 PostgreSQL 逻辑解码插件的参数。连接器会将指定的参数原样传递给插件。有效的参数取决于所配置的逻辑解码插件。与所配置插件无关的参数将被忽略。

decoderbufs 插件

如果连接器使用 decoderbufs 插件,你可以设置 debug-mode 参数,使复制槽记录逻辑解码的调试消息。当 Debezium 发送创建复制槽的 SQL 命令时,会传递 debug-mode 参数。例如:

debug-mode=true
要启用调试消息的查看,你必须将 PostgreSQL 日志级别(log_min_messsages)设置为 DEBUG1 或更高。

pgoutput 插件

对于使用 pgoutput 插件的连接器,请在 slot.stream.params 属性中添加 origin 参数,以指定订阅是请求发布者发送所有更改,还是仅发送本地更改。

例如,若要请求发布者发送所有更改,即既包括发布者本地产生的更改,也包括从其他源复制到发布者的更改,请设置以下值:

origin=any

slot.max.retries

6

如果连接复制槽失败,这是尝试连接的最大连续次数。

slot.retry.delay.ms

10000(10 秒)

当连接器无法连接到复制槽时,两次重试之间等待的毫秒数。

unavailable.value.placeholder

__debezium_unavailable_value

指定连接器提供的常量,用于表示原始值是数据库未提供的 toasted 值。如果 unavailable.value.placeholder 的设置以 hex: 前缀开头,则字符串的其余部分应表示十六进制编码的八位字节。有关更多信息,请参阅 toasted 值。

provide.transaction.metadata

false

决定连接器是否生成带有事务边界的事件,并用事务元数据丰富变更事件信封。如果希望连接器执行此操作,请将其指定为 true。有关更多信息,请参阅事务元数据。

publish.via.partition.root

false

指定连接器如何捕获和发布其从分区表中捕获的变更事件。此设置仅在 publication.autocreate.mode 属性设置为 all_tables 或 filtered,且 Debezium 为被捕获的表创建发布(publication)时才适用。

可设置以下选项之一:

true

连接器将所有分区的变更事件发布到以基表命名的主题中。

当连接器创建发布时,它会提交一条 CREATE PUBLICATION 语句,其中 publish_via_partition_root 参数设置为 true。因此,该发布会忽略变更所在的分区,仅记录源表的名称。

false

连接器将每个源分区的变更发布到以该分区名称命名的主题中。

当连接器创建发布时,CREATE PUBLICATION 语句会省略 publish_via_partition_root 参数,从而让发布始终使用源分区的名称来发布变更事件。

flush.lsn.source

true

已弃用。 请改用 lsn.flush.mode。

为保持向后兼容,此属性会自动映射到 lsn.flush.mode:

  • flush.lsn.source=true 映射为 lsn.flush.mode=connector
  • flush.lsn.source=false 映射为 lsn.flush.mode=manual

lsn.flush.mode

connector

指定连接器如何管理将 LSN(日志序列号)刷新到 PostgreSQL 复制槽。此属性取代了已弃用的 flush.lsn.source 属性。

可设置以下选项之一:

manual

LSN 刷新由你的应用程序或另一种机制在外部管理。连接器不会刷新 LSN。

如果你未配置外部机制来刷新 LSN,WAL 文件会在数据库服务器上不断累积,可能导致存储耗尽和性能下降。

connector

Debezium 在处理每个逻辑复制变更事件后刷新 LSN。PostgreSQL JDBC 驱动程序的保活线程不会刷新 LSN。这是默认模式。

connector_and_driver

Debezium 和 PostgreSQL JDBC 驱动程序的保活线程都可以刷新 LSN。当连接器没有待刷新的 LSN 时,JDBC 驱动程序的保活机制可以刷新服务器报告的保活 LSN。该 LSN 反映了所有 WAL 活动,包括不会产生逻辑复制变更事件的未受监控活动,例如 CHECKPOINT、VACUUM 或 pg_switch_wal()。在低活动量的数据库中,被监控的表变更不频繁时,可使用此模式来防止 WAL 增长。

使用 connector_and_driver 模式时,JDBC 驱动程序的保活机制会刷新服务器报告的保活 LSN,当发生未受监控的 WAL 活动时,这可能会使复制槽推进到超过已存储的偏移量位置。如果你使用持久化偏移量存储,请将 offset.mismatch.strategy 配置为 trust_slot 或 trust_greater_lsn,以启用自动恢复。
由于 connector_and_driver 模式需要定期将 LSN 传播到 Kafka,当 heartbeat.interval.ms = 0 且 provide.transaction.metadata = false 时,连接器会在内部将 heartbeat.interval.ms 设置为 600000(10 分钟)。
connector_and_driver 模式是一项技术预览功能。技术预览功能不受 Red Hat 生产服务等级协议(SLA)支持,功能可能并不完整。Red Hat 不建议在生产环境中使用这些功能。这些功能可让客户提前访问即将推出的产品功能,从而在开发过程中测试功能并提供反馈。有关 Red Hat 技术预览功能支持范围的更多信息,请参阅技术预览功能支持范围。如果您使用的是临时偏移存储(例如 org.apache.kafka.connect.storage.MemoryOffsetBackingStore),则无需额外配置,因为连接器在启动时始终以复制槽位置为准。

lsn.flush.timeout.action

fail

指定 LSN 刷新操作超时时应采取的操作。

设置以下选项之一:

fail

使连接器失败。

warn

记录警告日志并继续处理。

ignore

继续处理并忽略超时。

lsn.flush.timeout.ms

30000(30 秒)

连接器等待 LSN 刷新操作完成的最长时间,以毫秒为单位。如果操作未在指定的时间间隔内完成,连接器将执行 lsn.flush.timeout.action 属性中配置的操作。

retriable.restart.connector.wait.ms

10000(10 秒)

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

skipped.operations

t

以逗号分隔的operation类型列表,指定连接器在流式传输过程中要跳过的操作类型。你可以将连接器配置为跳过以下类型的操作:

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

如果不想让连接器跳过任何操作,请将该值设置为 none。

当你使用 pgoutput 插件且 publication.autocreate.mode 属性设置为 all_tables 或 filtered 时,连接器会通过相应设置 publish 选项,将配置的跳过操作传播到 PostgreSQL 发布(publication)中。例如,如果 skipped.operations=d,连接器会创建或更新发布,设置为 WITH (publish = 'insert,update,truncate'),从而使 PostgreSQL 完全从复制流中省略被跳过的操作,减少不必要的 WAL 流量。

signal.data.collection

无默认值

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

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

<schemaName>.<tableName>

signal.enabled.channels

source

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

  • source
  • kafka
  • file
  • jmx

此外,你还可以实现自定义信号通道。

)notification.enabled.channels

无默认值

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

  • sink
  • log
  • jmx

此外,你还可以实现自定义通知通道。

)incremental.snapshot.chunk.size

1024

连接器在增量快照的一个数据块期间获取并读入内存的最大行数。增大数据块大小可以提高效率,因为快照执行的查询次数更少,而每次查询的数据量更大。但是,较大的数据块大小也需要更多内存来缓冲快照数据。请将数据块大小调整为在你的环境中能获得最佳性能的值。

)incremental.snapshot.watermarking.strategy

insert_insert

指定连接器在增量快照期间使用的水位标记机制,用于对事件去重,这些事件可能已被增量快照捕获,并在流式传输恢复后被再次捕获。你可以指定以下选项之一:

insert_insert

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

insert_delete

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

)read.only

false

指定连接器是否将水位标记写入信号数据集合,以跟踪增量快照的进度。将该值设置为 true,可使对数据库具有只读连接的连接器使用无需写入信号数据集合的增量快照水位标记策略。

)xmin.fetch.interval.ms

0

以毫秒为单位,表示多久从复制槽中读取一次 XMIN。XMIN 值提供了新的复制槽可能的起始位置的下界。默认值 0 表示禁用 XMIN 跟踪。

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.<fully_qualified_table_name> 设置为所需的值。

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

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

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

database.query.timeout.ms

600000(10 分钟)

指定连接器等待查询完成的时间,以毫秒为单位。将该值设置为 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 连接器的名称。

PostgreSQL 连接器的透传配置属性

该连接器支持透传属性,允许 Debezium 为 Apache Kafka 生产者和消费者指定自定义配置选项,以对其行为进行微调。有关 Kafka 生产者和消费者全部配置属性的信息,请参阅 Kafka 文档。

用于配置 PostgreSQL 连接器与 Kafka 信号主题交互方式的透传属性

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

下表描述了 Kafka signal 属性。

表 23. 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 信号消费者之前,会去除这些属性的前缀。

用于配置 PostgreSQL 连接器接收器通知通道的透传属性

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

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

表 24. 接收器通知配置属性

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

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

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

监控

Debezium PostgreSQL 连接器提供了两类指标,这是在 Kafka 和 Kafka Connect 内置的 JMX 指标支持之外额外提供的。

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

Debezium 监控文档详细介绍了如何通过 JMX 暴露这些指标。

自定义 MBean 名称

Debezium 连接器通过连接器的 MBean 名称暴露指标。这些指标针对每个连接器实例,提供有关该连接器快照、流式传输和模式历史进程行为的数据。

默认情况下,当你部署一个配置正确的连接器时,Debezium 会为不同的连接器指标分别生成唯一的 MBean 名称。要查看某个连接器进程的指标,你需要配置可观测性栈来监控其对应的 MBean。但这些默认的 MBean 名称依赖于连接器配置;配置的更改可能会导致 MBean 名称发生变化。MBean 名称一旦改变,连接器实例与 MBean 之间的关联就会中断,从而破坏监控活动。在这种情况下,如果要恢复监控,就必须重新配置可观测性栈以使用新的 MBean 名称。

为防止因 MBean 名称变更而导致的监控中断,你可以配置自定义指标标签。通过在连接器配置中添加 custom.metric.tags 属性即可配置自定义指标。该属性接受键值对,其中每个键表示 MBean 对象名称的一个标签,相应的值表示该标签的值。例如:k1=v1,k2=v2。Debezium 会将指定的标签追加到连接器的 MBean 名称之后。

为连接器配置 custom.metric.tags 属性后,你可以配置可观测性栈来获取与指定标签关联的指标。此后,可观测性栈将使用指定的标签(而非易变的 MBean 名称)来唯一标识连接器。之后,即使 Debezium 重新定义了 MBean 名称的构建方式,或者连接器配置中的 topic.prefix 发生变化,指标收集也不会中断,因为指标抓取任务会使用指定的标签模式来识别连接器。

使用自定义标签的另一个好处是,你可以使用能够反映数据管道架构的标签,从而按照符合你运维需求的方式组织指标。例如,你可以指定标签值,用来声明连接器活动的类型、应用上下文或数据源,比如 db1-streaming-for-application-abc。如果指定了多个键值对,所有指定的键值对都会追加到连接器的 MBean 名称中。

下面的示例说明了标签如何修改默认的 MBean 名称。

示例 3. 自定义标签如何修改连接器 MBean 名称

默认情况下,PostgreSQL 连接器用于流式指标的 MBean 名称如下:

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

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

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

快照指标

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

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

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

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

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

流处理指标

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

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

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

| [](/book/reader/de

出现问题时的行为

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

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

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

配置与启动错误

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

  • 连接器的配置无效。
  • 连接器无法使用指定的连接参数成功连接到 PostgreSQL。
  • 连接器正从之前记录的 PostgreSQL WAL 位置(使用 LSN)重新启动,但 PostgreSQL 已不再保留该历史记录。

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

PostgreSQL 变得不可用

连接器运行时,其所连接的 PostgreSQL 服务器可能因各种原因而变得不可用。如果发生这种情况,连接器会因错误而失败并停止。当服务器再次可用时,请重新启动连接器。

PostgreSQL 连接器以外部存储的形式记录最后处理的偏移量,其格式为 PostgreSQL LSN。连接器重新启动并连接到服务器实例后,会与服务器通信,以便从该特定偏移量继续进行流式传输。只要 Debezium 复制槽保持完好,该偏移量就始终可用。切勿删除主服务器上的复制槽,否则将丢失数据。有关复制槽已被删除的失败情形,请参阅下一节。

集群故障

Debezium PostgreSQL 连接器在集群故障后的行为取决于 PostgreSQL 的版本。PostgreSQL 17 及更高版本支持故障转移复制槽以实现自动恢复;更早的版本则需要人工干预才能恢复事件捕获。

PostgreSQL 15 或更早版本

在 PostgreSQL 15 或更早版本的集群中,只能在主服务器上创建逻辑复制槽。因此,在 PostgreSQL 15 环境中,Debezium PostgreSQL 连接器只能从集群中的活动主服务器捕获事件。在 PostgreSQL 15 集群中,主节点上的复制槽不会传播到副本服务器。如果主服务器宕机,必须将某个备用节点提升为主节点。

PostgreSQL 16 或更高版本

当你将 Debezium 与 PostgreSQL 16 或更高版本配合使用时,可以在副本上创建逻辑复制槽,但你必须手动将副本上的复制槽与主服务器上对应的复制槽进行同步。副本复制槽的同步不是自动完成的。

PostgreSQL 17 或更高版本

当你将 Debezium 与 PostgreSQL 17 或更高版本配合使用时,可以为主服务器上的复制槽配置自动故障转移,从而确保 Debezium 不会遗漏任何变更事件。当复制槽配置了故障转移功能后,PostgreSQL 会自动将复制槽从主服务器同步到副本,使 Debezium 能够在副本被提升为新的主服务器后继续从该复制槽读取数据。

某些托管型 PostgreSQL 服务(例如 AWS RDS 和 GCP CloudSQL)使用磁盘复制来实现到备用实例的复制。因此,这些服务会自动复制复制槽,使其在故障转移后依然可用。

新的主服务器必须安装逻辑解码插件,并且已配置一个可供该插件使用、且对应于你想要捕获变更的数据库的复制槽。只有满足这些条件后,你才能将连接器指向新的服务器并重启连接器。

从 PostgreSQL 17 集群的故障中恢复

运行 PostgreSQL 17 或更高版本的环境支持使用故障转移复制槽。如果在 PostgreSQL 17 或更高版本的集群中发生故障,且备用实例已配置了故障转移复制槽,请完成以下步骤,使 Debezium 能够恢复捕获:

  1. 暂停 Debezium,直到你能够确认复制槽完好无损且未丢失数据。
  2. 在应用程序开始向新的主服务器写入数据之前,重新创建 Debezium 复制槽。虽然连接器在启动后通常能够自动创建缺失的复制槽,但为了防止连接器在故障转移恢复后遗漏变更事件,你必须在客户端恢复向数据库写入之前手动重新创建缺失的复制槽。
  3. 重启连接器。
  4. 验证 Debezium 能够从复制槽中读取原始主服务器发生故障之前所有变更的 LSN。
    例如,从故障发生前一刻的备份中恢复故障的主服务器,并确定复制槽中记录的最后一个位置。虽然获取备份数据在管理上可能比较困难,但检查备份提供了一种可靠地确定 Debezium 是否已消费全部变更的机制。

在 PostgreSQL 15 或更早版本的集群发生故障后从新的主服务器捕获数据

在 PostgreSQL 15 或更早版本的集群中,当主服务器发生故障后,你可能需要决定配置 Debezium,使其从前从副本服务器之一采集数据,而不是从原来的主服务器采集。要让 Debezium 能从前从副本服务器采集数据,请完成以下步骤。

操作步骤

  1. 修复导致集群故障的问题。

  2. 在连接器停止的状态下,更新连接器配置中的属性值,以反映新服务器的详细信息。例如,确认配置中包含以下属性的正确值:

  3. 通过完成以下任务,配置新的主服务器以与 Debezium 配合工作:

    • 在写入恢复之前,确保服务器上存在所需的复制槽。在正常运行期间,如果复制槽尚不存在,连接器会在启动时自动创建;但在故障恢复期间,你可能需要手动创建该槽以保持连续性。

    • 确保 Debezium 能在该服务器上执行复制并创建发布。

      有关如何配置 PostgreSQL 服务器以与 Debezium 配合工作的更多信息,请参阅 PostgreSQL 设置。

  4. 将 PostgreSQL 备用节点提升为 primary。

  5. 重启连接器。

  6. 将快照模式设置为 always,并在新的主服务器上执行快照,以捕获该服务器上数据的初始状态,确保不会丢失任何数据。

Kafka Connect 进程正常停止

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

Kafka Connect 进程崩溃

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

由于在故障恢复过程中有可能产生重复事件,消费者应始终预期会出现一些重复事件。Debezium 的变更是幂等的,因此一连串事件始终会产生相同的状态。

在每条变更事件记录中,Debezium 连接器都会插入与来源相关的事件源信息,包括 PostgreSQL 服务器上的事件时间、服务器事务 ID,以及事务变更写入预写日志(WAL)时的位置。消费者可以跟踪这些信息(尤其是 LSN),以判断某个事件是否为重复事件。

Kafka 变得不可用

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

连接器被停止一段时间

如果连接器被正常停止,数据库仍可继续使用。任何变更都会记录在 PostgreSQL 的预写日志(WAL)中。当连接器重新启动时,它会从中断处继续流式传输变更。也就是说,它会为连接器停止期间发生的所有数据库变更生成变更事件记录。

经过适当配置的 Kafka 集群能够处理海量吞吐。Kafka Connect 是按照 Kafka 最佳实践编写的,在资源充足的情况下,Kafka Connect 连接器也能够处理数量非常庞大的数据库变更事件。正因如此,Debezium 连接器停止一段时间后重新启动时,很可能会快速追赶在停止期间发生的数据库变更。追赶速度取决于 Kafka 的能力和性能,以及 PostgreSQL 中数据的变更量。

评论

登录后参与评论

正在加载评论…