使用SMT谓词选择性地应用转换
选择性地应用转换
为连接器配置单消息转换(SMT)时,你可以为该转换定义一个断言。断言指定了如何将转换有条件地应用于连接器所处理的一部分消息。你可以为源连接器(例如 Debezium)以及接收连接器所配置的转换指定断言。
SMT 断言
Debezium 提供了多种单消息转换(SMT),你可以在 Kafka Connect 将事件记录保存到 Kafka 主题之前,使用这些转换来修改事件记录。默认情况下,当你为 Debezium 连接器配置其中某个 SMT 时,Kafka Connect 会对连接器发出的每一条记录应用该转换。不过,在某些情况下,你可能希望有选择地应用转换,使其仅修改那些具有共同特征的变更事件消息子集。
例如,对于 Debezium 连接器,你可能希望仅对来自特定表的事件消息,或包含特定标头键的事件消息运行转换。在运行 Apache Kafka 2.6 或更高版本的环境中,你可以在转换后附加一条断言语句,指示 Kafka Connect 仅对某些记录应用 SMT。在断言中,你指定一个条件,Kafka Connect 使用该条件来评估其处理的每条消息。当 Debezium 连接器发出变更事件消息时,Kafka Connect 会依据已配置的断言条件检查该消息。如果该事件消息满足条件,Kafka Connect 就会应用转换,然后将消息写入 Kafka 主题。不满足条件的消息则原样发送到 Kafka。
对于你为接收连接器 SMT 定义的断言,情况与此类似。连接器从 Kafka 主题读取消息,Kafka Connect 会依据断言条件对这些消息进行评估。如果某条消息满足条件,Kafka Connect 就会应用转换,然后将消息传递给接收连接器。
定义断言后,你可以重复使用它并将其应用于多个转换。断言还包含一个 negate 选项,可用于反转断言,使断言条件仅应用于不满足断言语句中所定义条件的记录。你可以使用 negate 选项将该断言与基于条件取反的其他转换搭配使用。
断言元素
断言包含以下元素:
predicates前缀- 别名(例如
isOutboxTable) - 类型(例如
org.apache.kafka.connect.transforms.predicates.TopicNameMatches)。Kafka Connect 提供了一组默认的断言类型,你可以通过定义自己的自定义断言来对其进行扩展。 - 条件语句以及任意附加配置属性(具体取决于断言类型,例如正则命名模式)
默认断言类型
默认提供以下谓词类型:
HasHeaderKey
指定事件消息头中你希望 Kafka Connect 评估的键名。对于包含指定名称头键的任何记录,该谓词的求值结果为 true。
RecordIsTombstone
匹配 Kafka 墓碑(tombstone)记录。对于值为 null 的任何记录,该谓词的求值结果为 true。将此谓词与过滤器 SMT 配合使用,即可删除墓碑记录。此谓词没有配置参数。
Kafka 中的墓碑是指带有一个键、负载为 0 字节且值为 null 的记录。当 Debezium 连接器处理源数据库中的删除操作时,连接器会为该删除操作发出两个变更事件:
一个删除操作(
"op" : "d")事件,其中提供了数据库记录的先前值。一个墓碑事件,它具有相同的键,但值为
null。墓碑代表该行的删除标记。当 Kafka 启用日志压缩时,在压缩过程中,Kafka 会移除所有与墓碑共享相同键的事件。日志压缩会周期性执行,压缩间隔由主题的
delete.retention.ms设置控制。虽然可以配置 Debezium 使其不发出墓碑事件,但最好允许 Debezium 发出墓碑,以便在日志压缩期间保持预期行为。抑制墓碑会阻止 Kafka 在日志压缩期间移除已删除键对应的记录。如果你的环境中包含无法处理墓碑的接收器连接器,可以配置该接收器连接器使用带有
RecordIsTombstone谓词的 SMT 来过滤掉墓碑记录。
TopicNameMatches
一个正则表达式,用于指定你希望 Kafka Connect 匹配的主题名称。当连接器记录的主题名称与指定的正则表达式匹配时,该谓词为真。使用此谓词,可根据源表的名称对记录应用 SMT。
其他资源
定义 SMT 谓词
配置 Kafka Connect 谓词与配置转换类似。你需要指定一个谓词别名,将该别名与某个转换关联,然后定义谓词的类型和配置。
前提条件
- Debezium 环境运行的是 Apache Kafka 2.6 或更高版本。
- 已为 Debezium 连接器配置了 SMT。
操作步骤
在 Debezium 连接器配置中,为
predicates参数指定一个谓词别名,例如IsOutboxTable。在连接器配置中,将该谓词别名附加到转换别名之后,从而把谓词别名与你希望有条件应用的转换关联起来:
transforms.<TRANSFORM_ALIAS>.predicate=<PREDICATE_ALIAS>
例如:
transforms.outbox.predicate=IsOutboxTable通过指定谓词的类型并为配置参数提供值来配置该谓词。
关于类型,请指定 Kafka Connect 中可用的以下默认类型之一:
HasHeaderKey
RecordIsTombstone
TopicNameMatches
例如:
predicates.IsOutboxTable.type=org.apache.kafka.connect.transforms.predicates.TopicNameMatches
对于 TopicNameMatch 或
HasHeaderKey谓词,请指定要匹配的主题或标头名称的正则表达式。例如:
predicates.IsOutboxTable.pattern=outbox.event.*
如果要对条件取反,请在转换别名后附加
negate关键字,并将其设置为true。例如:
transforms.outbox.negate=true
上述属性会反转谓词所匹配的记录集合,从而使 Kafka Connect 对任何不符合谓词指定条件的记录应用转换。
示例:用于 outbox 事件路由转换的 TopicNameMatch 谓词
以下示例展示了 Debezium 连接器配置,该配置仅对 Debezium 发送到 Kafka outbox.event.order 主题的消息应用 outbox 事件路由转换。
由于 TopicNameMatch 谓词仅对来自 outbox 表(outbox.event.*)的消息求值为 true,因此对数据库中其他表产生的消息不会应用该转换。
transforms=outbox
transforms.outbox.predicate=IsOutboxTable
transforms.outbox.type=io.debezium.transforms.outbox.EventRouter
predicates=IsOutboxTable
predicates.IsOutboxTable.type=org.apache.kafka.connect.transforms.predicates.TopicNameMatches
predicates.IsOutboxTable.pattern=outbox.event.*忽略墓碑事件
你可以控制 Debezium 是否发出墓碑事件,以及 Kafka 保留这些事件的时间长短。根据数据管道的不同,你可能需要为连接器设置 tombstones.on.delete 属性,使 Debezium 不发出墓碑事件。
是否启用 Debezium 发出墓碑事件,取决于你的环境中主题的消费方式以及汇端消费者的特性。有些汇端连接器依赖墓碑事件来从下游数据存储中删除记录。当汇端连接器依赖墓碑记录来指示何时删除下游数据存储中的记录时,应将 Debezium 配置为发出这些记录。
当你配置 Debezium 生成墓碑事件时,还需要进一步配置,以确保汇端连接器能够接收到这些事件。必须为主题设置合适的保留策略,以便连接器有足够的时间读取事件消息,然后再由 Kafka 在日志压缩过程中将其删除。主题在压缩前保留墓碑的时间长度由主题的 delete.retention.ms 属性控制。
默认情况下,连接器的 tombstones.on.delete 属性设置为 true,使连接器在每个删除事件之后生成一个墓碑。如果你将该属性设置为 false,以阻止 Debezium 将墓碑记录保存到 Kafka 主题中,那么缺少墓碑记录可能会导致意想不到的后果。Kafka 在日志压缩过程中依赖墓碑记录来删除与已删除键相关的记录。
如果你需要支持无法处理 null 值记录的汇端连接器或下游 Kafka 消费者,建议为连接器配置一个带有谓词的 SMT,该谓词使用 RecordIsTombstone 谓词类型,在消费者读取之前移除墓碑消息,而不是阻止 Debezium 发出墓碑事件。
操作步骤
要阻止 Debezium 为已删除的数据库记录发出墓碑事件,请将连接器选项
tombstones.on.delete设置为false。例如:
“tombstones.on.delete”: “false”
评论
登录后参与评论
KnowForge