消息过滤
消息过滤
默认情况下,Debezium 会将接收到的每个数据变更事件传递给 Kafka broker。但在许多情况下,你可能只对生产者发出的事件中的一部分感兴趣。为了让你只处理与自己相关的记录,Debezium 提供了 filter 单条消息转换(SMT)。
虽然可以使用 Java 编写自定义 SMT 来实现过滤逻辑,但使用自定义编码的 SMT 也有其缺点。例如:
- 必须预先编译转换逻辑,并将其部署到 Kafka Connect。
- 每次变更都需要重新编译和重新部署代码,导致运维缺乏灵活性。
filter SMT 支持与 JSR 223(Java™ 平台脚本语言规范)集成的脚本语言。目前对使用 Go 编写 SMT 的支持仍处于孵化阶段(TinyGo 和 WebAssembly)。
JSR 223
Debezium 本身不附带任何 JSR 223 API 的实现。要在 Debezium 中使用某种表达式语言,你必须下载该语言对应的 JSR 223 脚本引擎实现。例如,对于 Groovy 3,你可以从 https://groovy-lang.org/ 下载其 JSR 223 实现。GraalVM JavaScript 的 JSR223 实现可在 https://github.com/graalvm/graaljs 获取。获取脚本引擎文件后,将它们连同该语言实现所使用的其他 JAR 文件一起添加到你的 Debezium 连接器插件目录中。
Go
Debezumber 提供了一个由社区维护的辅助 PDK(插件开发套件),以便于在 Go 中使用。
你可以通过输入以下命令来获取 Debezium SMT Go PDK:
go get github.com/debezium/debezium-smt-go-pdk以下示例展示了如何使用 Go 实现最简过滤逻辑:
package main
import (
"github.com/debezium/debezium-smt-go-pdk"
)
//export process
func process(proxyPtr uint32) uint32 {
...
}
func main() {}在上例中,main 函数是 Wasm 编译目标所必需的,而 process 函数才是执行实际过滤逻辑的入口点。
要将过滤器代码编译为 Wasm,请使用较新版本的 TinyGo,示例如下:
docker run --rm \
-v ./:/src \
-w /src tinygo/tinygo:0.34.0 bash \
-c "tinygo build --no-debug -target=wasm-unknown -o /tmp/tmp.wasm myfilter.go && cat /tmp/tmp.wasm" > \
myfilter.wasm有关使用 Go 开发 Debezium 单消息转换的更多信息,请参阅 PDK 仓库:Debezium SMT Go PDK。
安装设置
出于安全考虑,过滤器 SMT 并未随 Debezium 连接器归档包一起提供,而是作为一个独立构件 debezium-scripting-3.6.3.Final.tar.gz 提供。
要在 Debezium 连接器插件中使用基于内容的路由 SMT,你必须显式地将该 SMT 构件添加到 Kafka Connect 环境中。重要提示:一旦 Kafka Connect 实例中存在过滤器 SMT,任何被允许向该实例添加连接器的用户都可以运行脚本表达式。为确保只有授权用户才能运行脚本表达式,请务必在添加过滤器 SMT 之前,先保护好 Kafka Connect 实例及其配置接口。
在已安装 Kafka、Kafka Connect 和一个或多个 Debezium 连接器的前提下,安装过滤器 SMT 的剩余步骤如下:
将归档包中的内容解压到 Kafka Connect 环境的 Debezium 插件目录中。
然后二选一:
- 获取一个 JSR-223 脚本引擎实现,并将其内容添加到 Kafka Connect 环境的 Debezium 插件目录中。
- 在磁盘上提供已编译的
.wasm文件。
重启 Kafka Connect 进程,以使新配置生效。
Groovy 语言需要在类路径中包含以下库:
groovygroovy-json(可选)groovy-jsr223
JavaScript 语言需要在类路径中包含以下库:
graalvm.jsgraalvm.js.scriptengine
示例:基本配置
你在 Debezium 连接器的 Kafka Connect 配置中配置过滤器转换。在配置中,通过定义基于业务规则的过滤条件来指定你所关注的事件。过滤器 SMT 处理事件流时,会将每个事件与已配置的过滤条件进行评估。只有满足过滤条件的事件才会被传递给代理。
要配置 Debezium 连接器以过滤变更事件记录,需要在该 Debezium 连接器的 Kafka Connect 配置中配置 Filter SMT。配置过滤器 SMT 时,你需要指定一个用于定义过滤条件的正则表达式。
例如,你可以在连接器配置中添加如下配置。
...
transforms=filter
transforms.filter.type=io.debezium.transforms.Filter
transforms.filter.language=jsr223.groovy
transforms.filter.condition=value.op == 'u' && value.before.id == 2
...前面的示例指定了使用 Groovy 表达式语言。正则表达式 value.op == 'u' && value.before.id == 2 会移除所有消息,只保留那些表示更新(u)记录且 id 值等于 2 的消息。
使用 Go 时:
...
transforms=filter
transforms.filter.type=io.debezium.transforms.Filter
transforms.filter.language=wasm.chicory
transforms.filter.expression=file://myfilter.wasm
...自定义配置
上面的示例展示了一个简单的 SMT 配置,该配置仅用于处理包含 op 字段的 DML 事件。连接器可能发出的其他类型消息(心跳消息、墓碑消息,以及关于模式变更和事务的元数据消息)并不包含该字段。为避免处理失败,你可以定义一条 SMT 谓词语句,将转换选择性地应用于特定事件。
过滤表达式中可用的变量
Debezium 会将特定变量绑定到过滤 SMT 的求值上下文中。在创建表达式来指定过滤条件时,你可以使用 Debezium 绑定到求值上下文中的变量。通过绑定变量,Debezium 使 SMT 在对表达式中的条件求值时能够查找并解析这些变量的值。
下表列出了 Debezium 绑定到过滤 SMT 求值上下文中的变量:
表 1. 过滤表达式变量
| 名称 | 说明 | 类型 |
|---|---|---|
key |
消息的键。 | org.apache.kafka.connect.data.Struct |
value |
消息的值。 | org.apache.kafka.connect.data.Struct |
keySchema |
消息键的模式。 | org.apache.kafka.connect.data.Schema |
valueSchema |
消息值的模式。 | org.apache.kafka.connect.data.Schema |
topic |
目标主题的名称。 | String |
header |
消息头的 Java 映射。键字段为消息头的名称。header 变量公开了以下属性:- value(类型为 Object)- schema(类型为 org.apache.kafka.connect.data.Schema) |
java.util.Map<String, io.debezium.transforms.scripting.RecordHeader> |
表达式可以对其变量调用任意方法。表达式应当解析为一个布尔值,该值决定 SMT 如何处置消息。当表达式中的过滤条件求值结果为 true 时,该消息会被保留;当过滤条件求值结果为 false 时,该消息会被移除。
表达式不应产生任何副作用,也就是说,它们不应修改其所传入的任何变量。
选择性应用转换的选项
除 Debezium 连接器在数据库发生变更时发出的变更事件消息外,连接器还会发出其他类型的消息,包括心跳消息以及关于模式变更和事务的元数据消息。由于这些其他消息的结构与 SMT 所设计处理的变更事件消息结构不同,最好配置连接器选择性地应用 SMT,使其仅处理预期的数据变更消息。你可以使用以下方法之一来配置连接器,使其选择性地应用 SMT:
- 为该转换配置 SMT 谓词。
- 为该 SMT 使用 topic.regex 配置选项。
语言特性
过滤条件的表达方式取决于你所使用的脚本语言。
例如,如基本配置示例中所示,当使用 Groovy 作为表达式语言时,以下表达式会移除除 id 值为 2 的更新记录之外的所有消息:
value.op == 'u' && value.before.id == 2其他语言使用不同的方法来表达相同的条件。
Debezium MongoDB 连接器会将 after 和 patch 字段作为序列化 JSON 文档发出,而不是作为结构体。要将过滤 SMT 与 MongoDB 连接器配合使用,必须首先将 JSON 中的数组字段展开为独立的文档。这可以通过应用 MongoDB ExtractNewDocumentState SMT 来实现。
你也可以采用在表达式中使用 JSON 解析器的方式,为每个数组元素生成独立的输出文档。例如,如果你使用 Groovy 作为表达式语言,请将 groovy-json 构件添加到类路径中,然后添加如下表达式:(new groovy.json.JsonSlurper()).parseText(value.after).last_name == 'Kretchmar'。
Javascript
如果使用 JavaScript 作为表达式语言,可以调用 Struct#get() 方法来指定过滤条件,如以下示例所示:
value.get('op') == 'u' && value.get('before').get('id') == 2Javascript 与 Graal.js
如果使用 JavaScript 配合 Graal.js 来定义过滤条件,其方式与使用 Groovy 时类似。例如:
value.op == 'u' && value.before.id == 2Go with TinyGo(使用 TinyGo)
如果你使用 Go 语言配合 TinyGo 编译器来定义过滤条件,就可以借助一个完全类型化的 API,以惰性方式访问字段。例如:
var op = debezium.GetString(debezium.Get(proxyPtr, "value.op"))
var beforeId = debezium.GetInt8(debezium.Get(proxyPtr, "value.before.id"))
return debezium.SetBool(op != "d" || beforeId != 2)配置选项
下表列出了可与 filter SMT 一起使用的配置选项。
表 2. filter SMT 配置选项
属性
默认值
描述
一个可选的正则表达式,用于评估事件的目标主题名称,以确定是否应用过滤逻辑。如果目标主题的名称与 topic.regex 中的值匹配,则该转换在将事件传递给主题之前会应用过滤逻辑。如果主题的名称与 topic.regex 中的值不匹配,则该 SMT 会将事件原样传递给主题。
编写表达式所使用的语言。对于 JSR223,必须以 jsr223. 开头,例如 jsr223.groovy 或 jsr223.graal.js。Debezium 通过 JSR 223 API(“Scripting for the Java ™ Platform”) 支持引导加载。对于基于 Go 的过滤器,其值应为 wasm.chicory 或 wasm.chicory-interpreter。
针对每条消息进行求值的表达式。该表达式必须求值为布尔值,结果为 true 时保留消息,结果为 false 时移除消息。
wasm 表达式所在的文件系统位置,用于针对每条消息进行求值。该 Go 函数必须求值为布尔值,结果为 true 时保留消息,结果为 false 时移除消息。
keep
指定该转换如何处理 null(墓碑)消息。可以指定以下选项之一:
keep
(默认)直接传递这些消息。
drop
完全移除这些消息。
evaluate
对这些消息应用过滤条件。
评论
登录后参与评论
KnowForge