嵌入式Debezium(1.2.0)仅启动时捕获变更,后续无捕获
我之前也碰到过类似的情况,结合你的配置和日志信息,咱们一步步拆解问题:
先还原你的场景:你在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

