You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

嵌入式Debezium(1.2.0)仅启动时捕获变更,后续无捕获

嵌入式Debezium 1.2.0仅在启动时捕获SQL Server变更,后续无响应的问题排查与解决

我之前也碰到过类似的情况,结合你的配置和日志信息,咱们一步步拆解问题:

先还原你的场景:你在Spring应用里运行嵌入式Debezium 1.2.0连接SQL Server,启动时能正常捕获最新的数据库变更,但日志输出Finished streaming和Connected metrics set to 'false'后,就再也抓不到后续变更了,只有重启应用才会再次触发捕获,全程没有错误日志。你的配置代码如下:

final Properties props = new Properties();
props.setProperty("name", "engine");
props.setProperty("connector.class", "io.debezium.connector.sqlserver.SqlServerConnector");
props.setProperty("offset.storage", "org.apache.kafka.connect.storage.FileOffsetBackingStore");
props.setProperty("offset.storage.file.filename", "/tmp/offsets.dat");
props.setProperty("offset.flush.interval.ms", "60000");
/* begin connector properties */
props.setProperty("database.hostname", "xxxx");
props.setProperty("database.port", "xxxx");
props.setProperty("database.user", "xxxx");
props.setProperty("database.password", "xxxx");
props.setProperty("database.server.id", "xxxx");
props.setProperty("database.server.name", "xxxx");
props.setProperty("database.dbname", "xxxx");
props.setProperty("database.history", "io.debezium.relational.history.FileDatabaseHistory");
props.setProperty("database.history.file.filename", "~logs/dbhistory.dat");
props.setProperty("snapshot.lock.timeout.ms", "-1");
try (DebeziumEngine<ChangeEvent<String, String>> engine = DebeziumEngine.create(Json.class)
        .using(props)
        .notifying(this::handleEvent)
        .build()) {
    // Run the engine asynchronously ...
    ExecutorService executor = Executors.newSingleThreadExecutor();
    executor.execute(engine);
    // Do something else or wait for a signal or an event
} catch (IOException | InterruptedException e) {
    logger.error("Unable to start debezium " + e);
}
private void handleEvent(ChangeEvent<String, String> changeEvent) {
    logger.info(changeEvent.toString());
}

最可能的元凶:Debezium Engine被提前关闭

你用了try-with-resources来创建DebeziumEngine,这个语法糖会在try块执行完毕后自动调用engine.close()。如果你的"Do something else or wait for a signal or an event"这段代码很快就跑完了,try块会立刻结束,直接关闭引擎——这就解释了为什么快照完成后就停了,引擎都没在运行了,怎么可能捕获后续变更?

解决方法:
要么去掉try-with-resources,手动管理引擎的生命周期;要么在try块里加阻塞逻辑,让主线程一直等待,直到收到停止信号(比如应用关闭)。给你改个示例:

DebeziumEngine<ChangeEvent<String, String>> engine = null;
ExecutorService executor = null;
try {
    final Properties props = new Properties();
    // 这里保留你的所有配置项...

    engine = DebeziumEngine.create(Json.class)
            .using(props)
            .notifying(this::handleEvent)
            .build();

    executor = Executors.newSingleThreadExecutor();
    executor.execute(engine);

    // 用CountDownLatch阻塞主线程,直到应用收到关闭信号
    CountDownLatch shutdownLatch = new CountDownLatch(1);
    Runtime.getRuntime().addShutdownHook(new Thread(() -> {
        shutdownLatch.countDown();
        logger.info("Received shutdown signal, stopping Debezium engine...");
    }));
    shutdownLatch.await(); // 一直等,直到应用关闭

} catch (IOException | InterruptedException e) {
    logger.error("Failed to start Debezium engine", e);
} finally {
    // 手动关闭引擎和线程池
    if (engine != null) {
        try {
            engine.close();
        } catch (IOException e) {
            logger.error("Error closing Debezium engine", e);
        }
    }
    if (executor != null) {
        executor.shutdown();
        try {
            if (!executor.awaitTermination(10, TimeUnit.SECONDS)) {
                executor.shutdownNow();
            }
        } catch (InterruptedException e) {
            executor.shutdownNow();
        }
    }
}

其他需要排查的点

1. Database History文件路径无效

你设置的database.history.file.filename是~logs/dbhistory.dat,Java不认识~这个用户主目录的简写,会把它当成普通的文件名,导致Debezium无法写入数据库历史文件。没有历史文件,连接器就没法跟踪已经处理过的快照和日志位置,快照完成后可能就停了。

解决方法:用绝对路径,或者通过System.getProperty("user.home")构建路径:

String dbHistoryPath = System.getProperty("user.home") + "/logs/dbhistory.dat";
props.setProperty("database.history.file.filename", dbHistoryPath);

记得提前创建logs目录,确保应用有读写权限。

2. SQL Server CDC未正确配置或权限不足

Debezium 1.2.0的SQL Server连接器依赖CDC来持续捕获变更,必须确保:

  • 目标数据库已经开启CDC:EXEC sys.sp_cdc_enable_db;
  • 你要捕获的表已经开启CDC:EXEC sys.sp_cdc_enable_table @source_schema = N'dbo', @source_name = N'你的表名', @role_name = NULL;
  • 连接数据库的用户至少拥有VIEW DATABASE STATE权限,以及CDC系统表的SELECT权限(最好是db_owner权限,避免踩权限坑)

3. Offset存储的权限问题

offset.storage.file.filename设为/tmp/offsets.dat,确保应用对/tmp目录有读写权限。如果无法写入偏移量,Debezium重启后会重新跑快照,运行时也没法更新偏移量,可能导致后续变更无法被正确跟踪。

4. 开启DEBUG日志排查细节

如果上面的方法都没用,可以开Debezium的DEBUG日志,看看有没有隐藏的警告或错误。在Spring的application.properties里加:

logging.level.io.debezium=DEBUG

这样能看到连接器和数据库交互的详细过程,比如是否连接到CDC日志、有没有权限问题等。

内容的提问来源于stack exchange,提问作者frisk0k

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.05.08 11:27:57