This is an automated email from the ASF dual-hosted git repository.

oscerd pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/camel.git


The following commit(s) were added to refs/heads/main by this push:
     new 2871cbdfc8e9 CAMEL-24889: camel-debezium - report a failed embedded 
engine (#26733)
2871cbdfc8e9 is described below

commit 2871cbdfc8e9f7d5e4bdcdfcbd2b96e7b1f59836
Author: Andrea Cosentino <[email protected]>
AuthorDate: Fri Sep 25 10:43:32 2026 +0200

    CAMEL-24889: camel-debezium - report a failed embedded engine (#26733)
    
    DebeziumEngine.create(Connect.class) resolves to an AsyncEmbeddedEngine, 
whose run() wraps its whole
    body in its own catch (Throwable), logs the failure and reports it through 
the CompletionCallback
    without ever rethrowing. The consumer registered no such callback, so its 
own catch around
    dbzEngine.run() was dead code and nothing in Camel learned that the engine 
had died: a connector that
    could not start left the route started and the health check UP while no 
change event was ever
    delivered again.
    
    Register a CompletionCallback that hands the failure to the consumer 
exception handler, and give the
    consumer a health check that reports DOWN carrying that error. Stopping a 
route whose engine has
    already failed no longer throws either, since close() rejects an engine 
that has already shut down and
    that used to skip the thread pool shutdown.
    
    The engine is still not restarted automatically; this only makes the 
failure observable.
    
    Closes #26733
    
    Co-authored-by: Claude Opus 5 <[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 +
 .../ROOT/pages/camel-4x-upgrade-guide-4_22.adoc    |  12 ++
 .../ROOT/pages/camel-4x-upgrade-guide-4_23.adoc    |  12 ++
 18 files changed, 390 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 0330fe1d62f3..bce0e3bf6f42 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 78a22dac3f4a..66908454ed35 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 7f07bb458ce5..e18e387d562e 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 8133cd569daa..e4a05e2b934d 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 fe88e3645c24..8b8dca6ad232 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 20a920fb0fc2..e7a48a170b6b 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 0330fe1d62f3..bce0e3bf6f42 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 78a22dac3f4a..66908454ed35 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 7f07bb458ce5..e18e387d562e 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 8133cd569daa..e4a05e2b934d 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 fe88e3645c24..8b8dca6ad232 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 20a920fb0fc2..e7a48a170b6b 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`
diff --git 
a/docs/user-manual/modules/ROOT/pages/camel-4x-upgrade-guide-4_22.adoc 
b/docs/user-manual/modules/ROOT/pages/camel-4x-upgrade-guide-4_22.adoc
index d429a7bbc767..025e8d4e3f68 100644
--- a/docs/user-manual/modules/ROOT/pages/camel-4x-upgrade-guide-4_22.adoc
+++ b/docs/user-manual/modules/ROOT/pages/camel-4x-upgrade-guide-4_22.adoc
@@ -49,6 +49,18 @@ are unaffected. Routes that set the header by its literal 
string name, or that u
 `allowTemplateFromHeader=true` with the old header names, must switch to the 
new `Camel`-prefixed
 names.
 
+=== camel-debezium - a failed embedded engine is now reported
+
+The Debezium consumers now register a `CompletionCallback` on the embedded 
engine. When the engine stops
+with an error - for example because the connector cannot reach the database, 
or the offset store cannot be
+read - the failure is passed to the consumer's `ExceptionHandler`, and the 
consumer's health check reports
+`DOWN` with that error.
+
+Previously the engine reported such failures only through its own logger, so 
the route stayed started and
+healthy while no longer receiving any change event. Deployments that use 
readiness or liveness probes will
+now see a Debezium route whose engine has died reported as `DOWN`, where it 
was previously reported as `UP`.
+The engine is still not restarted automatically.
+
 == Upgrading from 4.22.0 to 4.22.1
 
 === camel-docling
diff --git 
a/docs/user-manual/modules/ROOT/pages/camel-4x-upgrade-guide-4_23.adoc 
b/docs/user-manual/modules/ROOT/pages/camel-4x-upgrade-guide-4_23.adoc
index b3de381afe58..bfdb4e9a5260 100644
--- a/docs/user-manual/modules/ROOT/pages/camel-4x-upgrade-guide-4_23.adoc
+++ b/docs/user-manual/modules/ROOT/pages/camel-4x-upgrade-guide-4_23.adoc
@@ -2793,6 +2793,18 @@ are unaffected. Routes that set the header by its 
literal string name, or that u
 `allowTemplateFromHeader=true` with the old header names, must switch to the 
new `Camel`-prefixed
 names.
 
+=== camel-debezium - a failed embedded engine is now reported
+
+The Debezium consumers now register a `CompletionCallback` on the embedded 
engine. When the engine stops
+with an error - for example because the connector cannot reach the database, 
or the offset store cannot be
+read - the failure is passed to the consumer's `ExceptionHandler`, and the 
consumer's health check reports
+`DOWN` with that error.
+
+Previously the engine reported such failures only through its own logger, so 
the route stayed started and
+healthy while no longer receiving any change event. Deployments that use 
readiness or liveness probes will
+now see a Debezium route whose engine has died reported as `DOWN`, where it 
was previously reported as `UP`.
+The engine is still not restarted automatically.
+
 === camel-seda - purging the queue completes the discarded exchanges
 
 When a SEDA queue is purged (with `purgeWhenStopping=true` or the `purgeQueue` 
JMX operation), the discarded exchanges

Reply via email to