集成

Quarkus 的 Debezium 扩展

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

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

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 方法时,事件会通过回退链进行分发:

  1. 首先检查带有 filter 的方法,由过滤器的 shouldCapture 决定。
  2. 其次是为目标注册了反序列化器的方法(仅适用于单个事件)。
  3. 最后是带有 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) {
        //
    }
}

以下事件可用:

  • DebeziumNotification
  • SnapshotStarted
  • SnapshotInProgres
  • SnapshotTableScanCompleted
  • SnapshotAborted
  • SnapshotSkipped
  • SnapshotCompleted
  • SnapshotPaused
  • SnapshotResumed

数据类型转换器

可以在扩展中使用 @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

评论

登录后参与评论

正在加载评论…