This is an automated email from the ASF dual-hosted git repository.
oscerd pushed a commit to branch camel-4.22.x
in repository https://gitbox.apache.org/repos/asf/camel.git
The following commit(s) were added to refs/heads/camel-4.22.x by this push:
new 7579dcbb8199 [backport camel-4.22.x] CAMEL-24889: camel-debezium -
report a failed embedded engine (#27075)
7579dcbb8199 is described below
commit 7579dcbb819914d5ba278747ffe5f7e95935d6ae
Author: Andrea Cosentino <[email protected]>
AuthorDate: Tue Sep 29 13:06:25 2026 +0200
[backport camel-4.22.x] CAMEL-24889: camel-debezium - report a failed
embedded engine (#27075)
DebeziumEngine.create(Connect.class) builds an AsyncEmbeddedEngine, whose
run()
catches Throwable itself and reports through the CompletionCallback without
rethrowing. The consumer now registers a CompletionCallback, so a failed
engine
is reported and the route's health reflects it.
Backport of #26733 (cherry-picked from 2871cbdfc8e9 on main), without the
upgrade-guide edits: the per-release guides for every line live on main.
Co-Authored-By: Claude Opus 5 <[email protected]>
Signed-off-by: Andrea Cosentino <[email protected]>
---
.../camel/catalog/docs/debezium-db2-component.adoc | 5 +
.../catalog/docs/debezium-mongodb-component.adoc | 5 +
.../catalog/docs/debezium-mysql-component.adoc | 5 +
.../catalog/docs/debezium-oracle-component.adoc | 5 +
.../catalog/docs/debezium-postgres-component.adoc | 5 +
.../catalog/docs/debezium-sqlserver-component.adoc | 5 +
.../camel/component/debezium/DebeziumConsumer.java | 65 ++++++++-
.../debezium/DebeziumConsumerHealthCheck.java | 90 +++++++++++++
.../DebeziumConsumerEngineFailureTest.java | 147 +++++++++++++++++++++
.../component/debezium/DebeziumConsumerTest.java | 7 +
.../src/main/docs/debezium-db2-component.adoc | 5 +
.../src/main/docs/debezium-mongodb-component.adoc | 5 +
.../src/main/docs/debezium-mysql-component.adoc | 5 +
.../src/main/docs/debezium-oracle-component.adoc | 5 +
.../src/main/docs/debezium-postgres-component.adoc | 5 +
.../main/docs/debezium-sqlserver-component.adoc | 5 +
16 files changed, 366 insertions(+), 3 deletions(-)
diff --git
a/catalog/camel-catalog/src/generated/resources/org/apache/camel/catalog/docs/debezium-db2-component.adoc
b/catalog/camel-catalog/src/generated/resources/org/apache/camel/catalog/docs/debezium-db2-component.adoc
index 31542d6eb539..b4222c4c2bff 100644
---
a/catalog/camel-catalog/src/generated/resources/org/apache/camel/catalog/docs/debezium-db2-component.adoc
+++
b/catalog/camel-catalog/src/generated/resources/org/apache/camel/catalog/docs/debezium-db2-component.adoc
@@ -28,6 +28,11 @@ However, in case of an application crash (not having a
graceful shutdown),
the application will resume from the last recorded offset,
which may result in receiving duplicate events immediately after the restart.
Therefore, your downstream routes should be tolerant enough of such a case and
deduplicate events if needed.
+
+If the embedded engine stops with an error, for example because the connector
cannot reach the database or
+the offset store cannot be read, it is not restarted: the route stays started
but receives no further events.
+The failure is reported to the consumer exception handler and turns the route
health check `DOWN`, so it can
+be picked up by a readiness probe.
====
diff --git
a/catalog/camel-catalog/src/generated/resources/org/apache/camel/catalog/docs/debezium-mongodb-component.adoc
b/catalog/camel-catalog/src/generated/resources/org/apache/camel/catalog/docs/debezium-mongodb-component.adoc
index c68c362eed8b..89d635689338 100644
---
a/catalog/camel-catalog/src/generated/resources/org/apache/camel/catalog/docs/debezium-mongodb-component.adoc
+++
b/catalog/camel-catalog/src/generated/resources/org/apache/camel/catalog/docs/debezium-mongodb-component.adoc
@@ -34,6 +34,11 @@ However, in case of an application crash (not having a
graceful shutdown),
the application will resume from the last recorded offset,
which may result in receiving duplicate events immediately after the restart.
Therefore, your downstream routes should be tolerant enough of such a case and
deduplicate events if needed.
+
+If the embedded engine stops with an error, for example because the connector
cannot reach the database or
+the offset store cannot be read, it is not restarted: the route stays started
but receives no further events.
+The failure is reported to the consumer exception handler and turns the route
health check `DOWN`, so it can
+be picked up by a readiness probe.
====
Maven users will need to add the following dependency to their `pom.xml`
for this component.
diff --git
a/catalog/camel-catalog/src/generated/resources/org/apache/camel/catalog/docs/debezium-mysql-component.adoc
b/catalog/camel-catalog/src/generated/resources/org/apache/camel/catalog/docs/debezium-mysql-component.adoc
index 50e8d8f5e33f..244cab53a66b 100644
---
a/catalog/camel-catalog/src/generated/resources/org/apache/camel/catalog/docs/debezium-mysql-component.adoc
+++
b/catalog/camel-catalog/src/generated/resources/org/apache/camel/catalog/docs/debezium-mysql-component.adoc
@@ -27,6 +27,11 @@ However, in case of an application crash (not having a
graceful shutdown),
the application will resume from the last recorded offset,
which may result in receiving duplicate events immediately after the restart.
Therefore, your downstream routes should be tolerant enough of such a case and
deduplicate events if needed.
+
+If the embedded engine stops with an error, for example because the connector
cannot reach the database or
+the offset store cannot be read, it is not restarted: the route stays started
but receives no further events.
+The failure is reported to the consumer exception handler and turns the route
health check `DOWN`, so it can
+be picked up by a readiness probe.
====
Maven users will need to add the following dependency to their `pom.xml`
diff --git
a/catalog/camel-catalog/src/generated/resources/org/apache/camel/catalog/docs/debezium-oracle-component.adoc
b/catalog/camel-catalog/src/generated/resources/org/apache/camel/catalog/docs/debezium-oracle-component.adoc
index bfc2806cc0e5..0b3226534de8 100644
---
a/catalog/camel-catalog/src/generated/resources/org/apache/camel/catalog/docs/debezium-oracle-component.adoc
+++
b/catalog/camel-catalog/src/generated/resources/org/apache/camel/catalog/docs/debezium-oracle-component.adoc
@@ -28,6 +28,11 @@ However, in case of an application crash (not having a
graceful shutdown),
the application will resume from the last recorded offset,
which may result in receiving duplicate events immediately after the restart.
Therefore, your downstream routes should be tolerant enough of such a case and
deduplicate events if needed.
+
+If the embedded engine stops with an error, for example because the connector
cannot reach the database or
+the offset store cannot be read, it is not restarted: the route stays started
but receives no further events.
+The failure is reported to the consumer exception handler and turns the route
health check `DOWN`, so it can
+be picked up by a readiness probe.
====
Maven users will need to add the following dependency to their `pom.xml` for
this component.
diff --git
a/catalog/camel-catalog/src/generated/resources/org/apache/camel/catalog/docs/debezium-postgres-component.adoc
b/catalog/camel-catalog/src/generated/resources/org/apache/camel/catalog/docs/debezium-postgres-component.adoc
index 83f63cd29348..9dba1a052907 100644
---
a/catalog/camel-catalog/src/generated/resources/org/apache/camel/catalog/docs/debezium-postgres-component.adoc
+++
b/catalog/camel-catalog/src/generated/resources/org/apache/camel/catalog/docs/debezium-postgres-component.adoc
@@ -28,6 +28,11 @@ However, in case of an application crash (not having a
graceful shutdown),
the application will resume from the last recorded offset,
which may result in receiving duplicate events immediately after the restart.
Therefore, your downstream routes should be tolerant enough of such a case and
deduplicate events if needed.
+
+If the embedded engine stops with an error, for example because the connector
cannot reach the database or
+the offset store cannot be read, it is not restarted: the route stays started
but receives no further events.
+The failure is reported to the consumer exception handler and turns the route
health check `DOWN`, so it can
+be picked up by a readiness probe.
====
Maven users will need to add the following dependency to their `pom.xml` for
this component.
diff --git
a/catalog/camel-catalog/src/generated/resources/org/apache/camel/catalog/docs/debezium-sqlserver-component.adoc
b/catalog/camel-catalog/src/generated/resources/org/apache/camel/catalog/docs/debezium-sqlserver-component.adoc
index a2857dcbc154..994de4c39367 100644
---
a/catalog/camel-catalog/src/generated/resources/org/apache/camel/catalog/docs/debezium-sqlserver-component.adoc
+++
b/catalog/camel-catalog/src/generated/resources/org/apache/camel/catalog/docs/debezium-sqlserver-component.adoc
@@ -27,6 +27,11 @@ However, in case of an application crash (not having a
graceful shutdown),
the application will resume from the last recorded offset,
which may result in receiving duplicate events immediately after the restart.
Therefore, your downstream routes should be tolerant enough of such a case and
deduplicate events if needed.
+
+If the embedded engine stops with an error, for example because the connector
cannot reach the database or
+the offset store cannot be read, it is not restarted: the route stays started
but receives no further events.
+The failure is reported to the consumer exception handler and turns the route
health check `DOWN`, so it can
+be picked up by a readiness probe.
====
Maven users will need to add the following dependency to their `pom.xml`
diff --git
a/components/camel-debezium/camel-debezium-common/camel-debezium-common-component/src/main/java/org/apache/camel/component/debezium/DebeziumConsumer.java
b/components/camel-debezium/camel-debezium-common/camel-debezium-common-component/src/main/java/org/apache/camel/component/debezium/DebeziumConsumer.java
index 8a89c7b94f31..ebcf9ed3b996 100644
---
a/components/camel-debezium/camel-debezium-common/camel-debezium-common-component/src/main/java/org/apache/camel/component/debezium/DebeziumConsumer.java
+++
b/components/camel-debezium/camel-debezium-common/camel-debezium-common-component/src/main/java/org/apache/camel/component/debezium/DebeziumConsumer.java
@@ -23,8 +23,10 @@ import io.debezium.engine.ChangeEvent;
import io.debezium.engine.DebeziumEngine;
import org.apache.camel.Exchange;
import org.apache.camel.Processor;
+import org.apache.camel.RuntimeCamelException;
import
org.apache.camel.component.debezium.configuration.EmbeddedDebeziumConfiguration;
import org.apache.camel.support.DefaultConsumer;
+import org.apache.camel.util.URISupport;
import org.apache.kafka.connect.source.SourceRecord;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
@@ -38,6 +40,8 @@ public class DebeziumConsumer extends DefaultConsumer {
private ExecutorService executorService;
private DebeziumEngine<ChangeEvent<SourceRecord, SourceRecord>> dbzEngine;
+ private volatile Throwable engineFailure;
+ private volatile boolean engineStopped;
public DebeziumConsumer(DebeziumEndpoint endpoint, Processor processor) {
super(endpoint, processor);
@@ -45,10 +49,21 @@ public class DebeziumConsumer extends DefaultConsumer {
this.configuration = endpoint.getConfiguration();
}
+ @Override
+ protected void doBuild() throws Exception {
+ if (getHealthCheck() == null) {
+ setHealthCheck(new DebeziumConsumerHealthCheck(this, "consumer:" +
getRouteId()));
+ }
+ super.doBuild();
+ }
+
@Override
protected void doStart() throws Exception {
super.doStart();
+ engineFailure = null;
+ engineStopped = false;
+
// start a single threaded pool to monitor events
executorService = endpoint.createExecutor(this);
@@ -61,15 +76,28 @@ public class DebeziumConsumer extends DefaultConsumer {
try {
dbzEngine.run();
} catch (Throwable e) {
- LOG.error("Debezium engine has failed: {}",
e.getMessage(), e);
+ // the engine reports its own failures through the
completion callback and is not
+ // expected to throw, so this is only a safety net
+ onEngineCompleted(false, e.getMessage(), e);
}
});
}
@Override
protected void doStop() throws Exception {
- if (dbzEngine != null) {
- dbzEngine.close();
+ if (dbzEngine != null && !engineStopped) {
+ try {
+ dbzEngine.close();
+ } catch (IllegalStateException e) {
+ // close() rejects three states: the engine stopped on its own
between the check above and
+ // this call, which leaves nothing to close, but also tasks
still starting and a shutdown
+ // already in progress, and those two are worth seeing
+ if (engineStopped) {
+ LOG.debug("Debezium engine was already stopped: {}",
e.getMessage());
+ } else {
+ LOG.warn("Debezium engine could not be closed: {}",
e.getMessage());
+ }
+ }
}
// shutdown the thread pool gracefully
@@ -79,13 +107,44 @@ public class DebeziumConsumer extends DefaultConsumer {
super.doStop();
}
+ /**
+ * The failure that stopped the embedded engine, or <tt>null</tt> while
the engine is starting, running, or was
+ * stopped on request. Used by {@link DebeziumConsumerHealthCheck}.
+ */
+ Throwable getEngineFailure() {
+ return engineFailure;
+ }
+
private DebeziumEngine<ChangeEvent<SourceRecord, SourceRecord>>
createDbzEngine() {
return DebeziumEngine.create(Connect.class)
.using(configuration.createDebeziumConfiguration().asProperties())
+ .using(this::onEngineCompleted)
.notifying(this::onEventListener)
.build();
}
+ /**
+ * Called by the embedded engine once it has stopped, either because it
was closed or because it failed. The engine
+ * does not restart itself, so a failure means this consumer is
permanently dead while the route still reports as
+ * started, hence the failure is reported to the exception handler and
kept for the health check.
+ */
+ private void onEngineCompleted(boolean success, String message, Throwable
error) {
+ engineStopped = true;
+
+ if (success) {
+ LOG.debug("Debezium engine stopped: {}", message);
+ return;
+ }
+
+ final Throwable cause = error != null ? error : new
RuntimeCamelException(message);
+ engineFailure = cause;
+
+ getExceptionHandler().handleException(
+ "Debezium engine has failed and stopped, no more change events
are consumed from "
+ +
URISupport.sanitizeUri(endpoint.getEndpointUri()),
+ cause);
+ }
+
private void onEventListener(final ChangeEvent<SourceRecord, SourceRecord>
event) {
final Exchange exchange = endpoint.createDbzExchange(this,
event.value());
diff --git
a/components/camel-debezium/camel-debezium-common/camel-debezium-common-component/src/main/java/org/apache/camel/component/debezium/DebeziumConsumerHealthCheck.java
b/components/camel-debezium/camel-debezium-common/camel-debezium-common-component/src/main/java/org/apache/camel/component/debezium/DebeziumConsumerHealthCheck.java
new file mode 100644
index 000000000000..4d4ff2e08038
--- /dev/null
+++
b/components/camel-debezium/camel-debezium-common/camel-debezium-common-component/src/main/java/org/apache/camel/component/debezium/DebeziumConsumerHealthCheck.java
@@ -0,0 +1,90 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements. See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+package org.apache.camel.component.debezium;
+
+import java.util.Map;
+
+import org.apache.camel.health.HealthCheck;
+import org.apache.camel.health.HealthCheckResultBuilder;
+import org.apache.camel.util.URISupport;
+
+/**
+ * {@link HealthCheck} reporting the state of the embedded Debezium engine
that backs a {@link DebeziumConsumer}.
+ * <p>
+ * The engine runs on its own thread and does not restart itself, so once it
has stopped with an error the route no
+ * longer receives change events even though it is still started. This check
turns the route DOWN in that case.
+ */
+public class DebeziumConsumerHealthCheck implements HealthCheck {
+
+ private final DebeziumConsumer consumer;
+ private final String id;
+ private final String sanitizedUri;
+ private boolean enabled = true;
+
+ public DebeziumConsumerHealthCheck(DebeziumConsumer consumer, String id) {
+ this.consumer = consumer;
+ this.id = id;
+ this.sanitizedUri =
URISupport.sanitizeUri(consumer.getEndpoint().getEndpointUri());
+ }
+
+ @Override
+ public boolean isEnabled() {
+ return enabled;
+ }
+
+ @Override
+ public void setEnabled(boolean enabled) {
+ this.enabled = enabled;
+ }
+
+ @Override
+ public String getGroup() {
+ return "camel";
+ }
+
+ @Override
+ public String getId() {
+ return id;
+ }
+
+ @Override
+ public Result call(Map<String, Object> options) {
+ final HealthCheckResultBuilder builder =
HealthCheckResultBuilder.on(this);
+
+ // ensure to sanitize uri, so we do not show sensitive information
such as passwords
+ builder.detail(ENDPOINT_URI, sanitizedUri);
+
+ if (!isEnabled()) {
+ builder.message("Disabled");
+ builder.detail(CHECK_ENABLED, false);
+ return builder.unknown().build();
+ }
+
+ final Throwable failure = consumer.getEngineFailure();
+ if (failure == null) {
+ // the engine is either starting, running, or was stopped on
request, none of which is a failure
+ return builder.up().build();
+ }
+
+ builder.down();
+ builder.message(String.format(
+ "Debezium engine stopped after a failure, route: %s (%s) no
longer consumes change events",
+ consumer.getRouteId(), sanitizedUri));
+ builder.error(failure);
+ return builder.build();
+ }
+}
diff --git
a/components/camel-debezium/camel-debezium-common/camel-debezium-common-component/src/test/java/org/apache/camel/component/debezium/DebeziumConsumerEngineFailureTest.java
b/components/camel-debezium/camel-debezium-common/camel-debezium-common-component/src/test/java/org/apache/camel/component/debezium/DebeziumConsumerEngineFailureTest.java
new file mode 100644
index 000000000000..bd328327f5b3
--- /dev/null
+++
b/components/camel-debezium/camel-debezium-common/camel-debezium-common-component/src/test/java/org/apache/camel/component/debezium/DebeziumConsumerEngineFailureTest.java
@@ -0,0 +1,147 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements. See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+package org.apache.camel.component.debezium;
+
+import java.io.IOException;
+import java.nio.charset.StandardCharsets;
+import java.nio.file.Files;
+import java.nio.file.Path;
+import java.nio.file.Paths;
+import java.util.concurrent.CountDownLatch;
+import java.util.concurrent.TimeUnit;
+import java.util.concurrent.atomic.AtomicReference;
+
+import io.debezium.util.IoUtil;
+import org.apache.camel.CamelContext;
+import org.apache.camel.Exchange;
+import org.apache.camel.RoutesBuilder;
+import org.apache.camel.builder.RouteBuilder;
+import
org.apache.camel.component.debezium.configuration.FileConnectorEmbeddedDebeziumConfiguration;
+import org.apache.camel.health.HealthCheck;
+import org.apache.camel.spi.ExceptionHandler;
+import org.apache.camel.test.junit6.CamelTestSupport;
+import org.junit.jupiter.api.AfterEach;
+import org.junit.jupiter.api.BeforeEach;
+import org.junit.jupiter.api.Test;
+
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertNotNull;
+import static org.junit.jupiter.api.Assertions.assertTrue;
+
+/**
+ * The embedded engine runs on its own thread and reports a failure only
through its completion callback, so without
+ * that callback a connector that cannot start leaves the route started,
healthy and silent.
+ */
+public class DebeziumConsumerEngineFailureTest extends CamelTestSupport {
+
+ private static final String ROUTE_ID = "debezium-failing-engine";
+ private static final Path TEST_FILE_PATH
+ = Paths.get("target/data",
"camel-debezium-engine-failure-input.txt").toAbsolutePath();
+ private static final Path TEST_OFFSET_STORE_PATH
+ = Paths.get("target/data",
"camel-debezium-engine-failure-offset-store.txt").toAbsolutePath();
+
+ @BeforeEach
+ public void beforeEach() throws IOException {
+ IoUtil.createFile(TEST_FILE_PATH);
+ // an offset store the engine cannot read, which is what a corrupted
offset file looks like; the
+ // content must be long enough not to be mistaken for an empty store
+ Files.write(IoUtil.createFile(TEST_OFFSET_STORE_PATH).toPath(),
+ "this is not a serialized offset
store".getBytes(StandardCharsets.UTF_8));
+ }
+
+ @AfterEach
+ public void afterEach() throws IOException {
+ IoUtil.delete(TEST_FILE_PATH);
+ IoUtil.delete(TEST_OFFSET_STORE_PATH);
+ }
+
+ @Test
+ void engineFailureIsReportedToTheExceptionHandlerAndTheHealthCheck()
throws Exception {
+ final DebeziumConsumer consumer = (DebeziumConsumer)
context.getRoute(ROUTE_ID).getConsumer();
+
+ final CapturingExceptionHandler exceptionHandler = new
CapturingExceptionHandler();
+ consumer.setExceptionHandler(exceptionHandler);
+
+ context.getRouteController().startRoute(ROUTE_ID);
+
+ assertTrue(exceptionHandler.latch.await(30, TimeUnit.SECONDS),
+ "the engine failure should be handed to the consumer exception
handler");
+ assertNotNull(exceptionHandler.captured.get(), "the failure cause
should be reported");
+
+ final HealthCheck healthCheck = consumer.getHealthCheck();
+ assertNotNull(healthCheck, "the consumer should expose a health
check");
+ assertEquals(HealthCheck.State.DOWN, healthCheck.call().getState(),
+ "a consumer whose engine has died should not report as
healthy");
+ }
+
+ @Override
+ protected CamelContext createCamelContext() throws Exception {
+ final CamelContext context = super.createCamelContext();
+ final DebeziumTestComponent component = new
DebeziumTestComponent(context);
+
+ component.setConfiguration(initConfiguration());
+ context.addComponent("debezium", component);
+ context.disableJMX();
+
+ return context;
+ }
+
+ @Override
+ protected RoutesBuilder createRouteBuilder() {
+ return new RouteBuilder() {
+ @Override
+ public void configure() {
+ // started by the test, so that the exception handler is in
place before the engine runs
+ from("debezium:dummy").routeId(ROUTE_ID).autoStartup(false)
+ .to("mock:result");
+ }
+ };
+ }
+
+ private FileConnectorEmbeddedDebeziumConfiguration initConfiguration() {
+ final FileConnectorEmbeddedDebeziumConfiguration configuration = new
FileConnectorEmbeddedDebeziumConfiguration();
+ configuration.setName("engine_failure_dummy");
+ configuration.setTopicConfig("engine_failure_topic");
+ configuration.setTestFilePath(TEST_FILE_PATH);
+
configuration.setOffsetStorageFileName(TEST_OFFSET_STORE_PATH.toString());
+ configuration.setOffsetFlushIntervalMs(0);
+
+ return configuration;
+ }
+
+ private static final class CapturingExceptionHandler implements
ExceptionHandler {
+
+ private final CountDownLatch latch = new CountDownLatch(1);
+ private final AtomicReference<Throwable> captured = new
AtomicReference<>();
+
+ @Override
+ public void handleException(Throwable exception) {
+ handleException(null, null, exception);
+ }
+
+ @Override
+ public void handleException(String message, Throwable exception) {
+ handleException(message, null, exception);
+ }
+
+ @Override
+ public void handleException(String message, Exchange exchange,
Throwable exception) {
+ captured.compareAndSet(null, exception);
+ latch.countDown();
+ }
+ }
+}
diff --git
a/components/camel-debezium/camel-debezium-common/camel-debezium-common-component/src/test/java/org/apache/camel/component/debezium/DebeziumConsumerTest.java
b/components/camel-debezium/camel-debezium-common/camel-debezium-common-component/src/test/java/org/apache/camel/component/debezium/DebeziumConsumerTest.java
index e648941b47f9..bffde43c2c4e 100644
---
a/components/camel-debezium/camel-debezium-common/camel-debezium-common-component/src/test/java/org/apache/camel/component/debezium/DebeziumConsumerTest.java
+++
b/components/camel-debezium/camel-debezium-common/camel-debezium-common-component/src/test/java/org/apache/camel/component/debezium/DebeziumConsumerTest.java
@@ -33,6 +33,7 @@ import org.apache.camel.builder.RouteBuilder;
import
org.apache.camel.component.debezium.configuration.EmbeddedDebeziumConfiguration;
import
org.apache.camel.component.debezium.configuration.FileConnectorEmbeddedDebeziumConfiguration;
import org.apache.camel.component.mock.MockEndpoint;
+import org.apache.camel.health.HealthCheck;
import org.apache.camel.test.junit6.CamelTestSupport;
import org.junit.jupiter.api.AfterAll;
import org.junit.jupiter.api.AfterEach;
@@ -41,6 +42,8 @@ import org.junit.jupiter.api.Test;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
+import static org.junit.jupiter.api.Assertions.assertEquals;
+
public class DebeziumConsumerTest extends CamelTestSupport {
private static final Logger LOG =
LoggerFactory.getLogger(DebeziumConsumerTest.class);
@@ -93,6 +96,10 @@ public class DebeziumConsumerTest extends CamelTestSupport {
// verify the first records if they being consumed
to.assertIsSatisfied(50);
+ // a consumer whose engine is running must report as healthy
+ final DebeziumConsumer consumer = (DebeziumConsumer)
context.getRoutes().get(0).getConsumer();
+ assertEquals(HealthCheck.State.UP,
consumer.getHealthCheck().call().getState());
+
// send another batch
appendLinesToSource(NUMBER_OF_LINES);
diff --git
a/components/camel-debezium/camel-debezium-db2/src/main/docs/debezium-db2-component.adoc
b/components/camel-debezium/camel-debezium-db2/src/main/docs/debezium-db2-component.adoc
index 31542d6eb539..b4222c4c2bff 100644
---
a/components/camel-debezium/camel-debezium-db2/src/main/docs/debezium-db2-component.adoc
+++
b/components/camel-debezium/camel-debezium-db2/src/main/docs/debezium-db2-component.adoc
@@ -28,6 +28,11 @@ However, in case of an application crash (not having a
graceful shutdown),
the application will resume from the last recorded offset,
which may result in receiving duplicate events immediately after the restart.
Therefore, your downstream routes should be tolerant enough of such a case and
deduplicate events if needed.
+
+If the embedded engine stops with an error, for example because the connector
cannot reach the database or
+the offset store cannot be read, it is not restarted: the route stays started
but receives no further events.
+The failure is reported to the consumer exception handler and turns the route
health check `DOWN`, so it can
+be picked up by a readiness probe.
====
diff --git
a/components/camel-debezium/camel-debezium-mongodb/src/main/docs/debezium-mongodb-component.adoc
b/components/camel-debezium/camel-debezium-mongodb/src/main/docs/debezium-mongodb-component.adoc
index c68c362eed8b..89d635689338 100644
---
a/components/camel-debezium/camel-debezium-mongodb/src/main/docs/debezium-mongodb-component.adoc
+++
b/components/camel-debezium/camel-debezium-mongodb/src/main/docs/debezium-mongodb-component.adoc
@@ -34,6 +34,11 @@ However, in case of an application crash (not having a
graceful shutdown),
the application will resume from the last recorded offset,
which may result in receiving duplicate events immediately after the restart.
Therefore, your downstream routes should be tolerant enough of such a case and
deduplicate events if needed.
+
+If the embedded engine stops with an error, for example because the connector
cannot reach the database or
+the offset store cannot be read, it is not restarted: the route stays started
but receives no further events.
+The failure is reported to the consumer exception handler and turns the route
health check `DOWN`, so it can
+be picked up by a readiness probe.
====
Maven users will need to add the following dependency to their `pom.xml`
for this component.
diff --git
a/components/camel-debezium/camel-debezium-mysql/src/main/docs/debezium-mysql-component.adoc
b/components/camel-debezium/camel-debezium-mysql/src/main/docs/debezium-mysql-component.adoc
index 50e8d8f5e33f..244cab53a66b 100644
---
a/components/camel-debezium/camel-debezium-mysql/src/main/docs/debezium-mysql-component.adoc
+++
b/components/camel-debezium/camel-debezium-mysql/src/main/docs/debezium-mysql-component.adoc
@@ -27,6 +27,11 @@ However, in case of an application crash (not having a
graceful shutdown),
the application will resume from the last recorded offset,
which may result in receiving duplicate events immediately after the restart.
Therefore, your downstream routes should be tolerant enough of such a case and
deduplicate events if needed.
+
+If the embedded engine stops with an error, for example because the connector
cannot reach the database or
+the offset store cannot be read, it is not restarted: the route stays started
but receives no further events.
+The failure is reported to the consumer exception handler and turns the route
health check `DOWN`, so it can
+be picked up by a readiness probe.
====
Maven users will need to add the following dependency to their `pom.xml`
diff --git
a/components/camel-debezium/camel-debezium-oracle/src/main/docs/debezium-oracle-component.adoc
b/components/camel-debezium/camel-debezium-oracle/src/main/docs/debezium-oracle-component.adoc
index bfc2806cc0e5..0b3226534de8 100644
---
a/components/camel-debezium/camel-debezium-oracle/src/main/docs/debezium-oracle-component.adoc
+++
b/components/camel-debezium/camel-debezium-oracle/src/main/docs/debezium-oracle-component.adoc
@@ -28,6 +28,11 @@ However, in case of an application crash (not having a
graceful shutdown),
the application will resume from the last recorded offset,
which may result in receiving duplicate events immediately after the restart.
Therefore, your downstream routes should be tolerant enough of such a case and
deduplicate events if needed.
+
+If the embedded engine stops with an error, for example because the connector
cannot reach the database or
+the offset store cannot be read, it is not restarted: the route stays started
but receives no further events.
+The failure is reported to the consumer exception handler and turns the route
health check `DOWN`, so it can
+be picked up by a readiness probe.
====
Maven users will need to add the following dependency to their `pom.xml` for
this component.
diff --git
a/components/camel-debezium/camel-debezium-postgres/src/main/docs/debezium-postgres-component.adoc
b/components/camel-debezium/camel-debezium-postgres/src/main/docs/debezium-postgres-component.adoc
index 83f63cd29348..9dba1a052907 100644
---
a/components/camel-debezium/camel-debezium-postgres/src/main/docs/debezium-postgres-component.adoc
+++
b/components/camel-debezium/camel-debezium-postgres/src/main/docs/debezium-postgres-component.adoc
@@ -28,6 +28,11 @@ However, in case of an application crash (not having a
graceful shutdown),
the application will resume from the last recorded offset,
which may result in receiving duplicate events immediately after the restart.
Therefore, your downstream routes should be tolerant enough of such a case and
deduplicate events if needed.
+
+If the embedded engine stops with an error, for example because the connector
cannot reach the database or
+the offset store cannot be read, it is not restarted: the route stays started
but receives no further events.
+The failure is reported to the consumer exception handler and turns the route
health check `DOWN`, so it can
+be picked up by a readiness probe.
====
Maven users will need to add the following dependency to their `pom.xml` for
this component.
diff --git
a/components/camel-debezium/camel-debezium-sqlserver/src/main/docs/debezium-sqlserver-component.adoc
b/components/camel-debezium/camel-debezium-sqlserver/src/main/docs/debezium-sqlserver-component.adoc
index a2857dcbc154..994de4c39367 100644
---
a/components/camel-debezium/camel-debezium-sqlserver/src/main/docs/debezium-sqlserver-component.adoc
+++
b/components/camel-debezium/camel-debezium-sqlserver/src/main/docs/debezium-sqlserver-component.adoc
@@ -27,6 +27,11 @@ However, in case of an application crash (not having a
graceful shutdown),
the application will resume from the last recorded offset,
which may result in receiving duplicate events immediately after the restart.
Therefore, your downstream routes should be tolerant enough of such a case and
deduplicate events if needed.
+
+If the embedded engine stops with an error, for example because the connector
cannot reach the database or
+the offset store cannot be read, it is not restarted: the route stays started
but receives no further events.
+The failure is reported to the consumer exception handler and turns the route
health check `DOWN`, so it can
+be picked up by a readiness probe.
====
Maven users will need to add the following dependency to their `pom.xml`