【发布时间】:2020-10-19 10:42:12
【问题描述】:
我在 Spring 应用程序中运行嵌入式 Debezium (1.2.0),但它仅在启动时捕获更改
我的设置如下所示:
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());
}
当我启动应用程序时,它会捕获最新的更改,但以
结尾INFO i.d.p.ChangeEventSourceCoordinator - Finished streaming
INFO i.d.p.m.StreamingChangeEventSourceMetrics - Connected metrics set to 'false'
然后在下次应用程序重新启动之前不会捕获任何后续更改
不会抛出任何错误
【问题讨论】:
标签: spring-boot debezium