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`

Reply via email to