应用程序接口与服务提供者接口

Debezium 引擎

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

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

Debezium 引擎

Debezium 连接器通常通过部署到 Kafka Connect 服务来运行,并配置一个或多个连接器来监控上游数据库,为其所观察到的上游数据库中的所有变更生成数据变更事件。这些数据变更事件会被写入 Kafka,在那里可以被许多不同的应用程序独立消费。Kafka Connect 提供了出色的容错能力和可扩展性,因为它以分布式服务的方式运行,并确保所有已注册和已配置的连接器始终处于运行状态。例如,即使集群中的某个 Kafka Connect 端点发生故障,其余的 Kafka Connect 端点也会重新启动此前运行在已终止端点上的连接器,从而最大限度地减少停机时间并免去管理操作。

并非每个应用程序都需要这种级别的容错性和可靠性,而且它们可能也不希望依赖外部的 Kafka broker 和 Kafka Connect 服务集群。相反,有些应用程序更希望将 Debezium 连接器直接嵌入到应用程序空间中。它们仍然需要相同的数据变更事件,但更希望由连接器将事件直接发送给应用程序,而不是将其持久化存储在 Kafka 中。

这个 debezium-api 模块定义了一套小型 API,使应用程序能够使用 Debezium 引擎轻松地配置和运行 Debezium 连接器。

从 2.6.0 版本开始,Debezium 提供了 DebeziumEngine 接口的两种实现。较早的 EmbeddedEngine 实现运行单个连接器,且仅使用一个任务。该连接器按顺序发出所有记录。

在 Debezium 3.1.0.Final 及更早版本中,EmbeddedEngine 是默认实现。从 Debezium 3.2.0.Alpha1 版本开始,默认实现改为 AsyncEmbeddedEngine,EmbeddedEngine 实现不再可用。

从 2.6.0 版本开始,提供了新的 AsyncEmbeddedEngine 实现。该实现同样只运行单个连接器,但它可以在多个线程中处理记录,并在连接器支持的情况下运行多个任务(目前只有 SQL Server 和 MongoDB 的连接器支持在单个连接器内运行多个任务)。由于这两种引擎实现的是同一个接口并共享相同的 API,因此下文中的代码示例对任一引擎都适用。两种实现支持相同的配置选项。

不过,新的 AsyncEmbeddedEngine 提供了若干新的配置选项,用于设置和微调并行处理。有关这些新配置选项的说明,请参阅异步引擎属性。若想了解开发 AsyncEmbeddedEngine 的动机及其实施细节,请参阅异步嵌入式引擎设计文档。

依赖

要使用 Debezium Engine 模块,请将 debezium-api 模块添加到应用程序的依赖中。该 API 有一个现成的实现位于 debezium-embedded 模块中,同样也需要将其添加到依赖中。对于 Maven 来说,这意味着需要在应用程序的 POM 中添加以下内容:

<dependency>
    <groupId>io.debezium</groupId>
    <artifactId>debezium-api</artifactId>
    <version>${version.debezium}</version>
</dependency>
<dependency>
    <groupId>io.debezium</groupId>
    <artifactId>debezium-embedded</artifactId>
    <version>${version.debezium}</version>
</dependency>

其中 ${version.debezium} 可以是所使用的 Debezium 版本号,也可以是一个值为 Debezium 版本字符串的 Maven 属性。

同样地,还要为你的应用程序将要使用的每个 Debezium 连接器添加依赖项。例如,可以在应用程序的 Maven POM 文件中添加以下内容,以便应用程序使用 MySQL 连接器:

<dependency>
    <groupId>io.debezium</groupId>
    <artifactId>debezium-connector-mysql</artifactId>
    <version>${version.debezium}</version>
</dependency>

或者对于 MongoDB 连接器:

<dependency>
    <groupId>io.debezium</groupId>
    <artifactId>debezium-connector-mongodb</artifactId>
    <version>${version.debezium}</version>
</dependency>

本文档的其余部分介绍如何将 MySQL 连接器嵌入到你的应用程序中。其他连接器的使用方式类似,只是连接器的配置、主题和事件不同。

打包你的项目

Debezium 通过 ServiceLoader 使用 SPI 来加载实现。该实现可以基于连接器类型,也可以是自定义实现。

某些接口有多个实现。例如,io.debezium.snapshot.spi.SnapshotLock 在核心模块中有一个默认实现,同时每个连接器都有各自的特定实现。为确保 Debezium 能够找到所需的实现,你必须显式配置构建工具以合并 META-INF/services 文件。

例如,如果你使用 Maven shade 插件,请添加 ServicesResourceTransformer 转换器,如以下示例所示:

...
<configuration>
 <transformers>
    ...
    <transformer implementation="org.apache.maven.plugins.shade.resource.ServicesResourceTransformer" />
    ...
 </transformers>
...
</configuration>

或者,如果你使用 Maven Assembly 插件,则可以使用 metaInf-services 容器描述符处理器。

代码实现

你的应用程序需要为每个想要运行的连接器实例设置一个嵌入式引擎。io.debezium.engine.DebeziumEngine<R> 类作为任何 Debezium 连接器的易用包装器,完全管理连接器的生命周期。你需要使用其构建器 API 来创建 DebeziumEngine 实例,并提供以下内容:

  • 你希望接收消息的格式,例如 JSON、Avro 或 Kafka Connect 的 SourceRecord(参见输出消息格式)
  • 配置属性(可能从属性文件加载),用于定义引擎和连接器的运行环境
  • 一个方法,连接器产生的每个数据变更事件都会调用该方法

以下是一个配置并运行嵌入式 MySQL 连接器的代码示例:

// Define the configuration for the Debezium Engine with MySQL connector...
final Properties props = new Properties();
props.setProperty("name", "engine");
props.setProperty("connector.class", "io.debezium.connector.mysql.MySqlConnector");
props.setProperty("offset.storage", "org.apache.kafka.connect.storage.FileOffsetBackingStore");
props.setProperty("offset.storage.file.filename", "/path/to/storage/offsets.dat");
props.setProperty("offset.flush.interval.ms", "60000");
/* begin connector properties */
props.setProperty("database.hostname", "localhost");
props.setProperty("database.port", "3306");
props.setProperty("database.user", "mysqluser");
props.setProperty("database.password", "mysqlpw");
props.setProperty("database.server.id", "85744");
props.setProperty("topic.prefix", "my-app-connector");
props.setProperty("schema.history.internal", "io.debezium.storage.file.history.FileSchemaHistory");
props.setProperty("schema.history.internal.file.filename", "/path/to/storage/schemahistory.dat");

// Create the engine with this configuration ...
try (DebeziumEngine<ChangeEvent<String, String>> engine = DebeziumEngine.create(Json.class)
        .using(props)
        .notifying(record -> {
            System.out.println(record);
        }).build()
    ) {
    // Run the engine asynchronously ...
    ExecutorService executor = Executors.newSingleThreadExecutor();
    executor.execute(engine);

    // Do something else or wait for a signal or an event
}
// Engine is stopped when the main code is finished

我们来更详细地研究这段代码,先从前几行开始(这里我们再次列出):

// Define the configuration for the Debezium Engine with MySQL connector...
final Properties props = new Properties();
props.setProperty("name", "engine");
props.setProperty("connector.class", "io.debezium.connector.mysql.MySqlConnector");
props.setProperty("offset.storage", "org.apache.kafka.connect.storage.FileOffsetBackingStore");
props.setProperty("offset.storage.file.filename", "/path/to/storage/offsets.dat");
props.setProperty("offset.flush.interval.ms", "60000");

这会创建一个标准的 Properties 对象,用于设置引擎所需的若干字段,无论使用哪种连接器都适用。第一个字段是引擎的名称,它会被用于连接器产生的源记录及其内部状态中,因此请在应用中使用一个有意义的名称。connector.class 字段定义了继承 Kafka Connect org.apache.kafka.connect.source.SourceConnector 抽象类的类名;在本例中,我们指定 Debezium 的 MySqlConnector 类。

当 Kafka Connect 连接器运行时,它会从源中读取信息,并定期记录"偏移量"(offset),以标示它已处理了多少信息。如果连接器重启,它将使用最后记录的偏移量来确定应从源信息的哪个位置继续读取。由于连接器不知道也不关心偏移量如何存储,因此提供存储和恢复这些偏移量的方式是引擎的职责。配置中的接下来几个字段指定了我们的引擎应使用 FileOffsetBackingStore 类,将偏移量存储在本地文件系统的 /path/to/storage/offset.dat 文件中(文件名和存储位置均可任意指定)。此外,尽管连接器在产生每条源记录时都会记录偏移量,但引擎会定期将偏移量刷新到后端存储中(在我们的示例中是每分钟一次)。这些字段可根据应用的需要进行调整。

接下来的几行定义了该连接器特有的字段(具体说明见各连接器的文档),在本例中即 MySqlConnector 连接器:

    /* begin connector properties */
    props.setProperty("database.hostname", "localhost");
    props.setProperty("database.port", "3306");
    props.setProperty("database.user", "mysqluser");
    props.setProperty("database.password", "mysqlpw");
    props.setProperty("database.server.id", "85744");
    props.setProperty("topic.prefix", "my-app-connector");
    props.setProperty("schema.history.internal", "io.debezium.storage.file.history.FileSchemaHistory");
    props.setProperty("schema.history.internal.file.filename", "/path/to/storage/schemahistory.dat");

在这里,我们设置了运行 MySQL 数据库服务器的主机名和端口号,并定义了用于连接 MySQL 数据库的用户名和密码。请注意,对于 MySQL,用户名和密码应对应一个已获得以下 MySQL 权限的数据库用户:

  • SELECT
  • RELOAD
  • SHOW DATABASES
  • REPLICATION SLAVE
  • REPLICATION CLIENT

前三个权限是在读取数据库一致性快照时所必需的。后两个权限则允许数据库读取服务器的 binlog,该 binlog 通常用于 MySQL 复制。

配置中还包含一个用于 server.id 的数字标识符。由于 MySQL 的 binlog 是 MySQL 复制机制的一部分,为了读取 binlog,MySqlConnector 实例必须加入 MySQL 服务器组,这意味着该服务器 ID 必须在构成 MySQL 服务器组的所有进程中是唯一的,取值范围是 1 到 2³²-1 之间的任意整数。在我们的代码中,我们将其设置为一个相当大但略显随机的值,该值仅用于我们的应用。

配置中还指定了 MySQL 服务器的逻辑名称。连接器会将该逻辑名称包含在其生成的每条源记录的 topic 字段中,使你的应用能够识别这些记录的来源。我们的示例使用服务器名称 "products",推测是因为数据库中存储了产品信息。当然,你可以为它取任何对你的应用有意义的名称。

当 MySqlConnector 类运行时,它会读取 MySQL 服务器的 binlog,其中包含该服务器上托管的所有数据库的数据变更和模式变更。由于所有数据变更都是基于变更记录时所属表的模式来组织的,因此连接器需要跟踪所有模式变更,以便正确解码变更事件。连接器会记录模式信息,这样即使连接器重启并从上次记录的偏移量继续读取时,它也能确切知道该偏移量处的数据库模式是什么样的。连接器如何记录数据库模式历史由我们配置中的最后两个字段定义,即我们的连接器应使用 FileSchemaHistory 类将数据库模式历史变更存储在本地文件系统的 /path/to/storage/schemahistory.dat 文件中(同样,该文件可以任意命名并存储在任意位置)。

最后,使用 build() 方法构建不可变的配置。(顺便说一下,我们也可以不通过编程方式构建,而是使用某个 Configuration.read(…​) 方法从属性文件中读取配置。)

现在我们有了配置,就可以创建引擎了。以下是相关的代码行:

// Create the engine with this configuration ...
try (DebeziumEngine<ChangeEvent<String, String>> engine = DebeziumEngine.create(Json.class)
        .using(props)
        .notifying(record -> {
            System.out.println(record);
        })
        .build()) {
}

所有变更事件都会传递给给定的处理程序方法,该方法必须匹配 java.util.function.Consumer<R> 函数式接口的签名,其中 <R> 必须与调用 create() 时所指定的格式类型一致。请注意,你的应用程序的处理函数不应抛出任何异常;如果抛出了异常,引擎会记录该方法抛出的任何异常,并继续处理下一条源记录,但你的应用程序将不会再有机会处理引发该异常的特定源记录,这意味着你的应用程序可能与数据库变得不一致。

至此,我们已经获得了一个配置完毕、可以运行的 DebeziumEngine 对象,但它还不会做任何事情。DebeziumEngine 设计为由 Executor 或 ExecutorService 异步执行:

// Run the engine asynchronously ...
ExecutorService executor = Executors.newSingleThreadExecutor();
executor.execute(engine);

// Do something else or wait for a signal or an event

你的应用可以通过调用引擎的 close() 方法来安全、优雅地停止它:

// At some later time ...
engine.close();

或者,由于引擎实现了 Closeable 接口,离开 try 块时会自动调用它。

引擎的连接器将停止从源系统读取信息,把所有剩余的变更事件转发给你的处理函数,并将最新的偏移量刷新到偏移存储中。只有在这一切全部完成之后,引擎的 run() 方法才会返回。如果你的应用需要在退出前等待引擎完全停止,可以使用 ExcecutorService 的 shutdown 和 awaitTermination 方法来实现:

try {
    executor.shutdown();
    while (!executor.awaitTermination(5, TimeUnit.SECONDS)) {
        logger.info("Waiting another 5 seconds for the embedded engine to shut down");
    }
}
catch ( InterruptedException e ) {
    Thread.currentThread().interrupt();
}

你也可以在创建 DebeziumEngine 时注册 CompletionCallback 作为回调,以便在引擎终止时收到通知。

请记住,JVM 关闭时只会等待非守护线程。因此,当你在守护线程上运行引擎时,如果应用程序退出,务必等待引擎进程完成。

你的应用程序应始终正确地停止引擎,以确保优雅、完整地关闭,并保证每条源记录都恰好发送给应用程序一次。例如,不要依赖于关闭 ExecutorService,因为这会中断正在运行的线程。尽管当其线程被中断时 DebeziumEngine 确实会终止,但引擎可能无法干净地终止,并且当你的应用程序重新启动时,可能会再次看到关闭前刚处理过的部分源记录。

如前所述,DebeziumEngine 接口有两种实现。这两种实现使用相同的 API,上面的代码示例对两者都适用。唯一的例外是创建 DebeziumEngine 实例的方式。正如简介中也提到的,默认情况下使用 AsyncEmbeddedEngine 实现。因此,方法 DebeziumEngine.create(Json.class) 在内部会使用 AsyncEmbeddedEngine 实例。

输出消息格式

DebeziumEngine#create() 可以接受多个不同的参数,这些参数会影响消费者接收消息的格式。允许的值有:

  • Connect.class - 输出值是包装 Kafka Connect 的 SourceRecord 的变更事件
  • Json.class - 输出值是以 JSON 字符串编码的键值对
  • JsonByteArray.class - 输出值是格式化为 JSON 并编码为 UTF-8 字节数组的键值对
  • Avro.class - 输出值是以 Avro 序列化记录编码的键值对(详见 Avro 序列化)
  • CloudEvents.class - 输出值是以 Cloud Events 消息编码的键值对

在调用 DebeziumEngine#create() 时还可以指定头信息格式。允许的值有:

  • Json.class - 头信息值以 JSON 字符串编码
  • JsonByteArray.class - 头信息值格式化为 JSON 并编码为 UTF-8 字节数组

在内部,引擎会将数据转换委托给 Kafka Connect 或 Apicurio 转换器实现,并采用最适合执行该转换的算法。可以通过引擎属性对转换器进行参数化,以修改其行为。

下面是 JSON 输出格式的示例:

final Properties props = new Properties();
...
props.setProperty("converter.schemas.enable", "false"); // don't include schema in message
...
final DebeziumEngine<ChangeEvent<String, String>> engine = DebeziumEngine.create(Json.class)
    .using(props)
    .notifying((records, committer) -> {

        for (ChangeEvent<String, String> r : records) {
            System.out.println("Key = '" + r.key() + "' value = '" + r.value() + "'");
            committer.markProcessed(r);
        }
...

其中 ChangeEvent 数据类型即键/值对。

消息转换

在消息投递给处理器之前,可以先让其经过一条 Kafka Connect 单条消息转换(SMT)流水线。每个 SMT 可以原样传递消息、修改消息,或者将其过滤掉。该转换链通过 transforms 属性进行配置,该属性包含一个以逗号分隔的、待应用转换的逻辑名称列表。随后,transforms.<logical_name>.type 属性定义每个转换的实现类名称,而 transforms.<logical_name>.* 则定义传递给该转换的配置选项。

以下是配置示例

final Properties props = new Properties();
...
props.setProperty("transforms", "filter, router");                                               // (1)
props.setProperty("transforms.router.type", "org.apache.kafka.connect.transforms.RegexRouter");  // (2)
props.setProperty("transforms.router.regex", "(.*)");                                            // (3)
props.setProperty("transforms.router.replacement", "trf$1");                                     // (3)
props.setProperty("transforms.filter.type", "io.debezium.embedded.ExampleFilterTransform");      // (4)
  1. 定义了两个转换:filter 和 router
  2. router 转换的实现类是 org.apache.kafka.connect.transforms.RegexRouter
  3. router 转换有两个配置选项:regex 和 replacement
  4. filter 转换的实现类是 io.debezium.embedded.ExampleFilterTransform

消息转换谓词

可以为转换应用谓词,从而使转换变为可选。

配置示例如下:

final Properties props = new Properties();
...
props.setProperty("transforms", "filter");                                                 // (1)
props.setProperty("predicates", "headerExists");                                           // (2)
props.setProperty("predicates.headerExists.type", "org.apache.kafka.connect.transforms.predicates.HasHeaderKey"); //(3)
props.setProperty("predicates.headerExists.name", "header.name");                          // (4)
props.setProperty("transforms.filter.type", "io.debezium.embedded.ExampleFilterTransform");// (5)
props.setProperty("transforms.filter.predicate", "headerExists");                          // (6)
props.setProperty("transforms.filter.negate", "true");                                     // (7)
  1. 定义了一个转换 - filter
  2. 定义了一个谓词 - headerExists
  3. headerExists 谓词的实现类是 org.apache.kafka.connect.transforms.predicates.HasHeaderKey
  4. headerExists 谓词有一个配置项 - name
  5. filter 转换的实现类是 io.debezium.embedded.ExampleFilterTransform
  6. filter 转换需要谓词 headerExists
  7. filter 转换期望对谓词的结果取反,从而让该谓词用于判断头部是否存在

高级记录消费

对于某些用例,例如需要批量写入记录或对接异步 API 时,上述函数式接口可能会带来困难。在这种情况下,使用 io.debezium.engine.DebeziumEngine.ChangeConsumer<R>. 接口可能更为简便。

该接口只有一个函数,其签名如下:

/**
  * Handles a batch of records, calling the {@link RecordCommitter#markProcessed(Object)}
  * for each record and {@link RecordCommitter#markBatchFinished()} when this batch is finished.
  * @param records the records to be processed
  * @param committer the committer that indicates to the system that we are finished
  */
 void handleBatch(List<R> records, RecordCommitter<R> committer) throws InterruptedException;

正如 Javadoc 中所述,RecordCommitter 对象需要对每条记录调用一次,并在每个批次结束时调用一次。RecordCommitter 接口是线程安全的,这使得记录的处理更加灵活。

你可以选择性地覆盖所处理记录的偏移量。具体做法是:首先通过调用 RecordCommitter#buildOffsets() 构建一个新的 Offsets 对象,然后使用 Offsets#set(String key, Object value) 更新偏移量,最后调用 RecordCommitter#markProcessed(SourceRecord record, Offsets sourceOffsets) 并传入更新后的 Offsets。

要使用 ChangeConsumer API,你必须将该接口的一个实现传递给 notifying API,如下所示:

class MyChangeConsumer implements DebeziumEngine.ChangeConsumer<RecordChangeEvent<SourceRecord>> {
  public void handleBatch(List<RecordChangeEvent<SourceRecord>> records, RecordCommitter<RecordChangeEvent<SourceRecord>> committer) throws InterruptedException {
    ...
  }
}
// Create the engine with this configuration ...
DebeziumEngine<RecordChangeEvent<SourceRecord>> engine = DebeziumEngine.create(ChangeEventFormat.of(Connect.class))
        .using(props)
        .notifying(new MyChangeConsumer())
        .build();

如果使用 JSON 格式(其他格式同理),代码如下:

class JsonChangeConsumer implements DebeziumEngine.ChangeConsumer<ChangeEvent<String, String>> {
  public void handleBatch(List<ChangeEvent<String, String>> records,
    RecordCommitter<ChangeEvent<String, String>> committer) throws InterruptedException {
    ...
  }
}
// Create the engine with this configuration ...
DebeziumEngine<ChangeEvent<String, String>> engine = DebeziumEngine.create(Json.class)
        .using(props)
        .notifying(new JsonChangeConsumer())
        .build();

引擎属性

以下配置属性是必需的,除非提供了默认值(为便于排版,Java 类的包名以下 <…​> 形式替代)。

属性默认值描述
name连接器实例的唯一名称。
connector.class连接器的 Java 类名,例如 MySQL 连接器为 <…​>.MySqlConnector。
offset.storage<…​>.FileOffsetBackingStore负责持久化连接器偏移量的 Java 类名。它必须实现 <…​>.OffsetBackingStore 接口。
offset.storage.file.filename""存储偏移量的文件路径。当 offset.storage 设置为 <…​>.FileOffsetBackingStore 时必填。
offset.storage.topic""存储偏移量的 Kafka topic 名称。当 offset.storage 设置为 <…​>.KafkaOffsetBackingStore 时必填。
offset.storage.partitions""创建偏移量存储 topic 时使用的分区数。当 offset.storage 设置为 <…​>.KafkaOffsetBackingStore 时必填。
offset.storage.replication.factor""创建偏移量存储 topic 时使用的副本因子。当 offset.storage 设置为 <…​>.KafkaOffsetBackingStore 时必填。
offset.commit.policy<…​>.PeriodicCommitOffsetPolicy提交策略的 Java 类名。它根据已处理的事件数量以及距离上次提交所经过的时间,定义何时触发偏移量提交。该类必须实现 <…​>.OffsetCommitPolicy 接口。默认是基于时间间隔的周期性提交策略。
offset.flush.interval.ms60000尝试提交偏移量的时间间隔。默认为 1 分钟。
offset.flush.timeout.ms5000在取消本次操作并将偏移量数据留待下次尝试提交之前,等待记录刷新并将分区偏移量数据提交到偏移量存储的最大毫秒数。默认为 5 秒。
errors.max.retries-1连接错误发生时在最终失败前的最大重试次数(-1 = 不限制,0 = 禁用,> 0 = 重试次数)。
errors.retry.delay.initial.ms300遇到连接错误时重试的初始延迟(毫秒)。每次重试该值会翻倍,但不会超过 errors.retry.delay.max.ms。
errors.retry.delay.max.ms10000遇到连接错误时两次重试之间的最大延迟(毫秒)。

异步引擎属性

属性

默认值

说明

record.processing.threads

按需分配线程,取决于工作负载和可用 CPU 核心数。

可用于处理变更事件记录的线程数量。如果未指定值(默认情况),引擎将使用 Java ThreadPoolExecutor,根据当前工作负载动态调整线程数。最大线程数为给定机器的 CPU 核心数。如果指定了值,引擎将使用 Java 固定线程池方法,创建具有指定线程数的线程池。要使用给定机器上的所有可用核心,请将占位符值设置为 AVAILABLE_CORES。

record.processing.shutdown.timeout.ms

1000

在任务关闭后,等待已提交记录完成处理的最长时间(毫秒)。

record.processing.order

ORDERED

确定记录的产出方式。

ORDERED

记录按顺序处理;即按照从数据库获取的顺序产出记录。

UNORDERED

记录按非顺序方式处理;即可以与源数据库不同的顺序产出记录。

UNORDERED 选项的非顺序处理可获得更好的吞吐量,因为记录在完成任何 SMT 处理和消息序列化后立即产出,无需等待其他记录。当向引擎提供 ChangeConsumer 方法时,此选项不产生任何效果。

record.processing.with.serial.consumer

false

指定是否应从所提供的 Consumer 创建默认的 ChangeConsumer,从而实现串行的 Consumer 处理。如果在使用 API 创建引擎时指定了 ChangeConsumer 接口,则此选项不产生任何效果。

task.management.timeout.ms

180,000(3 分钟)

引擎等待任务生命周期管理操作(启动和停止)完成的时间,单位为毫秒。

数据库模式历史属性

部分连接器还需要一组额外的属性来配置数据库模式历史:

  • MySQL
  • SQL Server
  • Oracle
  • Db2

如果未正确配置数据库模式历史,连接器将拒绝启动。默认配置要求有一个可用的 Kafka 集群。对于其他部署方式,可以使用基于文件的数据库模式历史存储实现。

属性 默认值 说明

schema.history.internal

<…​>.KafkaSchemaHistory

负责持久化数据库模式历史的 Java 类名。
该类必须实现 <…​>.SchemaHistory 接口。

schema.history.internal.file.filename

""

用于存储数据库模式历史记录的文件路径。
当 schema.history.internal 设置为 <…​>.FileSchemaHistory 时,此项为必填。

schema.history.internal.kafka.topic

""

用于存储数据库模式历史记录的 Kafka 主题。
当 schema.history.internal 设置为 <…​>.KafkaSchemaHistory 时,此项为必填。

schema.history.internal.kafka.bootstrap.servers

""

要连接的 Kafka 集群服务器的初始列表。该集群提供用于存储数据库模式历史记录的主题。
当 schema.history.internal 设置为 <…​>.KafkaSchemaHistory 时,此项为必填。

故障处理

当引擎执行时,其连接器会在每条源记录中积极记录源偏移量,引擎会定期将这些偏移量刷新到持久化存储中。当应用程序和引擎正常关闭或崩溃后,重新启动时,引擎及其连接器将从最后记录的偏移量处恢复读取源信息。

那么,当应用程序在嵌入式引擎运行期间发生故障时会怎样呢?最终的结果是,应用程序在重启后可能会收到一些它在崩溃前已经处理过的源记录。收到多少条取决于引擎刷新偏移量到存储的频率(通过 offset.flush.interval.ms 属性设置),以及特定连接器在单次批次中返回的源记录数量。最理想的情况是每次都刷新偏移量(例如 offset.flush.interval.ms 设置为 0),但即便如此,嵌入式引擎仍然只会在从连接器接收到每批源记录之后才刷新偏移量。

例如,MySQL 连接器使用 max.batch.size 来指定批次中可以出现的最大源记录数。即使 offset.flush.interval.ms 设置为 0,应用程序在崩溃后重启时仍可能看到多达 n 条重复记录,其中 n 为批次大小。如果 offset.flush.interval.ms 属性设置得更高,那么应用程序可能会看到多达 n * m 条重复记录,其中 n 为批次的最大大小,m 为在单个偏移量刷新间隔内可能累积的批次数。(显然,可以配置嵌入式连接器不使用批处理并始终刷新偏移量,从而使应用程序永远不会收到重复的源记录。然而,这会大幅增加开销并降低连接器的吞吐量。)

最重要的是,使用嵌入式连接器时,应用程序在正常运行期间(包括正常关闭后的重启)每条源记录只会收到一次,但在崩溃或异常关闭后重启时,需要能够容忍立即收到重复事件。如果应用程序需要更严格的精确一次(exactly-once)行为,则应使用完整的 Debezium 平台,该平台可以提供精确一次保证(即使在崩溃和重启之后)。

评论

登录后参与评论

正在加载评论…