源连接器

MongoDB

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

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

Debezium MongoDB 连接器

Debezium 的 MongoDB 连接器会跟踪 MongoDB 副本集或 MongoDB 分片集群中数据库和集合的文档变更,并将这些变更记录为 Kafka 主题中的事件。该连接器会自动处理分片集群中分片的添加或移除、每个副本集成员的变化、每个副本集内的选举,以及等待通信问题的解决。

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

概述

MongoDB 的复制机制提供了冗余和高可用性,是在生产环境中运行 MongoDB 的首选方式。MongoDB 连接器会捕获副本集或分片集群中的变更。

MongoDB 副本集由一组拥有相同数据副本的服务器组成,复制机制确保客户端对副本集主节点上文档所做的所有变更都能正确地应用到该副本集的其他服务器上,这些服务器称为从节点。MongoDB 复制的工作方式是:主节点将其变更记录到oplog(即操作日志)中,然后每个从节点读取主节点的 oplog,并按顺序将所有操作应用到自身的文档上。当有新服务器加入副本集时,该服务器首先对主节点上的所有数据库和集合执行一次快照,然后读取主节点的 oplog,以应用自开始快照以来可能发生的全部变更。当该新服务器追上主节点 oplog 的末尾时,它就成为一个从节点(并能够处理查询)。

变更流

尽管 Debezium MongoDB 连接器并不属于副本集的一部分,但它使用类似的复制机制来获取 oplog 数据。主要区别在于,该连接器并不直接读取 oplog。相反,它将 oplog 数据的捕获和解码工作委托给 MongoDB 的变更流功能。通过变更流,MongoDB 服务器会将集合中发生的变更作为事件流暴露出来。Debezium 连接器监视该事件流,然后将变更传递到下游。连接器首次检测到副本集时,会检查 oplog 以获取最后记录的事务,然后对主节点的数据库和集合执行快照。连接器完成数据复制后,会从先前读取到的 oplog 位置开始创建一个变更流。

MongoDB 连接器在处理变更时,会周期性地记录事件在 oplog 流中产生的位置。当连接器停止时,它会记录最后处理到的 oplog 流位置,以便在重启后能从该位置继续流式传输。换句话说,连接器可以被停止、升级或维护,稍后再重启,并且总能准确地从中断处继续,不会丢失任何一个事件。当然,MongoDB 的 oplog 通常有最大容量限制,因此如果连接器长时间停止运行,oplog 中的操作可能会在连接器有机会读取之前就被清除。在这种情况下,重启后连接器会检测到缺失的 oplog 操作,执行快照,然后继续流式传输变更。

MongoDB 连接器对副本集成员和主节点的变化、分片集群中分片的添加或移除,以及可能导致通信故障的网络问题也具有较强的容错能力。默认情况下,连接器从副本集的主节点流式传输变更,但你可以通过在连接字符串中设置读取偏好来配置它从从节点读取。如果连接器所读取的副本集成员变得不可用——例如副本集选举产生了新的主节点——连接器会停止流式传输,并根据配置的读取偏好自动重新连接到另一个合适的成员。同样,如果连接器无法与副本集通信,它会尝试使用指数退避策略重新连接,以免给网络或副本集造成过大压力。连接重新建立后,连接器会从最后捕获的事件继续流式传输变更。通过这种方式,连接器能够动态适应副本集成员的变化,并自动处理通信中断。

其他资源

读取偏好

连接器使用 MongoDB 变更流来捕获数据变更。由于变更流 API 可以从从节点读取数据,因此你可以配置连接器,使其从副本集的任意成员捕获变更。通过在 mongodb.connection.string 中配置 MongoDB 的 readPreference 参数,可以指定连接器读取的副本集成员。如果未指定读偏好,则会采用驱动程序的默认值 primary,连接器将从主节点流式传输变更。

为了减轻主节点的负载,可以通过设置 readPreference=secondaryPreferred 来配置连接器从从节点读取数据。在正常运行条件下,配置连接器从从节点读取数据不会显著增加数据捕获延迟。不过,由于资源争用或网络延迟,从从节点捕获的事件可能比从主节点捕获的事件更晚到达。此外,还可能出现额外的延迟,因为变更流只有在变更达到多数提交之后才会发出事件。当多数投票副本集成员确认已将变更写入各自的日志时,该变更即达到多数提交。因此,从节点只有在收到该变更已达多数提交的确认之后,才能发出事件。

MongoDB 连接器的工作原理

为了最佳地配置和运行 Debezium MongoDB 连接器,了解连接器支持的 MongoDB 拓扑结构、连接器如何使用逻辑名称并执行快照与增量快照,以及如何将变更事件流式传输到 Kafka 主题,都是非常有帮助的。

支持的 MongoDB 拓扑结构

Debezium MongoDB 连接器可以捕获 MongoDB 副本集或 MongoDB 分片集群中的变更。了解这些部署类型之间的差异,有助于你正确配置连接器的连接字符串和捕获范围。

MongoDB 副本集

Debezium MongoDB 连接器可以捕获单个 MongoDB 副本集中的变更。生产环境中的副本集至少需要 三个成员。

要将 MongoDB 连接器与副本集配合使用,你必须将连接器配置中 mongodb.connection.string 属性的值设置为副本集连接字符串。当连接器准备好开始从 MongoDB 变更流捕获变更时,它会启动一个连接任务。该连接任务随后使用指定的连接字符串与可用的副本集成员建立连接。

MongoDB 分片集群

MongoDB 分片集群由以下部分组成:

  • 一个或多个分片(shard),每个分片都以副本集的形式部署;
  • 一个单独的副本集,充当集群的配置服务器(configuration server)
  • 一个或多个路由器(router,也称为 mongos),客户端连接到路由器,由路由器将请求路由到相应的分片

要在分片集群上使用 MongoDB 连接器,请在连接器配置中将 mongodb.connection.string 属性的值设置为分片集群连接字符串。

MongoDB 独立服务器

MongoDB 连接器无法监控独立 MongoDB 服务器的更改,因为独立服务器没有 oplog。如果将独立服务器转换为包含一个成员的副本集,连接器即可正常工作。

MongoDB 不建议在生产环境中运行独立服务器。有关更多信息,请参阅 MongoDB 文档。

所需的用户权限

要从 MongoDB 捕获数据,Debezium 会以 MongoDB 用户的身份连接到数据库。为 Debezium 创建的 MongoDB 用户账户需要特定的数据库权限才能从数据库读取数据。连接器用户需要以下权限:

  • 从数据库读取数据。
  • 运行 hello 命令。

连接器用户可能还需要以下权限:

  • 读取 config.shards 系统集合。

数据库读取权限

连接器用户必须能够读取所有数据库,或读取特定数据库,具体取决于连接器的 capture.scope 属性的值。根据 capture.scope 的设置,为用户授予以下权限之一:

capture.scope 设置为 deployment

授予用户读取任意数据库的权限。

capture.scope 设置为 database

授予用户读取连接器的 capture.target 属性所指定数据库的权限。

capture.scope 设置为 collection

授予用户读取连接器的 capture.target 属性所指定集合的权限。

使用 MongoDB hello 命令的权限

无论 capture.scope 如何设置,用户都需要运行 MongoDB hello 命令的权限。

读取 config.shards 集合的权限

根据你的 Debezium 环境,要使连接器能够执行偏移量合并,必须向连接器用户授予显式权限以读取 config.shards 集合。以下连接器环境需要读取 config.shards 集合的权限:

  • 从 Debezium 2.5 或更早版本升级而来的连接器。
  • 配置为从分片 MongoDB 集群捕获变更的连接器。

逻辑连接器名称

连接器配置属性 topic.prefix 用作 MongoDB 副本集或分片集群的逻辑名称。连接器在多个方面使用该逻辑名称:作为所有主题名称的前缀,以及作为记录每个副本集变更流位置时的唯一标识符。

你应该为每个 MongoDB 连接器指定一个唯一且能清晰描述源 MongoDB 系统的逻辑名称。我们建议逻辑名称以字母或下划线字符开头,其余字符使用字母、数字或下划线。

偏移量合并

Debezium MongoDB 连接器不再支持对分片 MongoDB 部署使用 replica_set 连接。因此,使用 replica_set 连接模式的连接器版本所记录的偏移量与当前版本不兼容。

为了将连接模式变更的影响降到最低,并防止连接器执行不必要的快照,连接器在升级后重启时会运行一个用于合并偏移量的过程。在此偏移量合并过程中,连接器会完成以下步骤,以调和早期连接器版本记录的偏移量:

  1. 由 2.5 之后的连接器版本记录的偏移量按原样使用。

  2. 在 sharded 连接模式下从分片 MongoDB 部署或 MongoDB 副本集部署捕获的事件所对应的偏移量按原样使用。

  3. 由 2.5.x 及更早连接器版本记录的分片特定偏移量,如果同时满足以下条件,则按原样使用:

    • 所有当前数据库分片都存在对应的偏移量。
    • 偏移量失效已启用。如果偏移量失效被禁用,连接器将无法启动。
  4. 连接器在前述步骤中处理完现有偏移量后,会恢复流式传输变更,然后为其捕获的新事件提交偏移量。如果偏移量合并过程未检测到任何现有偏移量,连接器将执行初始快照。

执行快照

当 Debezium 任务开始使用副本集时,它会使用连接器的逻辑名称和副本集名称来查找一个 offset(偏移量),该 offset 描述了连接器之前停止读取变更的位置。如果能找到 offset 且该位置仍存在于 oplog 中,任务就会立即开始流式传输变更,从记录的 offset 位置开始。

但是,如果没有找到 offset,或者 oplog 中已不再包含该位置,任务就必须先通过执行快照来获取副本集内容的当前状态。此过程首先记录 oplog 的当前位置,并将其记录为 offset(同时附带一个表示快照已开始的标志)。然后任务会复制每个集合,并尽可能多地生成线程(最多不超过 snapshot.max.threads 配置属性指定的数量)以并行执行此工作。连接器会为其看到的每个文档记录一个单独的读取事件。每个读取事件都包含对象的标识符、对象的完整状态,以及关于该对象所在 MongoDB 副本集的来源信息。来源信息中还包含一个标志,表示该事件是在快照期间产生的。

快照会持续进行,直到复制完符合连接器过滤条件的所有集合。如果在任务的快照完成之前停止了连接器,重新启动时连接器会再次开始快照。

在连接器对任何副本集执行快照期间,请尽量避免任务重新分配和重新配置。连接器会生成日志消息来报告快照的进度。为了获得最大的控制权,请为每个连接器运行单独的 Kafka Connect 集群。

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

设置说明

always

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

initial

连接器启动后,执行初始数据库快照。

initial_only

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

never

已弃用,请参阅 no_data。

no_data

连接器捕获所有相关表的结构,但不会创建 READ 事件来表示连接器启动时的数据集。

when_needed

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

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

configuration_based

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

custom

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

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

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

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

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

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

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

临时快照信号指定要包含在快照中的集合。快照可以捕获整个数据库的内容,也可以仅捕获数据库中集合的子集。

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

字段默认值值
typeincremental指定要运行的快照类型。目前,你可以请求 incremental 或 blocking 快照。
data-collections不适用一个数组,包含与要包含在快照中的集合完全限定名称相匹配的正则表达式。对于 MongoDB 连接器,请使用以下格式指定集合的完全限定名称:database.collection。

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

触发即时增量快照

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

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

触发即时阻塞快照

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

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

增量快照

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

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

在增量快照进行过程中,Debezium 使用水位标记(watermark)来跟踪进度,并记录其捕获的每个集合行。与标准的初始快照流程相比,这种分阶段捕获数据的方式具有以下优势:

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

增量快照流程

运行增量快照时,Debezium 会按主键对每个集合进行排序,然后根据配置的分块大小将集合切分为多个分块。它逐个分块地工作,捕获分块中的每一行数据。对于捕获的每一行,快照都会发出一个 READ 事件。该事件表示该分块的快照开始时该行的值。

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

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

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

快照窗口

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

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

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

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

要让 Debezium 执行增量快照,必须授予连接器向信号集合写入的权限。

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

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

增量快照要求每张表的主键具有稳定的排序顺序。由于 String 字段可能包含特殊字符,并且会受到不同编码的影响,基于字符串的主键难以按一致且可预测的顺序进行排序。执行增量快照时,最好将主键设置为 String 以外的数据类型。

有关 MongoDB 中 BSON 字符串类型的更多信息,请参阅 MongoDB 文档)。

分片集群的增量快照

要在分片的 MongoDB 集群上使用增量快照,必须将 incremental.snapshot.chunk.size 设置为足够大的值,以抵消变更流管道复杂度增加所带来的影响。

触发增量快照

要发起增量快照,你可以向源数据库上的信号集合发送临时快照信号。

你可以使用 MongoDB 的 insert() 方法向信号集合中提交信号。

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

你提交的查询指定了要包含在快照中的集合,还可以选择指定快照操作的类型。目前,快照操作唯一有效的选项是 incremental 和 blocking。

要指定要包含在快照中的集合,请提供一个 data-collections 数组,其中列出各个集合,或者列出用于匹配集合的正则表达式数组,例如 {"data-collections": ["public.Collection1", "public.Collection2"]}。

增量快照信号的 data-collections 数组没有默认值。如果 data-collections 数组为空,Debezium 会检测到无需执行任何操作,从而不会执行快照。

如果要包含在快照中的集合,其数据库或表名称中含有点号(.),要将该集合添加到 data-collections 数组中,就必须用双引号对名称的每个部分进行转义。例如,要包含存在于 public 数据库中、名称为 My.Collection 的数据集合,请使用以下格式:"public"."My.Collection"。

前提条件

操作步骤

  1. 使用以下语法,将执行快照的信号文档插入到源信号集合中:

    <signalDataCollection>.insert({"id" : _<idNumber>,"type" : <snapshotType>, "data" : {"data-collections" ["<collectionName>", "<collectionName>"],"type": <snapshotType>, "additional-conditions" : [{"data-collections" : "<collectionName>", "filter" : "<additional-condition>"}] }});

例如,

db.debeziumSignal.insert({
"type" : "execute-snapshot",
"data" : {
"data-collections" ["\"public\".\"Collection1\"", "\"public\".\"Collection2\""],
"type": "incremental"}
"additional-conditions":[{"data-collection": "public.Collection1" ,"filter":"color=\'blue\'"}]}');
});

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

以下列表描述了增量快照示例中的各个参数:

db.debeziumSignal.insert

插入命令,用于指定要向其中添加快照信号文档的信号集合。该数据集合通过其全限定名来指定。

type

指定要发送的信号类型,此处为 execute-snapshot 信号。

_id

一个可选参数,用于指定任意字符串作为信号请求的标识符。在上面的示例中,信号省略了可选的 _id 参数。由于信号未显式地为该参数赋值,MongoDB 会自动为文档分配一个随机 id。该 id 就成为信号请求的标识符。Debezium 不会使用这个标识字符串;但你可以使用 id 将日志消息与信号集合中的文档关联起来。

data.data-collections

执行快照信号的 data 字段中的一个必填组件,用于指定一个集合名称或正则表达式数组,表示要包含在快照中的集合,此处为 public.Collection1 和 public.Collection2。

data.type

信号的 data 字段中一个可选的 type 组件,用于指定要运行的快照操作类型。有效值为 incremental 或 blocking。

additional-conditions

一个可选数组,用于指定一组附加条件,连接器根据这些条件评估要包含在快照中的记录子集。你可以指定一个或多个数据集合,并为每个集合设置过滤参数。每个过滤条件指示快照仅捕获包含指定字段值的文档。在此示例中,该参数指示信号仅当 public.Collection1 集合中的文档包含值为 'blue' 的 color 字段时才捕获它们。

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

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

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

以下列表描述了前面事件消息示例中的部分参数:

snapshot

指定快照操作的类型,在此例中为增量快照。

op

指定事件类型。快照事件的值为 r,表示 READ 操作。

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

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

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

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

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

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

字段 默认值

type

incremental

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

data-collections

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> …​.

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

Key = `test_connector`

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

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

Key = `test_connector`

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

停止增量快照

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

要向信号集合提交停止快照信号,你需要向该集合插入一个停止快照信号文档。你提交的停止快照信号会将快照操作的 type 指定为 incremental,并可选择指定你希望从当前正在运行的快照中排除的集合。Debezium 检测到信号集合发生变化后,会读取该信号,并在增量快照操作正在进行时将其停止。

其他资源

你还可以通过向 Kafka 信号主题 发送 JSON 消息来停止增量快照。

先决条件

操作步骤

  1. 向信号集合插入一个停止快照信号文档:

    <signalDataCollection>.insert({"id" : _<idNumber>,"type" : "stop-snapshot", "data" : {"data-collections" ["<collectionName>", "<collectionName>"],"type": "incremental"}});

例如,

db.debeziumSignal.insert({
"type" : "stop-snapshot",
"data" : {
"data-collections" ["\"public\".\"Collection1\"", "\"public\".\"Collection2\""],
"type": "incremental"}
});

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

以下列表描述了示例中的各参数:

db.debeziumSignal.insert

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

type

指定要发送的信号种类,本例中为 stop-snapshot 信号。

_id

此可选字段在上面的信号中被省略。由于文档未显式地为该参数赋值,MongoDB 会自动为文档分配一个任意 id。这个 _id 成为信号请求的标识符。Debezium 不使用该标识字符串;但你可以使用 id 来将日志消息与信号集合中的文档关联起来。

data.data-collections

信号 data 字段中的可选组件,用于指定一个集合名称或正则表达式数组,表示要从快照中排除的集合。数组中的正则表达式以 database.collection 格式匹配集合的完全限定名称。

data.type

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

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

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

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

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

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

字段默认值说明
typeincremental要执行的快照类型。目前 Debezium 仅支持 incremental 类型。详情请参阅下一节。
data-collections无可选的字符串数组,包含用于匹配要从快照中移除的表的完全限定名称的正则表达式,或用于匹配要移除的集合名称的集合名称数组或正则表达式数组。请使用 database.collection 格式指定集合名称。

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

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

Key = `test_connector`

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

自定义快照器 SPI

要自定义标准快照模式之外的快照行为,你可以实现一个或多个 Debezium 快照器 SPI 接口。这些接口可以控制是否执行快照、如何查询数据,以及是否锁定表。

io.debezium.snapshot.spi.Snapshotter

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

io.debezium.snapshot.spi.SnapshotQuery

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

io.debezium.snapshot.spi.SnapshotLock

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

}

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

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

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

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

}

阻塞式快照

阻塞式快照允许你在连接器运行时按需捕获集合的完整、一致的快照,并在快照完成之前临时暂停流式传输。阻塞式快照依赖于 Debezium 的向 Debezium 连接器发送信号机制。

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

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

  • 你新增了一个集合,并希望在连接器运行期间完成快照。
  • 你新增了一个大型集合,并希望快照能比增量快照更快完成。

阻塞式快照流程

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

配置快照

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

  • data-collections:指定哪些集合必须进行快照。

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

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

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

例如:

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

可能的重复

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

流式传输变更

在副本集的连接器任务记录偏移量之后,它会使用该偏移量来确定应在 oplog 中的哪个位置开始流式传输变更。随后任务会打开一个 MongoDB 变更流,并从该偏移位置开始流式传输变更。根据 capture.scope 设置的值,任务会从整个部署、特定数据库或某个集合流式传输变更。变更流由与配置的读偏好(默认为 primary)相匹配的副本集成员提供服务。有关更多信息,请参阅 read-preference。任务会处理所有 create、insert 和 delete 操作,并将它们转换为 Debezium 变更事件。每个变更事件都包含该操作在 oplog 中的位置,连接器会定期将其记录为最新的偏移量。偏移量的记录间隔由 Kafka Connect 工作进程配置属性 offset.flush.interval.ms 控制。

当连接器被优雅地停止时,最后处理的偏移量会被记录下来,这样在重新启动时,连接器就能从中断的位置继续。然而,如果连接器的任务意外终止,那么任务可能在最后一次记录偏移量之后、但在该偏移量被记录之前已经处理并生成了事件;重新启动时,连接器将从最后已记录的偏移量开始,可能会再次生成崩溃前刚刚生成过的部分相同事件。

当 Kafka 管道中的所有组件都正常运行时,Kafka 消费者会恰好一次接收每条消息。但是,当出现问题时,Kafka 只能保证消费者至少一次接收到每条消息。为避免意外结果,消费者必须能够处理重复消息。

如读取偏好中所述,连接器从 MongoDB 连接字符串中 readPreference 参数所指定的副本集成员流式传输变更。默认情况下,连接器从主节点捕获变更,这样捕获操作的延迟最低。你也可以选择配置连接器从从节点读取,但从从节点捕获可能因复制延迟而引入额外的延迟。有关更多信息,请参阅读取偏好。如果连接器读取的副本集成员变得不可用,连接器会停止流式传输变更,重新连接到另一个合适的成员,并从相同的位置恢复流式传输。同样,如果连接器在与副本集成员通信时遇到任何问题,它会尝试重新连接,使用指数退避以避免给副本集造成过大压力;连接成功后,它会从中断处继续流式传输变更。通过这种方式,连接器能够动态适应副本集成员的变化,并自动处理通信故障。

总而言之,MongoDB 连接器在大多数情况下都会持续运行。通信问题可能会导致连接器等待,直到问题得到解决。

预映像支持

在 MongoDB 6.0 及更高版本中,你可以配置变更流,使其发出文档的预映像状态,以便填充 MongoDB 变更事件的 before 字段。要在 MongoDB 中启用预映像的使用,必须通过 db.createCollection()、create 或 collMod 为集合设置 changeStreamPreAndPostImages。要使 Debezium MongoDB 连接器在变更事件中包含预映像,请将连接器的 capture.mode 设置为 *_with_pre_image 选项之一。

MongoDB 变更流事件的大小限制

MongoDB 变更流事件的大小上限为 16 兆字节。因此,使用预映像会增加超过该阈值的可能性,从而导致失败。有关如何避免超过变更流限制的信息,请参阅 MongoDB 文档。

主题名称

MongoDB 连接器将每个集合中所有文档的插入、更新和删除操作所产生的事件写入到单个 Kafka 主题中。Kafka 主题的名称始终采用 logicalName.databaseName.collectionName 的形式,其中 logicalName 是通过 topic.prefix 配置属性指定的连接器逻辑名称,databaseName 是发生操作的数据库名称,collectionName 是受影响文档所在的 MongoDB 集合名称。

例如,考虑一个 MongoDB 副本集,其中包含一个 inventory 数据库,该数据库有四个集合:products、products_on_hand、customers 和 orders。如果监控该数据库的连接器被赋予逻辑名称 fulfillment,则该连接器将在这四个 Kafka 主题上产生事件:

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

请注意,主题名称中不包含副本集名称或分片名称。因此,对分片集合(每个分片包含该集合文档的子集)的所有更改都会发送到同一个 Kafka 主题。

你可以配置 Kafka 以便按需自动创建主题。如果没有启用自动创建,则必须在启动连接器之前使用 Kafka 管理工具手动创建主题。

分区

MongoDB 连接器不会对事件的主题分区方式做出任何显式决定。相反,它允许 Kafka 根据事件键来决定主题的分区方式。你可以通过在 Kafka Connect 工作器配置中定义 Partitioner 实现的名称来更改 Kafka 的分区逻辑。

Kafka 仅保证写入单个主题分区的事件具有全序性。按键对事件进行分区意味着具有相同键的所有事件总是发送到同一分区。这确保了特定文档的所有事件始终保持全序。

事务元数据

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

Debezium 接收事务元数据的限制

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

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

status

BEGIN 或 END

id

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

event_count(针对 END 事件)

该事务发出的事件总数。

data_collections(针对 END 事件)

一个由 data_collection 和 event_count 组成的配对数组,提供源自给定数据集合的变更所发出的事件数量。

以下示例展示了一条典型的消息:

{
  "status": "BEGIN",
  "id": "1462833718356672513",
  "event_count": null,
  "data_collections": null
}

{
  "status": "END",
  "id": "1462833718356672513",
  "event_count": 2,
  "data_collections": [
    {
      "data_collection": "rs0.testDB.collectiona",
      "event_count": 1
    },
    {
      "data_collection": "rs0.testDB.collectionb",
      "event_count": 1
    }
  ]
}

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

变更数据事件的丰富处理

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

id

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

total_order

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

data_collection_order

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

下面是消息示例:

{
  "after": "{\"_id\" : {\"$numberLong\" : \"1004\"},\"first_name\" : \"Anne\",\"last_name\" : \"Kretchmar\",\"email\" : \"[email protected]\"}",
  "source": {
...
  },
  "op": "c",
  "ts_ms": "1580390884335",
  "ts_us": "1580390884335486",
  "ts_ns": "1580390884335486281",
  "transaction": {
    "id": "1462833718356672513",
    "total_order": "1",
    "data_collection_order": "1"
  }
}

数据变更事件

Debezium MongoDB 连接器会针对每个文档级的插入、更新或删除操作生成一个数据变更事件。每个事件都包含一个键和一个值,键和值的结构取决于被修改的集合。

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

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

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

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

schema(第一个)

第一个 schema 字段是事件键的一部分。它指定了一个 Kafka Connect schema,用于描述事件键 payload 部分中的内容。换句话说,第一个 schema 字段描述了被更改文档的键的结构。

payload(第一个)

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

schema(第二个)

第二个 schema 字段是事件值的一部分。它指定了描述事件值 payload 部分内容的 Kafka Connect schema。换句话说,第二个 schema 描述了被更改文档的结构。通常,该 schema 包含嵌套的 schema。

payload(第二个)

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

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

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

如果逻辑服务器名称、数据库名或集合名包含无效字符,而唯一能够区分这些名称的字符恰好是无效的并因此被替换为下划线,就可能导致意外的冲突。

变更事件键

变更事件的键包含被更改文档的键所对应的 schema 以及被更改文档的实际键。对于给定的集合,该 schema 及其对应的 payload 都只包含一个 id 字段。该字段的值是文档的标识符,表示为一个字符串,派生自 MongoDB 扩展 JSON 序列化严格模式。

考虑一个逻辑名称为 fulfillment 的连接器,其副本集中包含一个 inventory 数据库,以及一个包含如下文档的 customers 集合。

示例文档

{
  "_id": 1004,
  "first_name": "Anne",
  "last_name": "Kretchmar",
  "email": "[email protected]"
}

更改事件键示例

每个捕获 customers 集合变更的更改事件都具有相同的事件键架构。只要 customers 集合保留前面的定义,每个捕获 customers 集合变更的更改事件就具有以下键结构。用 JSON 表示如下:

{
  "schema": {
    "type": "struct",
    "name": "fulfillment.inventory.customers.Key",
    "optional": false,
    "fields": [
      {
        "field": "id",
        "type": "string",
        "optional": false
      }
    ]
  },
  "payload": {
    "id": "1004"
  }
}

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

schema

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

fulfillment.inventory.customers.Key

定义键 payload 结构的 schema 名称。此 schema 描述被更改文档的键结构。键 schema 名称的格式为 connector-name.database-name.collection-name.Key。在前面的示例中,name 字段包含以下值:

fulfillment

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

inventory

包含被操作更改的集合的数据库。

customers

包含被操作更新的文档的集合。

schema.optional

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

schema.fields

指定 payload 中预期的每个字段,包括每个字段的名称、类型以及是否为必填项。

schema.payload

包含生成此更改事件的文档的键。在此示例中,键包含一个类型为 string、值为 1004 的 id 字段。

此示例使用带整数标识符的文档,但任何有效的 MongoDB 文档标识符的工作方式都相同,包括文档标识符。对于文档标识符,事件键的 payload.id 值是一个字符串,该字符串以使用严格模式的 MongoDB 扩展 JSON 序列化形式表示被更新文档的原始 _id 字段。下表提供了不同类型的 _id 字段表示方式的示例。

类型MongoDB _id 值键的 payload
整数1234{ "id" : "1234" }
浮点数12.34{ "id" : "12.34" }
字符串"1234"{ "id" : "\"1234\"" }
文档{ "hi" : "kafka", "nums" : [10.0, 100.0, 1000.0] }{ "id" : "{\"hi\" : \"kafka\", \"nums\" : [10.0, 100.0, 1000.0]}" }
ObjectIdObjectId("596e275826f08b2730779e1f"){ "id" : "{\"$oid\" : \"596e275826f08b2730779e1f\"}" }
二进制BinData("a2Fma2E=",0){ "id" : "{\"$binary\" : \"a2Fma2E=\", \"$type\" : \"00\"}" }

表 5. 在事件键 payload 中表示文档 _id 字段的示例

更改事件值

变更事件中的值比键要复杂一些。与键一样,值也包含 schema 部分和 payload 部分。schema 部分包含描述 payload 部分 Envelope 结构的架构,其中包括其嵌套字段。针对创建、更新或删除数据的操作所产生的变更事件,其值负载都具有信封结构。

来看一下之前用于展示变更事件键示例的同一份示例文档:

示例文档

{
  "_id": 1004,
  "first_name": "Anne",
  "last_name": "Kretchmar",
  "email": "[email protected]"
}

针对此文档的变更,各事件类型的变更事件 value 部分说明如下:

create 事件

每当有新文档插入到集合中时,连接器就会发出 create 事件。事件负载(payload)在 after 字段中包含新文档的完整状态,同时还包含用于标识副本集、数据库和集合的源(source)元数据。

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

{
    "schema": {
      "type": "struct",
      "fields": [
        {
          "type": "string",
          "optional": true,
          "name": "io.debezium.data.Json",
          "version": 1,
          "field": "after"
        },
        {
          "type": "string",
          "optional": true,
          "name": "io.debezium.data.Json",
          "version": 1,
          "field": "patch"
        },
        {
          "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": "rs"
            },
            {
              "type": "string",
              "optional": false,
              "field": "collection"
            },
            {
              "type": "int32",
              "optional": false,
              "field": "ord"
            },
            {
              "type": "int64",
              "optional": true,
              "field": "h"
            }
          ],
          "optional": false,
          "name": "io.debezium.connector.mongo.Source",
          "field": "source"
        },
        {
          "type": "string",
          "optional": true,
          "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": "dbserver1.inventory.customers.Envelope"
      },
    "payload": {
      "after": "{\"_id\" : {\"$numberLong\" : \"1004\"},\"first_name\" : \"Anne\",\"last_name\" : \"Kretchmar\",\"email\" : \"[email protected]\"}",
      "source": {
        "version": "3.6.3.Final",
        "connector": "mongodb",
        "name": "fulfillment",
        "ts_ms": 1558965508000,
        "ts_ms": 1558965508000000,
        "ts_ms": 1558965508000000000,
        "snapshot": false,
        "db": "inventory",
        "rs": "rs0",
        "collection": "customers",
        "ord": 31,
        "h": 1546547425148721999
      },
      "op": "c",
      "ts_ms": 1558965515240,
      "ts_us": 1558965515240142,
      "ts_ns": 1558965515240142879,
    }
  }

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

schema

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

name(io.debezium.data.Json)

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

io.debezium.data.Json 是负载中 after、patch 和 filter 字段的 schema。该 schema 专属于 customers 集合。create 事件是唯一包含 after 字段的事件类型。update 事件包含 filter 字段和 patch 字段。delete 事件包含 filter 字段,但不包含 after 字段和 patch 字段。

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

io.debezium.connector.mongo.Source 是负载中 source 字段的 schema。该 schema 专属于 MongoDB 连接器。连接器对其生成的所有事件都使用该 schema。

name(dbserver1.inventory.customers.Envelope)

dbserver1.inventory.customers.Envelope 是负载整体结构的 schema,其中 dbserver1 是连接器名称,inventory 是数据库,customers 是集合。该 schema 专属于该集合。

payload

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

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

payload.after

一个可选字段,用于指定事件发生后文档的状态。在本示例中,after 字段包含新文档的 _id、first_name、last_name 和 email 字段的值。after 的值始终是一个字符串。按照惯例,它包含文档的 JSON 表示形式。MongoDB oplog 条目仅在 create 事件中包含文档的完整状态,在 update 事件中也仅当 capture.mode 选项设置为 change_streams_update_full 时才包含完整状态;换句话说,无论 capture.mode 选项如何设置,create 事件都是唯一包含 after 字段的事件类型。

payload.source

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

  • Debezium 版本。
  • 生成该事件的连接器名称。
  • MongoDB 副本集的逻辑名称,它构成所生成事件的命名空间,并用于连接器写入的 Kafka 主题名称。
  • 包含新文档的集合与数据库名称。
  • 该事件是否属于快照的一部分。
  • 数据库中发生变更的时间戳,以及该事件在同一时间戳内的序号。
  • MongoDB 操作的唯一标识符(oplog 事件中的 h 字段)。
  • 该变更是在事务中执行时,MongoDB 会话 lsid 与事务号 txnNumber 的唯一标识符(仅适用于变更流捕获模式)。

payload.source.op

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

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

ts_ms、ts_us、ts_ns

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

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

update(更新)事件

每当集合中的现有文档被修改时,连接器就会发出 update 事件。事件有效负载的内容取决于 capture.mode 连接器属性的值,该属性用于控制是否包含完整文档状态和前像。

变更流捕获模式

在示例 customers 集合中,某个更新的变更事件值具有与该集合的 create(创建)事件相同的模式。同样,事件值的有效负载结构也相同。不过,update 事件中事件值有效负载包含的值有所不同。只有当 capture.mode 选项设置为 change_streams_update_full 时,update 事件才会包含 after 值。如果 capture.mode 选项设置为某个 *_with_pre_image 选项,则会提供 before 值。在这种情况下,会有一个新的结构化字段 updateDescription,其中包含几个额外字段:

  • updatedFields 是一个字符串字段,包含更新后文档字段及其值的 JSON 表示形式
  • removedFields 是一个列表,包含从文档中删除的字段名
  • truncatedArrays 是一个列表,包含文档中被截断的数组

以下是连接器为 customers 集合中的某次更新生成的变更事件值示例:

{
    "schema": { ... },
    "payload": {
      "op": "u",
      "ts_ms": 1465491461815,
      "ts_us": 1465491461815698,
      "ts_ns": 1465491461815698142,
      "before":"{\"_id\": {\"$numberLong\": \"1004\"},\"first_name\": \"unknown\",\"last_name\": \"Kretchmar\",\"email\": \"[email protected]\"}",
      "after":"{\"_id\": {\"$numberLong\": \"1004\"},\"first_name\": \"Anne Marie\",\"last_name\": \"Kretchmar\",\"email\": \"[email protected]\"}",
      "updateDescription": {
        "removedFields": null,
        "updatedFields": "{\"first_name\": \"Anne Marie\"}",
        "truncatedArrays": null
      },
      "source": {
        "version": "3.6.3.Final",
        "connector": "mongodb",
        "name": "fulfillment",
        "ts_ms": 1558965508000,
        "ts_us": 1558965508000000,
        "ts_ns": 1558965508000000000,
        "snapshot": false,
        "db": "inventory",
        "rs": "rs0",
        "collection": "customers",
        "ord": 1,
        "h": null,
        "tord": null,
        "stxnid": null,
        "lsid":"{\"id\": {\"$binary\": \"FA7YEzXgQXSX9OxmzllH2w==\",\"$type\": \"04\"},\"uid\": {\"$binary\": \"47DEQpj8HBSa+/TImW+5JCeuQeRkm5NMpJWZG3hSuFU=\",\"$type\": \"00\"}}",
        "txnNumber":1
      }
    }
  }

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

payload.op

必需的字符串字段,用于描述导致连接器生成该事件的操作类型。在本示例中,u 表示该操作更新了一个文档。

ts_ms、ts_us、ts_ns

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

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

payload.before

包含变更前实际 MongoDB 文档的 JSON 字符串表示。

如果捕获模式未设置为 *_with_preimage 选项之一,update 事件值中将不包含 before 字段。

payload.after

包含实际 MongoDB 文档的 JSON 字符串表示。

如果捕获模式未设置为 change_streams_update_full,update 事件值中将不包含 after 字段。

payload.updateDescription.updatedFields

包含文档已更新字段值的 JSON 字符串表示。连接器使用 json.serialization.mode 连接器属性指定的模式序列化该字段。在本示例中,此次更新将 first_name 字段改为一个新值。

payload.source

必需字段,用于描述事件的源元数据。该字段包含与同一集合的 create 事件相同的信息,但由于该事件来自 oplog 中的不同位置,因此值有所不同。源元数据包括:

  • Debezium 版本。
  • 生成该事件的连接器名称。
  • MongoDB 副本集的逻辑名称,它为生成的事件构成命名空间,并用于连接器写入的 Kafka 主题名称中。
  • 包含已更新文档的集合和数据库名称。
  • 该事件是否属于快照的一部分。
  • 数据库中发生更改的时间戳以及该事件在该时间戳内的序号。
  • MongoDB 会话的唯一标识符 lsid 和事务编号 txnNumber(如果该更改是在事务中执行的)。

事件中的 after 值应视为文档在特定时间点的值。该值并非动态计算得出,而是从集合中获取的。因此,如果多次更新紧密连续发生,所有 update 事件可能包含相同的 after 值,该值代表文档中最后存储的值。

如果你的应用依赖于逐步演进的变更,则应仅依赖 updateDescription。

delete 事件

当集合中的某个文档被删除时,连接器会发出一个 delete 事件。该事件的负载中既不包含 after 字段,也不包含 updateDescription 字段。

delete 变更事件的值与同一集合的 create 和 update 事件具有相同的 schema 部分。而 delete 事件的 payload 部分所包含的值则与同一集合的 create 和 update 事件不同。特别是,delete 事件既不包含 after 值,也不包含 updateDescription 值。下面是一个 customers 集合中文档的 delete 事件示例:

{
    "schema": { ... },
    "payload": {
      "op": "d",
      "ts_ms": 1465495462115,
      "ts_us": 1465495462115748,
      "ts_ns": 1465495462115748263,
      "before":"{\"_id\": {\"$numberLong\": \"1004\"},\"first_name\": \"Anne Marie\",\"last_name\": \"Kretchmar\",\"email\": \"[email protected]\"}",
      "source": {
        "version": "3.6.3.Final",
        "connector": "mongodb",
        "name": "fulfillment",
        "ts_ms": 1558965508000,
        "ts_us": 1558965508000000,
        "ts_ns": 1558965508000000000,
        "snapshot": true,
        "db": "inventory",
        "rs": "rs0",
        "collection": "customers",
        "ord": 6,
        "h": 1546547425148721999
      }
    }
  }

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

payload.op

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

ts_ms、ts_us、ts_ns

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

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

payload.before

包含更改前实际 MongoDB 文档的 JSON 字符串表示。

如果捕获模式未设置为 *_with_preimage 选项之一,则 update 事件值中不会包含 before 字段。

payload.source

必填字段,用于描述事件的源元数据。该字段包含与同一集合的 create 或 update 事件相同的信息,但由于此事件来自 oplog 中的不同位置,因此值有所不同。源元数据包括:

  • Debezium 版本。
  • 生成该事件的连接器名称。
  • MongoDB 副本集的逻辑名称,它构成生成事件的命名空间,并用于连接器写入的 Kafka 主题名称中。
  • 包含被删除文档的集合和数据库名称。
  • 该事件是否属于快照的一部分。
  • 数据库中进行更改的时间戳,以及该事件在该时间戳内的序号。
  • MongoDB 操作的唯一标识符(oplog 事件中的 h 字段)。
  • 如果更改是在事务中执行的,则包含 MongoDB 会话 lsid 和事务号 txnNumber 的唯一标识符(仅限变更流捕获模式)。

MongoDB 连接器事件的设计目的是与 Kafka 日志压缩 协同工作。日志压缩允许删除一些较旧的消息,只要每个键至少保留最新的一条消息即可。这样 Kafka 就能在确保主题包含完整数据集、可用于重新加载基于键的状态的同时,回收存储空间。

墓碑事件

针对唯一标识的文档,所有 MongoDB 连接器事件都具有完全相同的键。当文档被删除时,delete 事件值仍然可以与日志压缩协同工作,因为 Kafka 可以删除所有具有相同键的较早消息。但是,为了让 Kafka 删除所有具有该键的消息,消息值必须为 null。为此,在 Debezium 的 MongoDB 连接器发出 delete 事件之后,连接器会发出一个特殊的墓碑事件,该事件具有相同的键但值为 null。墓碑事件会通知 Kafka 可以删除所有具有相同键的消息。

设置 MongoDB

MongoDB 连接器使用 MongoDB 的变更流(change streams)来捕获变更,因此该连接器仅适用于 MongoDB 副本集,或适用于每个分片都是独立副本集的分片集群。有关设置副本集或分片集群的方法,请参阅 MongoDB 文档。另外,请务必了解如何为副本集启用访问控制与身份验证。

你还需要一个 MongoDB 用户,该用户必须具有相应的角色,以便读取 oplog 所在的 admin 数据库。此外,该用户还必须能够读取分片集群中配置服务器的 config 数据库,并且必须拥有 listDatabases 权限操作。当使用变更流(默认方式)时,用户还必须拥有集群范围的 find 和 changeStream 权限操作。

如果你打算使用前像(pre-image)并填充 before 字段,则需要先通过 db.createCollection()、create 或 collMod 为某个集合启用 changeStreamPreAndPostImages。

云端中的 MongoDB

你可以将 Debezium 的 MongoDB 连接器与 MongoDB Atlas 搭配使用。请注意,MongoDB Atlas 仅支持通过 SSL 建立安全连接,即 +mongodb.ssl.enabled 连接器选项必须设置为 true。

最优 Oplog 配置

Debezium MongoDB 连接器通过读取变更流来获取副本集的 oplog 数据。由于 oplog 是大小固定的封顶集合(capped collection),如果其超出配置的最大大小,它就会开始覆盖最旧的记录。如果连接器因任何原因停止运行,当它重新启动时,它会尝试从上一次 oplog 流位置继续恢复流式传输。但是,如果上一次的流位置已被从 oplog 中移除,那么根据连接器 snapshot.mode 属性所指定的值,连接器可能会启动失败,并报告无效恢复令牌错误。一旦发生此类故障,你必须创建一个新的连接器,才能让 Debezium 继续从数据库中捕获记录。有关更多信息,请参阅如果 snapshot.mode 设置为 initial,连接器在长时间停止后会失败。

为确保 oplog 保留 Debezium 恢复流式传输所需的偏移值,你可以采用以下任一方法:

  • 增大 oplog 的大小。根据你的典型工作负载,将 oplog 大小设置为大于每小时 oplog 条目峰值的数量。
  • 增大 oplog 条目的最短保留小时数(MongoDB 4.4 及更高版本)。该设置基于时间,即使 oplog 达到其配置的最大大小,最近 n 小时内的条目也保证可用。虽然这通常是首选选项,但对于负载较高且已接近容量上限的集群,请指定最大 oplog 大小。

为了帮助防止因缺少 oplog 条目而导致的故障,跟踪反映复制行为的指标并优化 oplog 大小以支持 Debezium 至关重要。特别是,你应该监控 Oplog GB/Hour(每小时 oplog 增长量)和 Replication Oplog Window(复制 oplog 窗口)的值。如果 Debezium 离线的时间超过了复制 oplog 窗口的值,并且主节点 oplog 的增长速度超过 Debezium 消费条目的速度,就可能导致连接器故障。

有关如何监控这些指标的信息,请参阅 MongoDB 文档。

最好将最大 oplog 大小设置为基于 oplog 的预期每小时增长量(Oplog GB/Hour)乘以处理 Debezium 故障可能所需的时间。

即:

Oplog GB/Hour X average reaction time to Debezium failure

例如,如果 oplog 大小限制设置为 1GB,且 oplog 每小时增长 3GB,那么每小时会清理三次 oplog 条目。如果 Debezium 在这段时间内发生故障,其最后的 oplog 位置很可能会被删除。

如果 oplog 以 3GB/小时的速度增长,而 Debezium 离线两小时,那么你应将 oplog 大小设置为 3GB/小时 X 2 小时,即 6GB。

部署

要捕获来自 MongoDB 的变更事件,请将 Debezium MongoDB 连接器部署到 Kafka Connect 环境中,并将其配置为连接到 MongoDB 副本集或分片集群。

要部署 Debezium MongoDB 连接器,你需要安装 Debezium MongoDB 连接器归档文件、配置连接器,并通过将其配置添加到 Kafka Connect 来启动连接器。

前置条件

步骤

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

如果你使用的是不可变容器,可以参考 Debezium 的容器镜像,其中提供了已经安装 MongoDB 连接器、可直接运行的 Apache Kafka 和 Kafka Connect 镜像。

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

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

Debezium 的教程将引导你使用这些镜像,这是了解 Debezium 的绝佳方式。

MongoDB 连接器配置示例

以下是一个连接器实例的配置示例,用于从 192.168.99.100 上 27017 端口的 MongoDB 副本集 rs0 捕获数据,我们在逻辑上将其命名为 fullfillment。通常,你可以通过设置连接器可用的配置属性,在 JSON 文件中配置 Debezium MongoDB 连接器。

你可以选择为特定的 MongoDB 副本集或分片集群生成事件。还可以选择过滤掉不需要的集合。

{
  "name": "inventory-connector",
  "config": {
    "connector.class": "io.debezium.connector.mongodb.MongoDbConnector",
    "mongodb.connection.string": "mongodb://192.168.99.100:27017/?replicaSet=rs0",
    "topic.prefix": "fullfillment",
    "collection.include.list": "inventory[.]*"
  }
}

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

name

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

config.connector.class

MongoDB 连接器类的名称。

config.mongodb.connection.string

用于连接 MongoDB 副本集的连接串。

config.topic.prefix

MongoDB 副本集的逻辑名称,它构成了生成事件的命名空间,并用于连接器写入的所有 Kafka 主题名称、Kafka Connect 模式名称,以及使用 Avro 转换器时相应 Avro 模式的命名空间中。

config.collection.include.list

一组正则表达式,用于匹配要监控的所有集合的集合命名空间(例如 <dbName>.<collectionName>)。此项为可选。

有关 Debezium MongoDB 连接器可设置的全部配置属性,请参阅 MongoDB 连接器配置属性。

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

  • 连接到 MongoDB 副本集或分片集群。
  • 为每个副本集分配任务。
  • 必要时执行快照。
  • 读取变更流。
  • 将变更事件记录流式传输到 Kafka 主题。

添加连接器配置

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

前提条件

操作步骤

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

结果

连接器启动后,将完成以下操作:

  • 执行一致性快照,捕获 MongoDB 副本集中的集合。
  • 读取各副本集的变更流。
  • 为每个插入、更新和删除的文档生成变更事件。
  • 将变更事件记录流式传输到 Kafka 主题。

连接器属性

Debezium MongoDB 连接器提供了众多配置属性,你可以利用它们使连接器的行为符合应用的需求。许多属性都有默认值。属性相关信息按如下方式组织:

除非提供了默认值,否则以下配置属性均为必填项。

Debezium MongoDB 连接器必填配置属性

运行 Debezium MongoDB 连接器需要以下属性。对于所有没有默认值的必填属性,你都必须提供相应的值。

属性 默认值 说明
internal.mongodb.allow.offset.invalidation false 将此属性设置为 true,可使连接器能够使由早期连接器版本记录的分片特定偏移量失效,并对其进行整合。
name 无默认值 连接器的唯一名称。使用相同名称再次注册将会失败。(所有 Kafka Connect 连接器都要求提供此属性。)
connector.class 无默认值 连接器的 Java 类名。对于 MongoDB 连接器,始终使用 io.debezium.connector.mongodb.MongoDbConnector。
mongodb.connection.string 无默认值 指定连接器用于连接 MongoDB 副本集的连接字符串。

此属性取代了 MongoDB 连接器早期版本中可用的 mongodb.hosts 属性。
此属性允许你修改当前的默认行为。如果默认行为发生变更,允许连接器自动使由早期连接器版本记录的偏移量失效并进行整合,则该属性可能会在未来版本中被移除。

topic.prefix

无默认值

用于标识此连接器以及该连接器所监控的 MongoDB 副本集或分片集群的唯一名称。每个服务器最多只能由一个 Debezium 连接器进行监控,因为该服务器名称是源自该 MongoDB 副本集或集群的所有持久化 Kafka 主题的前缀。名称只能由字母、数字、连字符、点和下划线组成。

逻辑名称在所有其他连接器中应当是唯一的,因为该名称会作为前缀用于命名从此连接器接收记录的 Kafka 主题。

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

mongodb.authentication.class

DefaultMongoDbAuthProvider

一个完整的 Java 类名,实现了 io.debezium.connector.mongodb.connection.MongoDbAuthProvider 接口。该类负责设置 MongoDB 连接上的凭据(在每次应用启动时调用)。默认行为是按照 mongodb.user、mongodb.password 和 mongodb.authsource 各自的文档说明使用这些属性,但其他实现可能会以不同方式使用它们,或者完全忽略它们。请注意,mongodb.connection.string 中的任何设置都会覆盖由该类设置的配置

mongodb.user

无默认值

使用默认的 mongodb.authentication.class 时:连接 MongoDB 时要使用的数据库用户名。仅当 MongoDB 配置为启用身份验证时才需要此属性。

mongodb.password

无默认值

使用默认的 mongodb.authentication.class 时:连接 MongoDB 时要使用的密码。仅当 MongoDB 配置为启用身份验证时才需要此属性。

mongodb.authsource

admin

使用默认的 mongodb.authentication.class 时:包含 MongoDB 凭据的数据库(认证源)。只有当 MongoDB 配置为使用 admin 以外的认证数据库进行认证时,才需要设置此项。

mongodb.ssl.enabled

false

连接器将使用 SSL 连接到 MongoDB 实例。

mongodb.ssl.invalid.hostname.allowed

false

启用 SSL 时,此设置用于控制在连接阶段是否禁用严格的主机名检查。如果为 true,连接将无法阻止中间人攻击。

filters.match.mode

regex

用于根据要包含/排除的数据库和集合名称匹配事件的模式。将该属性设置为以下值之一:

regex

数据库和集合的包含/排除规则按逗号分隔的正则表达式列表进行求值。

literal

数据库和集合的包含/排除规则按逗号分隔的字符串字面量列表进行求值。这些字面量周围的空白字符会被去除。

database.include.list

空字符串

一个可选的、以逗号分隔的正则表达式或字面量列表,用于匹配要监控的数据库名称。

默认情况下,会监控所有数据库。当设置了 database.include.list 时,连接器仅监控该属性指定的数据库,其他数据库将不被监控。

为了匹配数据库名称,Debezium 会根据 filters.match.mode 属性的值执行以下操作之一:

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

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

database.exclude.list

空字符串

一个可选的、以逗号分隔的正则表达式或字面量列表,用于匹配要排除在监控之外的数据库名称。当设置了 database.exclude.list 时,连接器会监控除该属性指定的数据库之外的所有数据库。

要匹配数据库名称,Debezium 会根据 filters.match.mode 属性的值执行以下操作之一:

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

如果在配置中包含了该属性,请不要设置 database.include.list 属性。

collection.include.list

空字符串

一个可选的、以逗号分隔的正则表达式或字面值列表,用于匹配要监控的 MongoDB 集合的完全限定命名空间。默认情况下,连接器会监控除 local 和 admin 数据库中的集合之外的所有集合。当设置了 collection.include.list 时,连接器只监控该属性指定的集合,其他集合则不纳入监控。集合标识符的形式为 databaseName.collectionName。

要匹配命名空间的名称,Debezium 会根据 filters.match.mode 属性的值执行以下操作之一:

  • 将你指定的正则表达式作为锚定正则表达式应用。也就是说,指定的表达式会与命名空间的整个名称字符串进行匹配;它不会匹配名称中的子字符串。
  • 将你指定的字面值与命名空间的整个名称字符串进行比较。

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

collection.exclude.list

空字符串

一个可选的、以逗号分隔的正则表达式或字面值列表,用于匹配要排除在监控之外的 MongoDB 集合的完全限定命名空间。当设置了 collection.exclude.list 时,连接器会监控除该属性指定的集合之外的每个集合。集合标识符的形式为 databaseName.collectionName。

要匹配命名空间的名称,Debezium 会根据 filters.match.mode 属性的值执行以下操作之一:

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

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

capture.mode

change_streams_update_full

指定连接器从 MongoDB 服务器捕获 update 事件变更所使用的方式。将此属性设置为以下值之一:

change_streams

update 事件消息不包含完整文档。消息中不包含表示变更 before 文档状态的字段。

change_streams_update_full

update 事件消息包含完整文档。消息中不包含表示更新前文档状态的 before 字段。事件消息在 after 字段中返回文档的完整状态。设置 capture.mode.full.update.type 以指定连接器如何从数据库获取完整文档。

在某些情况下,当 capture.mode 配置为返回完整文档时,更新事件消息的 updateDescription 和 after 字段可能会报告不一致的值。当对同一文档在短时间内连续应用多次更新时,就可能出现这种差异。连接器只有在收到事件 updateDescription 字段中所描述的更新之后,才会向 MongoDB 数据库请求完整文档。如果在连接器能够从数据库检索该文档之前,后续的更新修改了源文档,那么连接器获取到的将是被该后续更新修改后的文档。

change_streams_update_full_with_pre_image

update 事件消息包含完整文档,并包含表示变更 before 文档状态的字段。设置 capture.mode.full.update.type 以指定连接器如何从数据库获取完整文档。

change_streams_with_pre_image

update 事件不包含完整文档,但包含表示变更 before 文档状态的字段。

capture.start.op.time

无默认值

指定 startAtOperationTime 变更流参数。将该值设置为 BSON 时间戳的 long 型表示。

能够设置 capture.start.op.time` 属性以覆盖基于偏移量的常规行为,并从特定时间戳开始流式传输,这是一项技术预览功能。技术预览功能不受 Red Hat 生产服务水平协议(SLA)的支持,功能可能并不完整。Red Hat 不建议在生产环境中使用这些功能。这些功能使您可以提前访问即将推出的产品功能,从而在开发过程中测试功能并提供反馈。

有关 Red Hat 技术预览功能支持范围的更多信息,请参阅技术预览功能支持范围。

capture.scope

deployment

指定连接器打开的变更流范围。将此属性设置为以下值之一:

deployment

为某个部署(副本集或分片集群)打开变更流游标,以监视除 admin、local 和 config 之外的所有数据库中全部非系统集合的变更。

database

为单个数据库打开变更流游标,以监视该数据库中所有非系统集合的变更。

为了支持 Debezium 信号,如果将 capture.scope 设置为 database,则信号数据集合必须位于由 capture.target 属性指定的数据库中。

collection

为单个集合打开变更流游标,以监视该集合的变更。

此功能目前处于孵化阶段。根据我们收到的反馈,其确切语义、配置选项等可能会发生变化。
将 capture.scope 属性的值设置为 collection 会阻止连接器使用默认的 source 信号通道。由于必须启用 source 通道才能允许连接器处理增量快照信号——即使信号是通过 Kafka、JMX 或 File 通道发送的——因此当 capture-scope 设置为 collection 时,连接器无法执行增量快照。

capture.target

指定连接器监控其变更的数据库。此属性仅在 capture.scope 设置为 database 时适用。

field.exclude.list

空字符串

一个可选的、以逗号分隔的字段完全限定名称列表,用于从变更事件消息值中排除这些字段。字段的完全限定名称格式为 databaseName.collectionName.fieldName.nestedFieldName,其中 databaseName 和 collectionName 可以包含通配符 (*),用于匹配任意字符。

field.renames

空字符串

一个可选的、以逗号分隔的字段完全限定替换列表,用于重命名变更事件消息值中的字段。字段的完全限定替换格式为 databaseName.collectionName.fieldName.nestedFieldName:newNestedFieldName,其中 databaseName 和 collectionName 可以包含通配符 (*),用于匹配任意字符;冒号字符 (:) 用于确定字段的重命名映射。列表中的下一个字段替换会作用于上一个字段替换的结果,因此在重命名同一路径中的多个字段时请注意这一点。

tombstones.on.delete

true

指定 delete 事件之后是否跟随一个墓碑(tombstone)事件。请设置以下选项之一:

true

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

false

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

源记录被删除后,发出墓碑事件(默认行为)允许 Kafka 完全删除与被删除行的键相关的所有事件,前提是对该主题启用了日志压缩。

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 命名。

以下高级配置属性具有良好的默认值,在大多数情况下都能适用,因此很少需要在连接器配置中指定。

Debezium MongoDB 连接器高级配置属性

以下高级属性控制连接器的行为,这些行为很少需要定制。默认值适用于大多数环境,但您可以调整它们以优化性能、调整快照行为或细化过滤。

属性 默认值 描述

capture.mode.full.update.type

lookup

当 capture.mode 设置为检索完整文档时,指定连接器如何查找更新后文档的完整值。当 capture.mode 设置为以下选项之一时,连接器会检索完整文档:

  • change_streams_update_full
  • change_streams_update_full_with_pre-image

要将此选项用于 MongoDB 变更流集合,必须将该集合配置为返回文档的前像和后像。只有在操作发生之前就完成所需配置,操作的前像和后像才可用。

将此属性设置为以下值之一:

lookup

连接器使用单独的查询来获取更新后的完整 MongoDB 文档。

如果查询过程未能检索到文档,就无法在事件负载的 after 状态中填充完整文档。在这种情况下,连接器会发出一条事件消息,其中 after 字段的值为 null。

查询失败可能发生在以下情况:删除操作在文档创建后立即将其删除,或者分片键的更改导致文档被移动到其他位置。当你修改组成分片键的任何属性时,都可能导致分片键发生变化。

post_image

连接器使用 MongoDB 后像来为事件填充完整的 MongoDB 文档。要使用此选项,数据库必须运行 MongoDB 6.0 或更高版本。

max.batch.size

2048

正整数值,指定该连接器每次迭代过程中应处理的每批事件的最大大小。默认值为 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 字节后,写入队列的操作就会被阻塞。

connect.max.attempts

16

一个正整数值,用于指定在发生异常并中止任务之前,连接副本集主节点的最大失败尝试次数。默认值为 16,在 connect.backoff.initial.delay.ms 和 connect.backoff.max.delay.ms 采用默认值的情况下,这意味着在失败之前会尝试略多于 20 分钟。

mongodb.ssl.keystore

无默认值

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

mongodb.ssl.keystore.password

无默认值

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

mongodb.ssl.keystore.type

无默认值

密钥库文件的类型。仅在配置了 mongodb.ssl.keystore 时才指定该类型。

mongodb.ssl.truststore

无默认值

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

mongodb.ssl.truststore.password

无默认值

信任库文件的密码。用于检查信任库的完整性并解锁信任库。仅在配置了 mongodb.ssl.truststore 时才指定该密码。

mongodb.ssl.truststore.type

无默认值

信任库文件的类型。仅在配置了 mongodb.ssl.truststore 时才指定该类型。

source.struct.version

v2

CDC 事件中 source 块的模式版本。Debezium 0.10 对 source 块的结构进行了一些不兼容的更改,以统一所有连接器对外暴露的结构。如果希望连接器生成使用早期版本 source 块结构的事件,请将此选项设置为 v1。请注意,不推荐使用此设置,并计划在未来的 Debezium 版本中移除。

heartbeat.interval.ms

0

指定连接器发送心跳消息的频率。

此属性包含一个以毫秒为单位的间隔值,用于定义连接器向心跳主题发送消息的频率。这可用于监控连接器是否仍在接收数据库的变更事件。当在较长时间内只有未被捕获集合中的记录发生更改时,你也应该利用心跳消息。在这种情况下,连接器会继续从数据库读取 oplog/变更流,但不会向 Kafka 发送任何变更消息,这意味着不会有偏移量更新提交到 Kafka。这将导致 oplog 文件被轮转出去,而连接器却不会察觉,因此在重启时某些事件将不再可用,从而需要重新执行初始快照。

默认设置(0)会阻止连接器发送心跳消息。

skipped.operations

t

一个以逗号分隔的操作类型列表,指定连接器在流式传输期间要跳过的操作类型。你可以配置连接器跳过以下类型的操作:

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

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

snapshot.collection.filter.overrides

无默认值

控制哪些集合项包含在快照中。此属性仅影响快照。请以 databaseName.collectionName 的形式指定以逗号分隔的集合名称列表。

对于你指定的每个集合,还需指定另一个配置属性:snapshot.collection.filter.overrides.databaseName.collectionName。例如,另一个配置属性的名称可能是:snapshot.collection.filter.overrides.customers.orders。将此属性设置为有效的过滤表达式,以仅检索你在快照中所需的项。当连接器执行快照时,它只会检索与过滤表达式匹配的项。

snapshot.delay.ms

无默认值

指定连接器启动后执行快照前需要等待的时间间隔(毫秒)。使用此设置可以在集群中启动多个连接器时避免快照被打断,从而防止连接器发生重新均衡。

streaming.delay.ms

0

指定连接器完成快照后延迟启动流式处理过程的时间间隔(毫秒)。设置延迟间隔有助于防止在快照刚刚完成、但流式处理尚未开始之际发生故障时,连接器重新执行快照。请将延迟值设置为高于为 Kafka Connect 工作进程配置的 offset.flush.interval.ms 属性值。

snapshot.fetch.size

0

指定执行快照时,每个集合单次最多读取的文档数量。连接器将按照该大小分多批读取集合内容。

默认设置(0)表示由服务端确定抓取大小。

snapshot.include.collection.list

collection.include.list 中指定的所有集合

一个可选的、以逗号分隔的正则表达式列表,用于匹配你希望包含在快照中的 schema 的全限定名称(<databaseName>.<collectionName>)。所列出的项必须已在连接器的 collection.include.list 属性中指定。只有当连接器的 snapshot.mode 属性设置为 no_data 以外的值时,此属性才会生效。

此属性不会影响增量快照的行为。

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

snapshot.max.threads

1

正整数值,指定对副本集中各集合执行初始同步时所使用的最大线程数。默认为 1。

snapshot.mode

initial

指定连接器启动时执行快照的条件。将该属性设置为以下值之一:

always

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

initial

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

initial_only

仅当该逻辑服务器名称尚未记录任何偏移量时,连接器才执行数据库快照。快照完成后,连接器停止运行,不会转入流式传输后续数据库更改的事件记录。

no_data

连接器执行的快照会捕获所有相关表的结构,但不会创建 READ 事件来表示连接器启动时的数据集。

when_needed

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

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

configuration_based

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

custom

custom 快照模式允许您注入自己的 io.debezium.spi.snapshot.Snapshotter 接口实现。将 snapshot.mode.custom.name 配置属性设置为您的实现中 name() 方法所提供的名称。

有关更多信息,请参阅自定义快照器 SPI。

snapshot.mode.configuration.based.snapshot.data

false

如果 snapshot.mode 设置为 configuration_based,请设置此属性以指定连接器在执行快照时是否包含表数据。

snapshot.mode.configuration.based.snapshot.schema

false

如果 snapshot.mode 设置为 configuration_based,请设置此属性以指定连接器在执行快照时是否包含表结构。

snapshot.mode.configuration.based.start.stream

false

如果 snapshot.mode 设置为 configuration_based,请设置此属性以指定连接器在快照完成后是否开始流式传输更改事件。

snapshot.mode.configuration.based.snapshot.on.schema.error

false

如果 snapshot.mode 设置为 configuration_based,可通过该属性指定:当 schema 历史主题不可用时,连接器是否在快照中包含表 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() 方法提供。连接器重启后,Debezium 会调用指定的自定义实现,以确定是否执行快照。有关更多信息,请参阅自定义快照器 SPI。

provide.transaction.metadata

false

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

更多详情请参阅事务元数据。

retriable.restart.connector.wait.ms

10000(10 秒)

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

mongodb.poll.interval.ms

30000

连接器轮询新增、移除或变更的副本集的时间间隔。

mongodb.connect.timeout.ms

10000(10 秒)

驱动程序在放弃新的连接尝试之前等待的毫秒数。

mongodb.heartbeat.frequency.ms

10000(10 秒)

集群监控器尝试连接每个服务器的频率。

mongodb.socket.timeout.ms

0

发送/接收在套接字上发生超时之前允许的毫秒数。设为 0 可禁用此行为。

mongodb.server.selection.timeout.ms

30000(30 秒)

驱动程序在超时并抛出错误之前等待选择服务器的毫秒数。

cursor.pipeline

无默认值

在流式传输变更时,此设置会作为标准 MongoDB 聚合流管道的一部分对变更流事件应用处理。管道是由一系列指令组成的 MongoDB 聚合管道,用于指示数据库过滤或转换数据。借此可以自定义连接器消费的数据。此属性的值必须是以 JSON 格式表示的、允许使用的聚合管道阶段数组。请注意,该管道会追加在连接器内部管道(例如过滤操作类型、数据库名称、集合名称等)之后。

cursor.pipeline.order

internal_first

用于构建实际 MongoDB 聚合流管道的顺序。将该属性设置为以下值之一:

internal_first

先应用由连接器定义的内部阶段。这意味着只有本应由连接器捕获的事件才会被送入用户定义的阶段(通过设置 cursor.pipeline 配置)。

user_first

先应用由 cursor.pipeline 属性定义的阶段。在此模式下,所有事件(包括未被连接器捕获的事件)都会被送入用户定义的管道阶段。如果 cursor.pipeline 的值包含复杂的操作,此模式可能会对性能产生负面影响。

user_only

由 cursor.pipeline 属性定义的阶段将替换连接器定义的内部阶段。此模式仅供专家用户使用,因为所有事件仅由用户定义的管道阶段进行处理。此模式可能会对连接器的性能和整体功能产生负面影响!

cursor.oversize.handling.mode

fail

用于处理超出指定 BSON 大小的文档的变更事件的策略。将该属性设置为以下值之一:

fail

如果变更事件的总大小超出最大 BSON 大小,则连接器失败。

skip

任何超出最大大小(由 cursor.oversize.skip.threshold 属性指定)的文档的变更事件都将被忽略。

split

超过最大 BSON 大小的变更事件将使用 $changeStreamSplitLargeEvent 聚合进行拆分。此选项需要 MongoDB 6.0.9 或更高版本。

cursor.oversize.skip.threshold

0

处理变更事件时,所存储文档允许的字节最大大小。这包括数据库操作前后的大小,更具体地说,它限制了 MongoDB 变更事件中 fullDocument 和 fullDocumentBeforeChange 字段的大小。

cursor.max.await.time.ms

0

指定 oplog/变更流游标在抛出执行超时异常之前,等待服务器产生结果的最大毫秒数。值为 0 表示使用服务器/驱动程序默认的等待超时时间。

signal.data.collection

无默认值

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

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

<databaseName>.<collectionName>

signal.enabled.channels

source

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

notification.enabled.channels

无默认值

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

incremental.snapshot.chunk.size

1024

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

incremental.snapshot.watermarking.strategy

insert_insert

指定连接器在增量快照期间使用的水印机制,用于对可能被增量快照捕获、并在流式传输恢复后再次捕获的事件进行去重。

您可以指定以下选项之一:

insert_insert

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

insert_delete

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

topic.naming.strategy

io.debezium.schema.DefaultTopicNamingStrategy

应使用的 TopicNamingStrategy 类的名称,用于确定数据变更、模式变更、事务、心跳事件等的 Topic 名称,默认为 DefaultTopicNamingStrategy。

topic.delimiter

.

指定 Topic 名称的分隔符,默认为 .。

topic.cache.size

10000

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

topic.heartbeat.prefix

__debezium-heartbeat

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

topic.heartbeat.prefix.topic.prefix

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

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

topic.heartbeat.name

空

指定连接器发送心跳消息的主题的显式完整名称,从而覆盖由 topic.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。

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

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

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 连接器的名称。

skip.messages.without.change

false

MongoDB 连接器不支持跳过未发生变化的事件。虽然该属性继承自共享连接器框架,但 MongoDB 连接器会忽略它。

statistics.metrics.enabled

true

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

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

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

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

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

下表描述了 Kafka signal 属性。

表 6. 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

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

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

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

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

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

监控

除 Kafka 和 Kafka Connect 内置的 JMX 指标支持外,Debezium MongoDB 连接器还提供两类指标。

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

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 名称

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

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

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

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

快照指标

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

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

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

Debezium MongoDB 连接器还提供以下自定义快照指标:

属性类型描述
NumberOfDisconnectslong数据库断开连接的次数。

流式传输指标

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

下表列出了可用于监控 Debezium 流式传输操作的 JMX 指标,包括按类型划分的事件数量、相对于源的延迟、队列容量以及连接状态。

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

| [LastTransactionId](/book/reader/debezium/mongodb#connectors-strm-metric-lasttransaction

Debezium MongoDB 连接器还提供以下自定义流式处理指标:

属性类型描述
NumberOfDisconnectslong数据库断开连接的次数。
NumberOfPrimaryElectionslong主节点选举的次数。

MongoDB 连接器常见问题

Debezium 是一个分布式系统,用于捕获多个上游数据库中的所有更改,绝不会遗漏或丢失任何事件。当系统正常运行并得到妥善管理时,Debezium 能够为每个更改事件提供*精确一次(exactly once)*的交付保证。

如果发生故障,系统不会丢失任何事件。但在从故障中恢复的过程中,可能会重复投递某些更改事件。在这种情况下,Debezium 与 Kafka 一样,对更改事件提供*至少一次(at least once)*的交付保证。

本节余下部分将介绍 Debezium 如何处理各类故障和问题。

配置与启动错误

在以下情况下,连接器在尝试启动时会失败,在日志中报告错误或异常,并停止运行:

  • 连接器的配置无效。
  • 连接器无法使用指定的连接参数成功连接到 MongoDB。

发生失败后,连接器会使用指数退避策略尝试重新连接。你可以配置重新连接的最大尝试次数。

在这些情况下,错误信息会提供更多关于问题的细节,可能还会给出建议的解决方法。当配置被修正或 MongoDB 的问题得到解决后,可以重新启动连接器。

MongoDB 变为不可用

如果连接器正在运行,而它读取的副本集成员变为不可用,连接器会持续尝试重新连接到另一个合适的成员,并使用指数退避策略以避免给网络或服务器造成过大的压力。如果连接器在配置的连接尝试次数内无法连接到任何副本集成员,连接器将失败。

重新连接的尝试由三个属性控制:

  • connect.backoff.initial.delay.ms - 首次尝试重新连接之前的延迟时间,默认为 1 秒(1000 毫秒)。
  • connect.backoff.max.delay.ms - 尝试重新连接之前的最长延迟时间,默认为 120 秒(120,000 毫秒)。
  • connect.max.attempts - 产生错误之前的最大尝试次数,默认为 16。

每次的延迟都是前一次延迟的两倍,直到达到最大延迟为止。根据默认值,下表显示了每次连接尝试失败时的延迟时间,以及失败前累积的总时间。

重连尝试次数尝试前的延迟(秒)尝试前的总延迟(分:秒)
1100:01
2200:03
3400:07
4800:15
51600:31
63201:03
76402:07
812004:07
912006:07
1012008:07
1112010:07
1212012:07
1312014:07
1412016:07
1512018:07
1612020:07

连接器无法启动 - InvalidResumeToken 或 ChangeStreamHistoryLost

长时间停止的连接器无法启动,并报告以下异常:

Command failed with error 286 (ChangeStreamHistoryLost): 'PlanExecutor error during aggregation :: caused by :: Resume of change stream was not possible, as the resume point may no longer be in the oplog

上述异常表明,与连接器的续传令牌(resume token)对应的条目已不再存在于 oplog 中。由于 oplog 中不再包含连接器已处理的最后一个偏移量,连接器无法恢复流式传输。

你可以使用以下任一选项从该故障中恢复:

  • 删除失败的连接器,并使用相同配置但不同连接器名称创建一个新的连接器。
  • 暂停连接器,然后移除偏移量,或者更改偏移量主题。

为帮助防止与缺失续传令牌相关的故障,请优化 oplog 配置。

Kafka Connect 进程正常停止

如果 Kafka Connect 以分布式模式运行,并且某个 Kafka Connect 进程被正常停止,那么在该进程关闭之前,Kafka Connect 会将该进程上的所有连接器任务迁移到同一组中的另一个 Kafka Connect 进程,新的连接器任务将从前一个任务中断的位置精确接续。在连接器任务被正常停止并在新进程上重新启动期间,处理会出现短暂延迟。

如果该组中只有一个进程且该进程被正常停止,那么 Kafka Connect 会停止连接器,并记录每个副本集的最后一个偏移量。重新启动后,副本集任务将从中断处精确继续。

Kafka Connect 进程崩溃

如果 Kafka Connector 进程意外停止,那么它正在运行的任何连接器任务都会在未记录其最近处理的偏移量的情况下终止。当 Kafka Connect 以分布式模式运行时,它会在其他进程上重新启动这些连接器任务。但是,MongoDB 连接器将从前序进程已记录的最后一个偏移量恢复,这意味着新接替的任务可能会重新生成崩溃前刚刚处理过的部分变更事件。重复事件的数量取决于偏移量刷写周期以及崩溃前数据变更的量。

由于在故障恢复期间可能会出现重复事件,使用者应始终预期某些事件可能重复。Debezium 的变更是幂等的,因此同一序列事件始终会产生相同的状态。

Debezium 还会在每条变更事件消息中包含与来源相关的信息,其中包括 MongoDB 事件的唯一事务标识符(h)以及时间戳(sec 和 ord)。使用者可以跟踪这些值,从而判断是否已经处理过某个特定事件。

Kafka 变得不可用

当连接器生成变更事件时,Kafka Connect 框架会使用 Kafka 生产者 API 将这些事件记录到 Kafka 中。Kafka Connect 还会按照你在 Kafka Connect 工作进程配置中指定的频率,周期性地记录这些变更事件中出现的最新偏移量。如果 Kafka 代理变得不可用,运行连接器的 Kafka Connect 工作进程将只是不断尝试重新连接到 Kafka 代理。换句话说,连接器任务将暂停,直到可以重新建立连接为止,届时连接器会从中断的地方准确恢复。

当 snapshot.mode 设置为 initial 时,连接器在长时间停止后会失败

如果连接器被正常停止,用户可能会继续对副本集成员执行操作。在连接器离线期间发生的变更会继续被记录到 MongoDB 的 oplog 中。在大多数情况下,连接器重新启动后,会读取 oplog 中的偏移量值,以确定它为每个副本集流式传输的最后一个操作,然后从该位置恢复流式传输变更。重新启动后,在连接器停止期间发生的数据库操作会照常发送到 Kafka,经过一段时间后,连接器会追上数据库的进度。连接器追上进度所需的时间取决于 Kafka 的能力和性能,以及数据库中发生的变更量。

然而,如果连接器停止的时间足够长,MongoDB 可能在连接器不活跃的期间清除了 oplog,从而导致连接器最后位置的信息丢失。连接器重新启动后无法恢复流式传输,因为 oplog 中不再包含标记连接器所处理的最后一个操作的先前偏移量值。连接器也无法执行快照,而通常当 snapshot.mode 属性设置为 initial 且不存在偏移量值时,它是会执行快照的。在这种情况下,出现了不匹配的情况,因为 oplog 中不包含先前的偏移量值,但连接器内部的 Kafka 偏移量主题中却存在该偏移量值。最终导致错误,连接器失败。

要从该故障中恢复,请删除失败的连接器,然后使用相同的配置但不同的连接器名称创建一个新的连接器。启动新连接器时,它会执行快照以摄取数据库的状态,然后恢复流式传输。

MongoDB 丢失写入

在某些故障情况下,MongoDB 可能会丢失提交,导致 MongoDB 连接器无法捕获这些丢失的变更。例如,如果主节点在应用某项变更并将其记录到 oplog 之后突然崩溃,那么在从节点读取其内容之前,该 oplog 可能已不可用。结果,被选举为新主节点的从节点可能缺少其 oplog 中最近的变更。

目前,MongoDB 中没有办法防止这种副作用。

评论

登录后参与评论

正在加载评论…