MongoDB 新文档状态提取
MongoDB 新文档状态提取
Debezium MongoDB 连接器会发出数据变更消息,用以表示 MongoDB 集合中发生的每一次操作。这些事件消息的结构复杂,忠实地还原了原始数据库事件的细节。然而,有些下游消费者可能无法处理原始格式的消息。例如,为了表示数据集合中的嵌套文档,连接器会发出包含嵌套字段格式的事件消息。为了支持 sink 连接器,或其他无法处理原始消息层次结构格式的消费者,你可以使用 Debezium MongoDB 事件展平(ExtractNewDocumentState)单消息转换(SMT)。该 SMT 会简化原始消息的结构,也可以通过其他方式修改消息,使数据更易于处理。
事件展平转换是一个 Kafka Connect SMT。
| 本章介绍的信息仅适用于 Debezium MongoDB 连接器的事件展平单消息转换(SMT)。有关在关系型数据库中使用的等效 SMT 的信息,请参阅新记录状态提取 SMT 的文档。 |
|---|
变更事件结构
Debezium MongoDB 连接器生成的变更事件结构复杂。每个事件消息包含以下部分:
源元数据
包括但不限于以下字段:
- 在集合中更改数据的操作类型(create/insert、update 或 delete)。
- 发生变更的数据库和集合的名称。
- 标识变更发生时间的时间戳。
- 可选的事务信息。
文档数据
before 数据
当 Debezium 连接器的 capture.mode 设置为以下值之一时,运行 MongoDB 6.0 及更高版本的环境中会存在此字段:
change_streams_with_pre_image。change_streams_update_full_with_pre_image。有关更多信息,请参阅 MongoDB 预映像支持
after 数据
表示当前操作之后文档中所存在值的 JSON 字符串。事件消息中是否包含 after 字段,取决于事件类型和连接器配置。对于 MongoDB insert 操作,无论 capture.mode 设置如何,生成的 create 事件始终包含 after 字段。对于 update 事件,只有当 capture.mode 设置为以下值之一时,才包含 after 字段:
change_streams_update_fullchange_streams_update_full_with_pre_image。变更事件消息中的
after值并不一定代表事件发生后文档的状态。该值不是动态计算的;相反,连接器在捕获变更事件后,会查询集合以获取该文档的当前值。例如,设想这样一种情况:多个操作
a、b和c接连快速地修改同一个文档。当连接器处理变更a时,它会查询集合以获取完整的文档。与此同时,变更b和c已经发生。当连接器收到针对变更a的完整文档查询响应时,它可能获得的是基于后续变更b或c的文档版本。有关详细信息,请参阅capture.mode属性的文档。
以下片段展示了 MongoDB insert 操作后连接器发出的 create 变更事件的基本结构:
{
"op": "c",
"after": "{\"field1\":\"newvalue1\",\"field2\":\"newvalue1\"}",
"source": { ... }
}前述示例中 after 字段的复杂格式提供了源数据库中所发生变更的详细信息。但是,某些消费者无法处理包含嵌套值的消息。若要将原始消息的复杂嵌套字段转换为更简单、通用兼容性更强的结构,请使用适用于 MongoDB 的事件展平 SMT。该 SMT 会展平消息中嵌套字段的结构,如下例所示:
{
"field1" : "newvalue1",
"field2" : "newvalue2"
}有关 Debezium MongoDB 连接器生成的默认消息结构的更多信息,请参阅连接器文档。
行为
MongoDB 的事件展平 SMT 会从 Debezium MongoDB 连接器发出的 create 或 update 变更事件消息中提取 after 字段。SMT 处理完原始变更事件消息后,会生成一个简化的版本,其中仅包含 after 字段的内容。
根据你的使用场景,你可以将 ExtractNewDocumentState SMT 应用于 Debezium MongoDB 连接器,也可以将其应用于消费 Debezium 连接器所生成消息的接收器连接器。如果你将 SMT 应用于 Debezium MongoDB 连接器,该 SMT 会在连接器发出的消息被发送到 Apache Kafka 之前对其进行修改。为了确保 Kafka 以原始格式保留完整的 Debezium 变更事件消息,请将 SMT 应用于接收器连接器。
当你使用事件展平 SMT 处理由 MongoDB 连接器发出的消息时,该 SMT 会将原始消息中记录的结构转换为类型正确的 Kafka Connect 记录,从而能够被典型的接收器连接器消费。例如,SMT 会将原始消息中表示 after 信息的 JSON 字符串转换为任何消费者都能处理的 schema 结构。
你还可以选择配置 MongoDB 的事件展平 SMT,使其在处理过程中以其他方式修改消息。有关更多信息,请参阅配置章节。
配置
为消费 Debezium MongoDB 连接器所发出消息的接收器连接器配置 MongoDB 的事件展平(ExtractNewDocumentState)SMT。
基本配置
要获得 SMT 的默认行为,请在不指定任何选项的情况下将该 SMT 添加到接收器连接器的配置中,如下例所示:
transforms=unwrap,...
transforms.unwrap.type=io.debezium.connector.mongodb.transforms.ExtractNewDocumentState与任何 Kafka Connect 连接器配置一样,你可以将 transforms= 设置为多个以逗号分隔的 SMT 别名。Kafka Connect 会按照你列出的顺序应用这些转换。
对于使用 MongoDB 事件扁平化 SMT 的连接器,你可以设置多个选项。以下示例展示了为连接器设置 delete.tombstone.handling.mode 和 add.headers 选项的配置:
transforms=unwrap,...
transforms.unwrap.type=io.debezium.connector.mongodb.transforms.ExtractNewDocumentState
transforms.unwrap.delete.tombstone.handling.mode=drop
transforms.unwrap.add.headers=op有关上述示例中配置选项的更多信息,请参阅配置主题。
自定义配置
连接器可能会发出多种类型的事件消息(例如心跳消息、墓碑消息或事务的元数据消息)。要将转换仅应用于部分事件,你可以定义一条 SMT 谓词语句,使其选择性地对特定事件应用转换。
数组编码
默认情况下,事件展平 SMT 会将 MongoDB 数组转换为与 Apache Kafka Connect 或 Apache Avro 架构兼容的数组。虽然 MongoDB 数组可以包含多种类型的元素,但 Kafka 数组中的所有元素必须具有相同的类型。
为确保 SMT 以符合你的环境需求的方式编码数组,你可以指定 array.encoding 配置选项。以下示例展示了设置数组编码的配置:
transforms=unwrap,...
transforms.unwrap.type=io.debezium.connector.mongodb.transforms.ExtractNewDocumentState
transforms.unwrap.array.encoding=<array|document>根据配置,SMT 使用以下其中一种编码方法处理源消息中的每个数组实例:
数组编码
如果 array.encoding 设置为 array(默认值),SMT 使用 array 数据类型来编码原始消息中的数组。为确保正确处理,数组实例中的所有元素必须为相同类型。此选项具有限制性,但它使下游客户端能够轻松处理数组。
文档编码
如果 array.encoding 设置为 document,SMT 会将源中的每个数组转换为 struct 组成的 struct,其方式类似于 BSON 序列化。主 struct 包含名为 _0、_1、_2 等的字段,其中每个字段名表示原始数组中元素的索引。SMT 使用从源数组中对应元素检索到的值填充这些索引字段。索引名称以下划线为前缀,因为 Avro 编码禁止以数字字符开头的字段名。
以下示例展示了 Debezium MongoDB 连接器如何表示包含多种数据类型的数组的数据库文档:
示例 1. 示例:包含多种数据类型的数组的文档编码
{
"_id": 1,
"a1": [
{
"a": 1,
"b": "none"
},
{
"a": "c",
"d": "something"
}
]
}如果将 array.encoding 设置为 document,SMT 会将前述文档转换为以下格式:
{
"_id": 1,
"a1": {
"_0": {
"a": 1,
"b": "none"
},
"_1": {
"a": "c",
"d": "something"
}
}
}document 编码选项使 SMT 能够处理由异构元素组成的任意数组。但是,在使用此选项之前,请务必验证接收端连接器以及其他下游消费者能够处理包含多种数据类型的数组。
嵌套结构扁平化
当数据库操作涉及嵌入式文档时,Debezium MongoDB 连接器发出的 Kafka 事件记录其结构反映了原始文档的层次结构。也就是说,事件消息以一组嵌套的字段结构来表示嵌套文档。在下游连接器无法处理包含嵌套结构的消息的环境中,你可以配置事件扁平化 SMT,将消息中的层次结构扁平化。扁平的消息结构更适合类似表的存储。
要配置 SMT 扁平化嵌套结构,请将 flatten.struct 配置选项设置为 true。在转换后的消息中,字段名的构造与文档来源保持一致。SMT 通过将父文档字段名与嵌套文档字段名拼接起来,为每个扁平化字段重新命名。由 flatten.struct.delimiter 选项定义的分隔符用于分隔名称的各个组成部分。struct.delimiter 的默认值是下划线字符(_)。
以下示例展示了用于指定 SMT 是否扁平化嵌套结构的配置:
transforms=unwrap,...
transforms.unwrap.type=io.debezium.connector.mongodb.transforms.ExtractNewDocumentState
transforms.unwrap.flatten.struct=<true|false>
transforms.unwrap.flatten.struct.delimiter=<string>以下示例展示了 MongoDB 连接器发出的事件消息。该消息包含一个文档 a 的字段,其中又包含两个嵌套文档 b 和 c 的字段:
{
"_id": 1,
"a": {
"b": 1,
"c": "none"
},
"d": 100
}下面的示例消息展示了 MongoDB 的 SMT(单消息转换)展平前述消息中的嵌套结构后的输出结果:
{
"_id": 1,
"a_b": 1,
"a_c": "none",
"d": 100
}在生成的消息中,原本嵌套在消息中的 b 和 c 字段会被展平并重命名。重命名后的字段由父文档 a 的名称与嵌套文档的名称拼接而成:a_b 和 a_c。新字段名的各个组成部分由下划线字符分隔,具体分隔符由 struct.delimiter 配置属性的设置决定。
MongoDB $unset 处理
在 MongoDB 中,$unset 操作符和 $rename 操作符都会从文档中移除字段。由于 MongoDB 集合是无模式的,在更新操作从文档中移除字段后,无法从更新后的文档中推断出被移除字段的名称。为了支持那些可能需要被移除字段信息的 sink 连接器或其他消费者,Debezium 会发出包含 removedFields 元素的更新消息,其中列出了被删除字段的名称。
以下示例展示了某次操作产生的更新消息的一部分,该操作导致字段 a 被移除:
"payload": {
"op": "u",
"ts_ms": "...",
"ts_us" : "...",
"ts_ns" : "...",
"before": "{ ... }",
"after": "{ ... }",
"updateDescription": {
"removedFields": ["a"],
"updatedFields": null,
"truncatedArrays": null
}
}在上面的示例中,before 和 after 分别表示源文档在更新前和更新后的状态。连接器发出的事件消息中是否包含这些字段,取决于该连接器的 capture.mode 设置,具体如下所述:
before 字段
提供变更发生前文档的状态。仅当 capture.mode 设置为以下值之一时,该字段才会存在:
change_streams_with_pre_imagechange_streams_update_full_with_pre_image
after 字段
提供变更后文档的完整状态。仅当 capture.mode 设置为以下值之一时,该字段才会存在:
change_streams_update_fullchange_streams_update_full_with_pre_image
假设连接器已配置为捕获完整文档,当 ExtractNewDocumentState SMT 收到 $unset 事件的 update 消息时,该 SMT 会重新编码消息,将被移除的字段表示为 null 值,如以下示例所示:
{
"id": 1,
"a": null
}对于未配置为捕获完整文档的连接器,当 SMT 收到 $unset 操作的更新事件时,会生成以下输出消息:
{
"a": null
}确定原始操作
SMT 展平事件消息后,得到的消息不再表明生成该事件的操作类型是 create、update 还是初始快照的 read。通常,你可以通过配置连接器,使其暴露与删除操作相关的墓碑(tombstone)或重写(rewrite)事件的信息,从而识别 delete 操作。有关配置连接器以在事件消息中暴露墓碑和重写信息的更多内容,请参阅 delete.tombstone.handling.mode 属性。
要在事件消息中报告数据库操作的类型,SMT 可以将 op 字段添加到以下元素之一:
- 事件消息体。
- 消息头。
例如,要添加一个显示原始操作类型的头部属性,请添加转换,然后在连接器配置中添加 add.headers 属性,如下例所示:
transforms=unwrap,...
transforms.unwrap.type=io.debezium.connector.mongodb.transforms.ExtractNewDocumentState
transforms.unwrap.add.headers=op根据上述配置,SMT 通过向消息添加 op 头并为其赋一个字符串值来报告事件类型。所赋的字符串值基于原始 MongoDB 变更事件消息 中的 op 字段值。
添加元数据字段
MongoDB 的事件展平 SMT 可以将原始变更事件消息中的元数据字段添加到简化后的消息中。添加的元数据字段带有双下划线("__")前缀。为事件记录添加元数据后,就可以包含诸如变更事件所在的集合名称等内容,或者包含连接器特有的字段,例如副本集名称。目前,SMT 只能从以下变更事件子结构中添加字段:source、transaction 和 updateDescription。
有关 MongoDB 变更事件结构的更多信息,请参阅 MongoDB 连接器文档。
例如,你可以指定以下配置,将副本集名称(rs)和变更事件的集合名称添加到最终的展平事件记录中:
transforms=unwrap,...
transforms.unwrap.type=io.debezium.connector.mongodb.transforms.ExtractNewDocumentState
transforms.unwrap.add.fields=rs,collection上述配置会将以下内容添加到展平后的记录中:
{ "__rs" : "rs0", "__collection" : "my-collection", ... }如果你希望 SMT 为 delete 事件添加元数据字段,请将 delete.tombstone.handling.mode 选项的值设置为 rewrite。
选择性应用转换的选项
除了 Debezium 连接器在数据库发生变化时发出的变更事件消息之外,连接器还会发出其他类型的消息,包括心跳消息,以及有关架构变更和事务的元数据消息。由于这些其他消息的结构与 SMT 设计要处理的变更事件消息的结构不同,因此最好将连接器配置为选择性地应用 SMT,使其仅处理预期的数据变更消息。
有关如何选择性应用 SMT 的更多信息,请参阅为转换配置 SMT 谓词。
配置选项
下表描述了 MongoDB 事件展平 SMT 的配置选项。
属性 默认值 说明
array
指定 SMT 在编码从原始事件消息中读取的数组时所使用的格式。可设置为以下选项之一:
array
SMT 使用 array 数据类型将 MongoDB 数组编码为与 Apache Kafka Connect 或 Apache Avro 架构兼容的格式。如果设置此选项,请验证每个数组实例中的元素类型是否相同。虽然 MongoDB 允许数组包含多种数据类型,但某些下游客户端无法处理这样的数组。
document
SMT 将每个 MongoDB 数组转换为由 struct 组成的 struct,其方式与 BSON 序列化类似。主 struct 包含名为 _0、_1、_2 等的字段。为了符合 Avro 命名规范,SMT 会为每个索引字段的数字名称添加下划线前缀。每个数字字段名表示该元素在原数组中的索引。SMT 会从源文档中取出指定数组元素的值,并填充到对应的索引字段中。
有关 array.coding 选项的更多信息,请参阅 MongoDB 事件消息的数组编码选项。
false
Mongodb 新文档状态提取
该 SMT 通过连接原始事件消息中嵌套属性的名称来扁平化结构(struct),属性名之间以可配置的分隔符分隔,从而形成简单的字段名。
_
当 flatten.struct 设置为 true 时,指定该转换在连接输入记录的字段名以生成输出记录的字段名时,插入其间的分隔符。
delete.tombstone.handling.mode
tombstone
Debezium 会为每个 DELETE 操作生成一条变更事件记录。此设置决定 MongoDB 事件扁平化 SMT 在流中如何处理 DELETE 事件。请设置为以下选项之一:
drop
SMT 会从流中同时移除 DELETE 事件和 `TOMBSTONE` 记录。
tombstone(默认)
SMT 会在流中保留 TOMBSTONE 记录。TOMBSTONE 记录仅包含以下值:"value": "null"。
rewrite
SMT 会在流中保留变更事件记录,并进行以下更改:
向记录添加一个
value字段,其中包含来自原始记录before字段的键/值对。向记录的
value中添加__deleted: true。移除
TOMBSTONE记录。此设置提供了另一种表明记录已被删除的方式。
rewrite-with-tombstone
SMT 的行为与选择 rewrite 选项时相同,但除此之外还会保留 TOMBSTONE 记录。
delete.tombstone.handling.mode.rewrite-with-id
false
设置为 true 且 delete.tombstone.handling.mode 为 rewrite 时,会从键中复制 id 字段,并在删除事件的负载中将其作为 _id 包含。
__(双下划线)
指定用于为消息头添加前缀的可选字符串。
无默认值
指定一个以逗号分隔、不含空格的字段列表,这些字段由该 SMT 添加到简化后的 Kafka 消息头中。如果原消息中存在同名字段,可以通过同时提供结构体名称和字段名称来指定要修改的字段,例如 source.ts_ms。
此外,你还可以按以下格式向列表中添加条目,从而覆盖字段的原始名称,并为其指定一个区分大小写的新名称:
<field_name>:<new_field_name>。
例如:
version:VERSION, connector:CONNECTOR, source.ts_ms:EVENT_TIMESTAMP当 SMT 向简化消息的头部添加元数据字段时,会为每个元数据字段名加上双下划线前缀。对于结构体(struct)形式的指定方式,SMT 还会在结构体名称与字段名称之间插入一个下划线。
如果你指定的字段不存在于原始变更事件消息的头部中,SMT 不会将该字段添加到头部。
__(双下划线)
指定一个可选字符串,作为字段名的前缀。
无默认值
指定一个以逗号分隔(不包含空格)的元数据字段列表,这些字段将由 SMT 添加到简化 Kafka 消息的 value 元素中。如果原始消息中存在重复的字段名,你可以通过同时提供结构体名称和字段名称来指明要修改的字段,例如 source.ts_ms。
你还可以选择性地覆盖字段的原始名称,并按以下格式向列表中添加一个条目,为其指定一个区分大小写的新名称:
<field_name>:<new_field_name>。
例如:
version:VERSION, connector:CONNECTOR, source.ts_ms:EVENT_TIMESTAMP当 SMT 向简化消息的 value 元素添加元数据字段时,它会在每个元数据字段名前加上双下划线。对于结构体(struct)形式的规范,SMT 还会在结构体名与字段名之间插入一个下划线。
如果你指定的字段在原始变更事件消息中不存在,SMT 仍会将该字段添加到修改后消息的 value 元素中。
none
指定 SMT 如何调整字段名和 schema 名,以兼容不同的消息转换器。请设置以下选项之一:
none(默认)
不执行任何调整。
avro
SMT 会将 Avro 字段名中无效的字符替换为下划线(_)。它还会逐一调整输出 schema 名中以点号分隔的各段,将开头的无效字符替换为下划线,以防止当集合名以数字开头时引发 Avro 的 SchemaParseException 错误。
avro_unicode
与 avro 类似,但会将无效字符替换为其 Unicode 等价形式(例如,开头的 1 会被替换为 _u0031)。
已知限制
由于 MongoDB 是无模式数据库,而接收端连接器要求特定的消息格式,因此在配置 MongoDB 事件展平转换时,需要注意以下约束。
- 由于 MongoDB 是无模式数据库,为了确保使用 Debezium 将变更流式传输到基于模式的关系型数据库时列定义一致,同一集合中名称相同的字段必须存储相同类型的数据。
- 请配置 SMT,使生成的消息格式与接收端连接器兼容。如果接收端连接器要求「扁平」的消息结构,但收到的消息将源 MongoDB 文档中的数组编码为结构体的结构体,则该接收端连接器将无法处理此消息。
评论
登录后参与评论
KnowForge