变形

消息过滤

qianmoQqianmoQ· 更新于 2026-09-28· 阅读 17 分钟· 0 次阅读

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

消息过滤

默认情况下,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 的剩余步骤如下:

  1. 下载脚本 SMT 归档包

  2. 将归档包中的内容解压到 Kafka Connect 环境的 Debezium 插件目录中。

  3. 然后二选一:

    1. 获取一个 JSR-223 脚本引擎实现,并将其内容添加到 Kafka Connect 环境的 Debezium 插件目录中。
    2. 在磁盘上提供已编译的 .wasm 文件。
  4. 重启 Kafka Connect 进程,以使新配置生效。

Groovy 语言需要在类路径中包含以下库:

  • groovy
  • groovy-json(可选)
  • groovy-jsr223

JavaScript 语言需要在类路径中包含以下库:

  • graalvm.js
  • graalvm.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:

语言特性

过滤条件的表达方式取决于你所使用的脚本语言。

例如,如基本配置示例中所示,当使用 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') == 2

Javascript 与 Graal.js

如果使用 JavaScript 配合 Graal.js 来定义过滤条件,其方式与使用 Groovy 时类似。例如:

value.op == 'u' && value.before.id == 2

Go 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 中的值匹配,则该转换在将事件传递给主题之前会应用过滤逻辑。如果主题的名称与 topic.regex 中的值不匹配,则该 SMT 会将事件原样传递给主题。

language

编写表达式所使用的语言。对于 JSR223,必须以 jsr223. 开头,例如 jsr223.groovy 或 jsr223.graal.js。Debezium 通过 JSR 223 API(“Scripting for the Java ™ Platform”) 支持引导加载。对于基于 Go 的过滤器,其值应为 wasm.chicory 或 wasm.chicory-interpreter。

condition

针对每条消息进行求值的表达式。该表达式必须求值为布尔值,结果为 true 时保留消息,结果为 false 时移除消息。

expression

wasm 表达式所在的文件系统位置,用于针对每条消息进行求值。该 Go 函数必须求值为布尔值,结果为 true 时保留消息,结果为 false 时移除消息。

null.handling.mode

keep

指定该转换如何处理 null(墓碑)消息。可以指定以下选项之一:

keep

(默认)直接传递这些消息。

drop

完全移除这些消息。

evaluate

对这些消息应用过滤条件。

评论

登录后参与评论

正在加载评论…