【问题标题】:Embedded debezium doesnt capture changes嵌入式 debezium 不捕获更改
【发布时间】: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


    【解决方案1】:

    您不能离开try 块,因为这会关闭引擎并停止流式传输。因此,代替评论 // Do something else or wait for a signal or an event 必须是某种等待逻辑,否则您不应该将 engine 放入 try 中,并带有用于自动浸泡的资源。

    【讨论】:

      猜你喜欢
      • 2020-04-07
      • 2021-05-18
      • 2020-09-17
      • 2017-09-05
      • 2019-12-18
      • 2020-03-17
      • 1970-01-01
      • 2019-05-13
      • 2017-10-13
      相关资源
      最近更新 更多