Quarkus 的 Debezium 扩展
Quarkus 的 Debezium 扩展
以下文档介绍你的 Quarkus 应用程序如何与 Debezium 集成。
简介
Quarkus 的 Debezium 扩展 将 Debezium 运行时 集成到 Quarkus 应用程序中,使开发人员能够在轻量级、云原生的应用程序中直接消费受支持数据库的 变更数据捕获(CDC)事件。
何时使用该扩展
- Quarkus 中的 CDC:将 Debezium 无缝嵌入 Quarkus 应用,非常适合微服务场景。
- 免 Kafka 方案:当只需要简单的流处理或进程内处理时,无需部署完整的 Kafka 基础设施
- 轻量快速:充分利用 Quarkus 原生的快速启动时间和低内存占用
- 对开发者友好:通过 Quarkus 原生的配置和依赖管理简化设置过程
特性
该扩展支持成熟的 Debezium 功能以及更多内容!以下是新增或重新审视的功能列表:
- Debezium 捕获监听器
- Debezium 自定义反序列化器
- Debezium 生命周期事件
- Debezium 心跳事件
- Debezium 自定义数据类型转换器
- Debezium 通知事件
- Debezium 后处理器
支持的连接器
支持以下连接器:
- 用于 Postgres 连接器的
debezium-quarkus-postgres扩展 - 用于 Mongodb 连接器的
debezium-quarkus-mongodb扩展 - 用于 MySQL 连接器的
debezium-quarkus-mysql扩展 - 用于 MariaDB 连接器的
debezium-quarkus-mariadb扩展 - 用于 MS Sql Server 连接器的
debezium-quarkus-sqlserver扩展
前置条件
- 已安装 JDK 21+ 并正确配置
JAVA_HOME - Apache maven 3.9.9
- Quarkus 版本 3.25.0
- Docker 或 Podman
- 如需构建原生可执行文件,可选择安装并正确配置 Mandrel 或 GraalVM
可运行示例
debezium-example 仓库 中提供了一个完整可用的示例,演示了该扩展的能力。它可以作为搭建和测试你自己的 Quarkus 应用程序的参考。
初始化 Quarkus 项目
如果你已经有一个现成的 Quarkus 应用程序,可以跳过本节。
创建新 Quarkus 项目的最简单方式是打开终端并运行以下命令:
mvn io.quarkus.platform:quarkus-maven-plugin:3.25.0:create \
-DprojectGroupId=org.acme \
-DprojectArtifactId=getting-started它会在 ./getting-started 中生成以下内容:
- Quarkus 样板文件
- 应用配置
将扩展添加到项目中
在项目中,将 Debezium 扩展依赖添加到应用的 pom.xml 中:
<dependencies>
<!-- [...] -->
<dependency>
<groupId>io.debezium.quarkus</groupId>
<artifactId>debezium-quarkus-postgres</artifactId>
<version>${version.debezium}</version>
</dependency>
<!-- [...] -->
</dependencies>配置 Debezium 连接器
Debezium 扩展自带一个 Connector,用于监听数据源的 CDC 事件,但必须进行最基本的配置。
该扩展使用 Quarkus(SmallRye Config API)提供所有与配置相关的机制。Debezium 的配置使用前缀 quarkus.debezium.*,如果我们使用位于 src/main/resources/application.properties 的应用配置文件,那么最小配置如下:
# Debezium CDC configuration
quarkus.debezium.offset.storage=org.apache.kafka.connect.storage.MemoryOffsetBackingStore
quarkus.debezium.name=native
quarkus.debezium.topic.prefix=native
quarkus.debezium.plugin.name=pgoutput
quarkus.debezium.snapshot.mode=initial
# datasource configuration
quarkus.datasource.db-kind=postgresql
quarkus.datasource.username=<your username>
quarkus.datasource.password=<your password>
quarkus.datasource.jdbc.url=jdbc:postgresql://localhost:5432/hibernate_orm_test
quarkus.datasource.jdbc.max-size=16可用的配置参数请参阅 Debezium 文档。此外,还必须按照 Debezium 运行时的要求指定数据源配置参数。
从 Debezium 捕获事件
基于上一节的最小配置,你的 Quarkus 应用可以直接接收 CDC 事件负载:
import io.debezium.runtime.CapturingEvent;
import jakarta.enterprise.context.ApplicationScoped;
import org.apache.kafka.connect.source.SourceRecord;
import io.debezium.runtime.Capturing;
@ApplicationScoped
public class ProductHandler {
@Capturing
public void capture(CapturingEvent<SourceRecord> record) {
// process your events
}
}CapturingEvent<T> 包含与数据库操作类型相关的信息:
@Capturing
public void capture(CapturingEvent<SourceRecord> record) {
switch (record) {
case Create<SourceRecord> event -> {}
case Delete<SourceRecord> event -> {}
case Message<SourceRecord> event -> {}
case Read<SourceRecord> event -> {}
case Truncate<SourceRecord> event -> {}
case Update<SourceRecord> event -> {}
}
}过滤变更数据捕获事件
按目标过滤
可以按 destination 过滤事件:
@Capturing(destination = "native.inventory.products")
public void capture(CapturingEvent<SourceRecord> record) {
// process your event
}默认情况下,Debezium 连接器的 destination 由配置中定义的 prefix 名称与数据库名称以及发生变更的表名组合而成。在某些情况下,destination 会通过 SMT 被重新定义。
通过自定义过滤策略
如需进行更高级的过滤,你可以提供一个自定义的 CapturingFilterStrategy,它接收完整的事件并决定是否捕获该事件:
@ApplicationScoped
public class ProductsOnlyFilter implements CapturingFilterStrategy {
@Override
public boolean shouldCapture(CapturingEvent<SourceRecord, SourceRecord> event) {
return "native.inventory.products".equals(event.destination());
}
}然后在注解中引用它:
@Capturing(filter = ProductsOnlyFilter.class)
public void capture(CapturingEvent<SourceRecord, SourceRecord> event) {
// only receives events accepted by the filter
}filter 属性与 destination 和 engine 互斥,此约束在构建时进行校验。 |
|---|
按字段变更过滤
内置的 CapturingFieldsFilterStrategy 可以根据特定被监听字段的状态来过滤事件,而该状态取决于事件类型:
- 创建 / 读取:只要记录结构中存在任一被监听字段,即被捕获
- 更新:仅当至少一个被监听字段的值在前后状态之间发生变化时,才被捕获
- 删除:只要 before 结构中存在任一被监听字段,即被捕获
- 截断 / 消息:不会被捕获(没有字段级数据)
继承该类并传入需要监听的字段名称:
@ApplicationScoped
public class PriceChangeFilter extends CapturingFieldsFilterStrategy {
public PriceChangeFilter() {
super(Set.of("name", "price"));
}
}然后使用它:
@Capturing(filter = PriceChangeFilter.class)
public void capture(CapturingEvent<SourceRecord, SourceRecord> event) {
// only receives events where 'name' or 'price' are impacted
}批量事件过滤
过滤策略同样适用于批量事件处理(CapturingEvents<BatchEvent>):
@Capturing(filter = PriceChangeFilter.class)
public void onBatch(CapturingEvents<BatchEvent> events) {
// only receives batch events where at least one record impacts 'name' or 'price'
}对于批量事件,批次中的每条记录都会单独由 shouldCapture 进行评估,只有匹配的记录才会被传递给过滤器处理器。
回退链
当定义了多个 @Capturing 方法时,事件会通过回退链进行分发:
- 首先检查带有
filter的方法,由过滤器的shouldCapture决定。 - 其次是为目标注册了反序列化器的方法(仅适用于单个事件)。
- 最后是带有
destination的方法(先精确匹配,再回退到通配符*)。
首个匹配者生效。不带 destination 的 @Capturing() 视为通配符,会捕获前面各层未匹配到的任何事件。如果未定义通配符,则未匹配的事件会被跳过。
同样的回退链同时适用于单个事件和批量事件的处理。
基于 Jackson 的反序列化器
Quarkus 内置了基于 Jackson 的 JSON 序列化与反序列化支持。已有现成的 ObjectMapperDeserializer,可用于通过 Jackson 反序列化所有数据对象。
需要对相应的反序列化器类进行子类化。因此,我们要创建一个继承自 ObjectMapperDeserializer 的 ProductDeserializer。
public class ProductDeserializer extends ObjectMapperDeserializer<Product> {
public ProductDeserializer() {
super(Product.class);
}
}最后,为特定目标配置捕获通道以使用 Jackson 反序列化器:
quarkus.debezium.capturing.products.destination=native.inventory.products
quarkus.debezium.capturing.products.deserializer=com.acme.product.jackson.ProductDeserializer并在代码中使用它,为反序列化器指定目标位置:
import io.debezium.runtime.CapturingEvent;
import jakarta.enterprise.context.ApplicationScoped;
import org.apache.kafka.connect.source.SourceRecord;
import io.debezium.runtime.Capturing;
@ApplicationScoped
public class ProductHandler {
@Capturing(destination = "native.inventory.products")
public void capture(CapturingEvent<Product> record) {
// process your events
}
}或者只获取不含 CapturingEvent<T> 的反序列化对象:
| 请注意,在这种情况下,你将无法获得与数据库操作相关的信息 |
|---|
import io.debezium.runtime.CapturingEvent;
import jakarta.enterprise.context.ApplicationScoped;
import org.apache.kafka.connect.source.SourceRecord;
import io.debezium.runtime.Capturing;
@ApplicationScoped
public class ProductHandler {
@Capturing(destination = "native.inventory.products")
public void capture(Product product) {
// process your events
}
}生命周期事件
可以获取与 Debezium 监听生命周期事件状态相关的信息:
import io.debezium.runtime.events.*;
import jakarta.enterprise.context.ApplicationScoped;
import jakarta.enterprise.event.Observes;
@ApplicationScoped
public class LifecycleListener {
public void started(@Observes ConnectorStartedEvent event) {
// your logic
}
public void stopped(@Observes ConnectorStoppedEvent connectorStoppedEvent) {
// your logic
}
public void tasksStarted(@Observes TasksStartedEvent tasksStartedEvent) {
// your logic
}
public void tasksStopped(@Observes TasksStoppedEvent tasksStoppedEvent) {
// your logic
}
public void pollingStarted(@Observes PollingStartedEvent pollingStartedEvent) {
// your logic
}
public void pollingStopped(@Observes PollingStoppedEvent pollingStoppedEvent) {
// your logic
}
public void completed(@Observes DebeziumCompletionEvent debeziumCompletionEvent) {
// your logic
}
}以下事件可用:
ConnectorStartedEvent在 Debezium 启动连接器时触发ConnectorStoppedEvent在 Debezium 停止连接器时触发TasksStartedEvent在连接器任务启动时触发TasksStoppedEvent在连接器任务停止时触发PollingStartedEvent在 Debezium 引擎开始轮询连接器变更时触发PollingStoppedEvent在 Debezium 引擎停止轮询连接器变更时触发DebeziumCompletionEvent在 Debezium 引擎完成关闭后触发。它包含此前执行是否成功或失败的全部信息,以及失败的原因和错误
心跳事件
你可以在 Quarkus 应用程序中监听心跳事件:
import io.debezium.runtime.events.DebeziumHeartbeat;
import jakarta.enterprise.context.ApplicationScoped;
import jakarta.enterprise.event.Observes;
@ApplicationScoped
public class HeartbeatListener {
public void heartbeat(@Observes DebeziumHeartbeat heartbeat) {
//
}
}DebeziumHeartbeat 包含与以下内容相关的信息:
- 连接器
- Debezium 状态
- 分区
- 偏移量
通知事件
Debezium 通知提供有关细粒度状态(snapshot 和 streaming)的事件,这些事件始终以 Jakarta 事件的形式提供:
import io.quarkus.debezium.notification.SnapshotEvent;
import io.quarkus.debezium.notification.DebeziumNotification;
import jakarta.enterprise.context.ApplicationScoped;
import jakarta.enterprise.event.Observes;
@ApplicationScoped
public class NotificationListener {
public void snapshot(@Observes SnapshotEvent event) {
//
}
public void notification(@Observes DebeziumNotification event) {
//
}
}以下事件可用:
DebeziumNotificationSnapshotStartedSnapshotInProgresSnapshotTableScanCompletedSnapshotAbortedSnapshotSkippedSnapshotCompletedSnapshotPausedSnapshotResumed
数据类型转换器
可以在扩展中使用 @CustomConverter 注解定义一个 Debezium 自定义转换器,并实例化一个 ConverterDefinition 来定义类型转换:
import io.debezium.relational.CustomConverterRegistry.ConverterDefinition;
import io.debezium.runtime.CustomConverter;
import io.debezium.spi.converter.ConvertedField;
import jakarta.enterprise.context.ApplicationScoped;
import org.apache.kafka.connect.data.SchemaBuilder;
@ApplicationScoped
public class StringConverter {
@CustomConverter
public ConverterDefinition<SchemaBuilder> bind(ConvertedField field) {
return new ConverterDefinition<>(SchemaBuilder.string(), String::valueOf);
}
}这种转换会应用于 CDC 事件中的所有字段。若只想对部分字段应用转换,可以为 CustomConverter 增加一个 FieldFilterStrategy,由它来筛选出需要处理的字段:
@CustomConverter(filter = CustomFieldFilterStrategy.class)
public ConverterDefinition<SchemaBuilder> filteredBind(ConvertedField field) {
return new ConverterDefinition<>(SchemaBuilder.string(), String::valueOf);
}
@ApplicationScoped
public static class CustomFieldFilterStrategy implements FieldFilterStrategy {
@Override
public boolean filter(ConvertedField field) {
// your logic
return false;
}
}后置处理器
后置处理器比 SMT 更早在事件流程中应用轻量级的逐消息变更,使它们能够在 Debezium 的上下文中修改消息。这使它们比转换更加高效。后置处理器可以通过两种方式定义:作为配置参数,或使用 @PostProcessing 注解。
关于配置,官方文档列出了可用的参数,例如 Reselect 后置处理器的参数:
quarkus.debezium.post.processors=reselector
quarkus.debezium.post.processors.reselector.type=io.debezium.processors.reselect.ReselectColumnsPostProcessor
quarkus.debezium.post.processors.reselector.reselect.unavailable.values=true
quarkus.debezium.post.processors.reselector.reselect.null.values=true
quarkus.debezium.post.processors.reselector.reselect.use.event.key=false
quarkus.debezium.post.processors.reselector.reselect.error.handling.mode=WARN对于代码,在扩展中提供了 @PostProcessing 注解,通过它可以访问 key 和 Struct:
import io.debezium.runtime.PostProcessing;
import jakarta.enterprise.context.ApplicationScoped;
import org.apache.kafka.connect.data.Struct;
@ApplicationScoped
public class PostProcessorHandler {
@PostProcessing
public void processing(Object key, Struct struct) {
// apply your logic
}
}DevService 支持
Quarkus 在开发模式和测试模式下会通过 Dev Services 自动配置未配置的服务。当引入一个扩展而未进行配置时,Quarkus 会启动所需的服务(通过 Testcontainers),并将其连接到你的应用。对于 Debezium 来说,所需的配置是 Quarkus 默认镜像所不支持的。该扩展已经内置了 Dev Service,并配置了适用于变更数据捕获的镜像,但该支持仍属于实验性功能,如果出现错误或问题,你可以通过以下属性将其禁用:
quarkus.datasource.devservices.enabled=false或者使用官方 Debezium 镜像进行覆盖
quarkus.datasource.devservices.image-name=quay.io/debezium/postgres:15评论
登录后参与评论
KnowForge