变形

基于内容的路由

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

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

基于内容的路由

默认情况下,Debezium 会将其从表中读取到的所有变更事件流式传输到单个静态主题中。不过,在某些情况下,你可能希望根据事件内容将选定的事件重新路由到其他主题。基于消息内容进行路由的过程在基于内容的路由消息模式中有所描述。要在 Debezium 中应用该模式,你需要使用基于内容的路由单消息转换(SMT)来编写针对每个事件进行求值的表达式。根据事件的求值结果,SMT 要么将事件消息路由到原始目标主题,要么将其重新路由到你在表达式中指定的主题。

虽然可以使用 Java 编写自定义 SMT 来实现路由逻辑,但使用自定义编码的 SMT 存在一些缺点。例如:

  • 必须预先编译转换逻辑并将其部署到 Kafka Connect。
  • 每次更改都需要重新编译和重新部署代码,导致运维缺乏灵活性。

基于内容的路由 SMT 支持可与 JSR 223(Java™ 平台脚本语言规范)集成的脚本语言。

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 连接器插件目录中。

设置

出于安全原因,基于内容的路由 SMT 并未包含在 Debezium 连接器归档文件中,而是作为独立构件 debezium-scripting-3.6.3.Final.tar.gz 提供。

要将基于内容的路由 SMT 与 Debezium 连接器插件一起使用,你必须显式地将该 SMT 构件添加到你的 Kafka Connect 环境中。重要提示:一旦 Kafka Connect 实例中存在路由 SMT,任何获准向该实例添加连接器的用户都可以运行脚本表达式。为确保只有经过授权的用户才能运行脚本表达式,请务必在添加路由 SMT 之前,保护好 Kafka Connect 实例及其配置界面。

在安装了 Kafka、Kafka Connect 以及一个或多个 Debezium 连接器之后,安装过滤器 SMT 的剩余任务如下:

  1. 下载脚本 SMT 归档文件
  2. 将归档文件的内容解压到 Kafka Connect 环境的 Debezium 插件目录中。
  3. 获取一个 JSR-223 脚本引擎实现,并将其内容添加到 Kafka Connect 环境的 Debezium 插件目录中。
  4. 重新启动 Kafka Connect 进程,以加载新的 JAR 文件。

Groovy 语言需要以下库位于类路径中:

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

JavaScript 语言需要以下库位于类路径中:

  • graalvm.js
  • graalvm.js.scriptengine

示例:基本配置

要配置 Debezium 连接器根据事件内容对变更事件记录进行路由,需要在连接器的 Kafka Connect 配置中配置 ContentBasedRouter SMT。

配置基于内容的路由 SMT 时,需要指定一个定义过滤条件的正则表达式。在配置中,您需要创建一个定义路由条件的正则表达式。该表达式定义了用于评估事件记录的模式,同时指定了目标主题的名称,与该模式匹配的事件将被路由到该主题。您指定的模式可以标识事件类型,例如表的插入、更新或删除操作。您也可以定义与特定列或行中的值相匹配的模式。

例如,要将所有更新(u)记录重新路由到 updates 主题,可以向连接器配置中添加以下配置:

...
transforms=route
transforms.route.type=io.debezium.transforms.ContentBasedRouter
transforms.route.language=jsr223.groovy
transforms.route.topic.expression=value.op == 'u' ? 'updates' : null
...

上面的示例指定了使用 Groovy 表达式语言。

不符合该模式的记录会被路由到默认主题。

自定义配置

上面的示例展示了一个简单的 SMT 配置,它被设计为仅处理包含 op 字段的 DML 事件。连接器可能发出的其他类型的消息(心跳消息、墓碑消息,或关于事务与模式变更的元数据消息)都不包含该字段。为避免处理失败,你可以定义一条 SMT 谓词语句,以便选择性地应用该转换,仅作用于特定事件。

用于基于内容的路由表达式的变量

Debezium 会将某些变量绑定到 SMT 的求值上下文中。当你创建表达式以指定条件来控制路由目标时,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 映射。键字段为消息头名称。headers 变量公开了以下属性:

  • 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 作为表达式语言时,下面的表达式会将所有更新(u)记录重新路由到 updates 主题,而将其他记录路由到默认主题:

value.op == 'u' ? 'updates' : null

其他语言使用不同的方法来表达相同的条件。

Debezium MongoDB 连接器会将 after 和 patch 字段作为序列化的 JSON 文档发出,而不是作为结构体。要在 MongoDB 连接器中使用 ContentBasedRouting SMT,必须首先将 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' ? 'updates' : null

使用 Graal.js 的 JavaScript

在使用带有 Graal.js 的 JavaScript 创建基于内容的路由条件时,你所采用的方式与使用 Groovy 时类似。例如:

value.op == 'u' ? 'updates' : null

选择 TinyGo

使用 Go 配合 TinyGo 编译器创建基于内容的路由条件时,你可以借助完全类型的 API 来延迟访问字段。例如:

var value = debezium.Get(proxyPtr, "value")
if !debezium.IsNull(value) {
    var op = debezium.GetString(debezium.Get(proxyPtr, "value.op"))
    if op == "u" {
        return debezium.SetString("updates")
    }
}
return debezium.SetNull()

配置选项

属性

默认值

描述

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。

topic.expression

无默认值

对每条消息进行求值的表达式。其结果必须是 String 值,结果为非空时会将消息重新路由到新主题,结果为 null 时则将消息路由到默认主题。使用 Go 时,此属性的值指定 wasm 表达式所在的文件系统位置,以便对每条消息进行求值。Go 函数必须求值为 String 值或 Null。

null.handling.mode

keep

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

keep

(默认)直接传递消息。

drop

完全移除消息。

evaluate

对消息应用条件逻辑。

评论

登录后参与评论

正在加载评论…