新建记录状态提取
新记录状态提取
Debezium 连接器会发出数据变更消息,以表示其从源数据库中捕获的每一次操作。连接器发送到 Apache Kafka 的消息具有复杂的结构,能够忠实呈现原始数据库事件的细节。
尽管这种复杂的消息格式能够准确详述系统中发生的变化信息,但该格式可能并不适用于某些下游消费者。Sink 连接器或 Kafka 生态系统的其他部分可能需要这样格式的消息:其字段名和值以简化的扁平结构呈现。
为了简化 Debezium 连接器产生的事件记录格式,你可以使用 Debezium 事件扁平化单消息转换(SMT)。配置该转换以支持那些要求 Kafka 记录格式比连接器默认输出格式更简单的消费者。根据你的具体用例,你可以将该 SMT 应用于 Debezium 连接器,或应用于消费 Debezium 连接器所产生消息的 sink 连接器。若要让 Apache Kafka 保留 Debezium 变更事件消息的原始格式,请为 sink 连接器配置该 SMT。
事件扁平化转换是一种 Kafka Connect SMT。
| 本章介绍的是适用于 Debezium 基于 SQL 的数据库连接器的事件扁平化单消息转换(SMT)。有关适用于 Debezium MongoDB 连接器的等效 SMT 的信息,请参阅 MongoDB 新文档状态提取。 |
|---|
变更事件结构
Debezium 生成的数据变更事件具有复杂的结构。每个事件由三部分组成:
元数据,包括但不限于:
- 导致数据变更的操作类型。
- 源信息,例如发生变更的数据库和表的名称。
- 标识变更发生时间的时间戳。
- 可选的事务信息。
变更前的行数据
变更后的行数据
以下示例展示了 UPDATE 变更事件消息结构的一部分:
{
"op": "u",
"source": {
...
},
"ts_ms" : "...",
"ts_us" : "...",
"ts_ns" : "...",
"before" : {
"field1" : "oldvalue1",
"field2" : "oldvalue2"
},
"after" : {
"field1" : "newvalue1",
"field2" : "newvalue2"
}
}有关连接器变更事件结构的更多信息,请参阅连接器文档。
在事件展平 SMT 处理上例中的消息之后,它会简化消息格式,得到如下示例中的消息:
{
"field1" : "newvalue1",
"field2" : "newvalue2"
}行为
事件展平 SMT 会从 Kafka 记录中的 Debezium 变更事件里提取 after 字段。该 SMT 用变更事件的 after 字段替换原始变更事件,从而生成一个简单的 Kafka 记录。
你可以为 Debezium 连接器配置事件展平 SMT,也可以为消费 Debezium 连接器所发出消息的接收器(sink)连接器进行配置。为接收器连接器配置事件展平的好处在于,存储在 Apache Kafka 中的记录只包含 Debezium 变更事件的 after 字段。究竟将该 SMT 应用于源连接器还是接收器连接器,取决于你的具体用例。
你可以将该转换配置为执行以下任一操作:
- 将变更事件中的元数据添加到简化的 Kafka 记录中。默认行为是 SMT 不添加任何元数据。
- 在流中保留包含
DELETE操作变更事件的 Kafka 记录。默认行为是 SMT 会丢弃DELETE操作变更事件对应的 Kafka 记录,因为大多数消费者目前还无法处理它们。
数据库中的 DELETE 操作会使 Debezium 生成两条 Kafka 记录:
- 一条记录包含
"op": "d",、before行数据以及其他一些字段。 - 一条墓碑(tombstone)记录,其键与被删除的行相同,值为
null。这条记录是 Apache Kafka 的一个标记,表示日志压缩可以移除所有具有该键的记录。
除了丢弃包含 before 行数据的记录外,你还可以配置事件展平 SMT 执行以下任一操作:
- 将该记录保留在流中,并将其编辑为只包含
"value": "null"字段。 - 将该记录保留在流中,并将其
value字段编辑为包含原before字段中的键值对,同时额外添加"__deleted": "true"条目。
类似地,除了丢弃墓碑记录外,你还可以配置事件展平 SMT 将墓碑记录保留在流中。
配置
通过在连接器的配置中添加 SMT 配置详情,可在 Kafka Connect 源连接器或接收器连接器中配置 Debezium 事件展平 SMT。例如,要获得该转换的默认行为,只需在连接器配置中添加它而不指定任何选项即可,如下例所示:
transforms=unwrap,...
transforms.unwrap.type=io.debezium.transforms.ExtractNewRecordState与任何 Kafka Connect 连接器配置一样,你可以将 transforms= 设置为多个以逗号分隔的 SMT 别名,并按你希望 Kafka Connect 应用这些 SMT 的顺序排列。
以下 .properties 示例设置了多个事件展平 SMT 选项:
transforms=unwrap,...
transforms.unwrap.type=io.debezium.transforms.ExtractNewRecordState
transforms.unwrap.delete.tombstone.handling.mode=rewrite
transforms.unwrap.add.fields=table,lsndelete.tombstone.handling.mode=rewrite
对于 DELETE 操作,移除墓碑记录,并通过将变更事件中的 value 字段扁平化来编辑 Kafka 记录。value 字段直接包含原先位于 before 字段中的键/值对。SMT 会添加 __deleted 字段并将其设置为 true,例如:
"value": {
"pk": 2,
"cola": null,
"__deleted": "true"
}add.fields=table,lsn
为简化的 Kafka 记录添加 table 和 lsn 字段的变更事件元数据。
自定义配置
连接器可能会发出多种类型的事件消息(心跳消息、墓碑消息,或有关事务与模式变更的元数据消息)。若要将转换仅应用于一部分事件,你可以定义一个 SMT 谓词语句,以便有选择地对特定事件应用该转换。
添加元数据
你可以配置事件扁平化 SMT,将原始变更事件元数据添加到简化的 Kafka 记录中。例如,你可能希望简化记录的头部或值包含以下任意内容:
- 导致变更的操作类型
- 被更改的数据库或表的名称
- 连接器特定的字段,例如 Postgres 的 LSN 字段
有关可用内容的详细信息,请参阅各连接器的文档。
要向简化 Kafka 记录的头部添加元数据,请指定 add.headers 选项;要向简化 Kafka 记录的值添加元数据,请指定 add.fields 选项。这两个选项都接受以逗号分隔的变更事件字段名列表,且不要包含空格。当存在重复的字段名时,若要为其中某个字段添加元数据,需要同时指定结构体(struct)和该字段。例如:
transforms=unwrap,...
transforms.unwrap.type=io.debezium.transforms.ExtractNewRecordState
transforms.unwrap.add.fields=op,table,lsn,source.ts_ms
transforms.unwrap.add.headers=db
transforms.unwrap.delete.tombstone.handling.mode=rewrite有了该配置,简化的 Kafka 记录将包含如下类似的内容:
{
...
"__op" : "c",
"__table": "MY_TABLE",
"__lsn": "123456789",
"__source_ts_ms" : "123456789",
...
}此外,简化的 Kafka 记录会带有一个 __db 头部。
在简化的 Kafka 记录中,SMT 会在元数据字段名前加上双下划线前缀。当你指定一个结构体(struct)时,SMT 还会在结构体名称与字段名之间插入一个下划线。
要为针对 DELETE 操作的简化 Kafka 记录添加元数据,你还必须配置 delete.tombstone.handling.mode=rewrite。
选择性应用事件展平转换的选项
除了 Debezium 连接器在数据库发生变化时发出的变更事件消息之外,连接器还会发出其他类型的消息,包括心跳消息,以及关于模式变更和事务的元数据消息。由于这些其他消息的结构与 SMT 设计用来处理的变更事件消息的结构不同,最好将连接器配置为选择性地应用 SMT,使其仅处理预期的数据变更消息。
有关如何选择性应用 SMT 的更多信息,请参阅为转换配置 SMT 谓词。
配置选项
下表描述了可用于配置事件展平 SMT 的选项。
表 1. 事件展平 SMT 配置选项说明
选项 默认值 说明
delete.tombstone.handling.mode
tombstone
Debezium 会为每个 DELETE 操作生成一条变更事件记录。此设置决定事件展平 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-to-tombstone
SMT 会从数据流中移除 TOMBSTONE 记录,并将 DELETE 记录转换为 TOMBSTONE 记录。
若要使用行数据来决定将记录路由到哪个主题,请将此选项设置为某个 after 字段属性。SMT 会将记录路由到名称与指定 after 字段属性值相匹配的主题。对于 DELETE 操作,请将此选项设置为某个 before 字段属性。
例如,配置 route.by.field=destination 会将记录路由到名称为 after.destination 值的主题。默认情况下,Debezium 连接器会将每条变更事件记录发送到这样一个主题:其名称由数据库名称和发生变更的表名组合而成。
如果你在接收器连接器上配置事件展平 SMT,当目标主题名称决定了要用简化的变更事件记录更新的数据库表名时,设置此选项可能会很有用。如果主题名称不适用于你的用例,你可以配置 route.by.field 来重新路由事件。
__(双下划线)
将此可选字符串设置为字段的前缀。
无默认值
指定一个以逗号分隔且不含空格的元数据字段列表,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 元素中。
__(双下划线)
为标头指定一个可选的前缀字符串。
无默认值
指定一个以逗号分隔(不含空格)的元数据字段列表,表示你希望 SMT 将这些字段添加到简化后 Kafka 消息的标头中。如果原始消息包含重名字段,你可以通过同时提供结构体名称与字段名称来指明要修改的字段,例如 source.ts_ms。
你还可以选择覆盖字段的原始名称,并通过在列表中添加以下格式的条目来为其指定一个区分大小写的新名称:
<field_name>:<new_field_name>。
例如:
version:VERSION, connector:CONNECTOR, source.ts_ms:EVENT_TIMESTAMP当 SMT 向简化消息的报头添加元数据字段时,它会在每个元数据字段名前加上双下划线。对于结构体(struct)形式的配置,SMT 还会在结构体名称与字段名之间插入一个下划线。
如果你指定的字段在原始变更事件消息的报头中不存在,SMT 不会将该字段添加到报头中。
没有默认值。
用于列出源消息中你希望从输出消息中删除的字段名的 Kafka 消息报头名称。
false
指定是否要 SMT 从事件的键中移除在 drop.fields.header.name 中列出的字段。
drop.fields.keep.schema.compatible
true
指定是否要 SMT 移除在 drop.fields.header.name 配置属性中包含的非可选字段。
默认情况下,SMT 只移除标记为 optional 的字段。
true
指定是否要 SMT 将记录的空值替换为源端定义的默认值。
评论
登录后参与评论
KnowForge