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 927b570d022d CAMEL-25024, CAMEL-25025: camel-mongodb - consumer 
failure handling and the change stream document key (#26961)
927b570d022d is described below

commit 927b570d022d0383678050d9a82b274d2457db97
Author: Andrea Cosentino <[email protected]>
AuthorDate: Mon Sep 28 21:44:06 2026 +0200

    CAMEL-25024, CAMEL-25025: camel-mongodb - consumer failure handling and the 
change stream document key (#26961)
    
    CAMEL-25024: both MongoDB consumers ignored the outcome of the route. A 
failed
    route leaves its exception on the exchange rather than throwing from 
process(),
    so the failure never reached the consumer's ExceptionHandler, and both 
consumers
    moved their position past a failed event. Both threads now route through
    processExchange(), which checks exchange.getException(), reports a failure 
to the
    ExceptionHandler, and moves the tail tracking value or the change stream 
resume
    token only on success. Exchanges are created with autoRelease=false and 
released
    once the outcome has been read.
    
    CAMEL-25025: the change streams consumer read the document key's _id 
blindly as an
    ObjectId, which threw for a string, numeric or compound _id and for events 
with no
    document key, before the exchange was created, so the regenerated cursor 
returned
    the same event forever. The key is now read defensively and the header 
keeps the
    id's natural Java type (javaType Object).
    
    Upgrade guide entries for both. ITs with a failing route for both 
consumers, and
    unit tests for the document key.
    
    Co-Authored-By: Claude Opus 5.5 (1M context) <[email protected]>
    Signed-off-by: Andrea Cosentino <[email protected]>
---
 .../apache/camel/catalog/components/mongodb.json   |  2 +-
 .../apache/camel/component/mongodb/mongodb.json    |  2 +-
 .../mongodb/MongoAbstractConsumerThread.java       | 38 +++++++++-
 .../mongodb/MongoDbChangeStreamsThread.java        | 65 ++++++++++++++---
 .../camel/component/mongodb/MongoDbConstants.java  | 10 +--
 .../component/mongodb/MongoDbTailingThread.java    | 13 ++--
 .../mongodb/MongoDbChangeStreamsThreadTest.java    | 84 ++++++++++++++++++++++
 .../MongoDbChangeStreamsConsumerIT.java            | 42 +++++++++++
 .../MongoDbTailableCursorConsumerIT.java           | 37 ++++++++++
 .../integration/RecordingExceptionHandler.java     | 62 ++++++++++++++++
 .../ROOT/pages/camel-4x-upgrade-guide-4_23.adoc    | 31 ++++++++
 .../dsl/MongoDbEndpointBuilderFactory.java         | 14 ++--
 12 files changed, 371 insertions(+), 29 deletions(-)

diff --git 
a/catalog/camel-catalog/src/generated/resources/org/apache/camel/catalog/components/mongodb.json
 
b/catalog/camel-catalog/src/generated/resources/org/apache/camel/catalog/components/mongodb.json
index e07db0812e4f..0118ae0c8885 100644
--- 
a/catalog/camel-catalog/src/generated/resources/org/apache/camel/catalog/components/mongodb.json
+++ 
b/catalog/camel-catalog/src/generated/resources/org/apache/camel/catalog/components/mongodb.json
@@ -53,7 +53,7 @@
     "CamelMongoDbDistinctQueryField": { "index": 19, "kind": "header", 
"displayName": "", "group": "producer", "label": "producer", "required": false, 
"javaType": "String", "deprecated": false, "deprecationNote": "", "autowired": 
false, "secret": false, "description": "The specified field name fow which we 
want to get the distinct values.", "constantName": 
"org.apache.camel.component.mongodb.MongoDbConstants#DISTINCT_QUERY_FIELD" },
     "CamelMongoDbAllowDiskUse": { "index": 20, "kind": "header", 
"displayName": "", "group": "producer findAll aggregate", "label": "producer 
findAll aggregate", "required": false, "javaType": "Boolean", "deprecated": 
false, "deprecationNote": "", "autowired": false, "secret": false, 
"description": "Sets allowDiskUse MongoDB flag. This is supported since MongoDB 
Server 4.3.1. Using this header with older MongoDB Server version can cause 
query to fail.", "constantName": "org.apache.camel. [...]
     "CamelMongoDbBulkOrdered": { "index": 21, "kind": "header", "displayName": 
"", "group": "producer bulkWrite", "label": "producer bulkWrite", "required": 
false, "javaType": "Boolean", "deprecated": false, "deprecationNote": "", 
"autowired": false, "secret": false, "defaultValue": "TRUE", "description": 
"Perform an ordered or unordered operation execution.", "constantName": 
"org.apache.camel.component.mongodb.MongoDbConstants#BULK_ORDERED" },
-    "_id": { "index": 22, "kind": "header", "displayName": "", "group": 
"consumer changeStreams", "label": "consumer changeStreams", "required": false, 
"javaType": "org.bson.types.ObjectId", "deprecated": false, "deprecationNote": 
"", "autowired": false, "secret": false, "description": "A document that 
contains the _id of the document created or modified by the insert, replace, 
delete, update operations (i.e. CRUD operations). For sharded collections, also 
displays the full shard key for [...]
+    "_id": { "index": 22, "kind": "header", "displayName": "", "group": 
"consumer changeStreams", "label": "consumer changeStreams", "required": false, 
"javaType": "Object", "deprecated": false, "deprecationNote": "", "autowired": 
false, "secret": false, "description": "The _id of the document created or 
modified by the insert, replace, delete or update operation (i.e. CRUD 
operations). It is an org.bson.types.ObjectId when MongoDB generated the id, 
and otherwise the id in its natural Ja [...]
     "CamelMongoDbStreamOperationType": { "index": 23, "kind": "header", 
"displayName": "", "group": "consumer changeStreams", "label": "consumer 
changeStreams", "required": false, "javaType": "String", "deprecated": false, 
"deprecationNote": "", "autowired": false, "secret": false, "description": "The 
type of operation that occurred. Can be any of the following values: insert, 
delete, replace, update, drop, rename, dropDatabase, invalidate.", 
"constantName": "org.apache.camel.component.m [...]
     "CamelMongoDbReturnDocumentType": { "index": 24, "kind": "header", 
"displayName": "", "group": "producer update one and return", "label": 
"producer update one and return", "required": false, "javaType": 
"com.mongodb.client.model.ReturnDocument", "enum": [ "BEFORE", "AFTER" ], 
"deprecated": false, "deprecationNote": "", "autowired": false, "secret": 
false, "description": "Indicates which document to return, the document before 
or after an update and return atomic operation.", "constan [...]
     "CamelMongoDbOperationOption": { "index": 25, "kind": "header", 
"displayName": "", "group": "producer update one and options", "label": 
"producer update one and options", "required": false, "javaType": "Object", 
"deprecated": false, "deprecationNote": "", "autowired": false, "secret": 
false, "description": "Options to use. When set, options set in the headers 
will be ignored.", "constantName": 
"org.apache.camel.component.mongodb.MongoDbConstants#OPTIONS" }
diff --git 
a/components/camel-mongodb/src/generated/resources/META-INF/org/apache/camel/component/mongodb/mongodb.json
 
b/components/camel-mongodb/src/generated/resources/META-INF/org/apache/camel/component/mongodb/mongodb.json
index e07db0812e4f..0118ae0c8885 100644
--- 
a/components/camel-mongodb/src/generated/resources/META-INF/org/apache/camel/component/mongodb/mongodb.json
+++ 
b/components/camel-mongodb/src/generated/resources/META-INF/org/apache/camel/component/mongodb/mongodb.json
@@ -53,7 +53,7 @@
     "CamelMongoDbDistinctQueryField": { "index": 19, "kind": "header", 
"displayName": "", "group": "producer", "label": "producer", "required": false, 
"javaType": "String", "deprecated": false, "deprecationNote": "", "autowired": 
false, "secret": false, "description": "The specified field name fow which we 
want to get the distinct values.", "constantName": 
"org.apache.camel.component.mongodb.MongoDbConstants#DISTINCT_QUERY_FIELD" },
     "CamelMongoDbAllowDiskUse": { "index": 20, "kind": "header", 
"displayName": "", "group": "producer findAll aggregate", "label": "producer 
findAll aggregate", "required": false, "javaType": "Boolean", "deprecated": 
false, "deprecationNote": "", "autowired": false, "secret": false, 
"description": "Sets allowDiskUse MongoDB flag. This is supported since MongoDB 
Server 4.3.1. Using this header with older MongoDB Server version can cause 
query to fail.", "constantName": "org.apache.camel. [...]
     "CamelMongoDbBulkOrdered": { "index": 21, "kind": "header", "displayName": 
"", "group": "producer bulkWrite", "label": "producer bulkWrite", "required": 
false, "javaType": "Boolean", "deprecated": false, "deprecationNote": "", 
"autowired": false, "secret": false, "defaultValue": "TRUE", "description": 
"Perform an ordered or unordered operation execution.", "constantName": 
"org.apache.camel.component.mongodb.MongoDbConstants#BULK_ORDERED" },
-    "_id": { "index": 22, "kind": "header", "displayName": "", "group": 
"consumer changeStreams", "label": "consumer changeStreams", "required": false, 
"javaType": "org.bson.types.ObjectId", "deprecated": false, "deprecationNote": 
"", "autowired": false, "secret": false, "description": "A document that 
contains the _id of the document created or modified by the insert, replace, 
delete, update operations (i.e. CRUD operations). For sharded collections, also 
displays the full shard key for [...]
+    "_id": { "index": 22, "kind": "header", "displayName": "", "group": 
"consumer changeStreams", "label": "consumer changeStreams", "required": false, 
"javaType": "Object", "deprecated": false, "deprecationNote": "", "autowired": 
false, "secret": false, "description": "The _id of the document created or 
modified by the insert, replace, delete or update operation (i.e. CRUD 
operations). It is an org.bson.types.ObjectId when MongoDB generated the id, 
and otherwise the id in its natural Ja [...]
     "CamelMongoDbStreamOperationType": { "index": 23, "kind": "header", 
"displayName": "", "group": "consumer changeStreams", "label": "consumer 
changeStreams", "required": false, "javaType": "String", "deprecated": false, 
"deprecationNote": "", "autowired": false, "secret": false, "description": "The 
type of operation that occurred. Can be any of the following values: insert, 
delete, replace, update, drop, rename, dropDatabase, invalidate.", 
"constantName": "org.apache.camel.component.m [...]
     "CamelMongoDbReturnDocumentType": { "index": 24, "kind": "header", 
"displayName": "", "group": "producer update one and return", "label": 
"producer update one and return", "required": false, "javaType": 
"com.mongodb.client.model.ReturnDocument", "enum": [ "BEFORE", "AFTER" ], 
"deprecated": false, "deprecationNote": "", "autowired": false, "secret": 
false, "description": "Indicates which document to return, the document before 
or after an update and return atomic operation.", "constan [...]
     "CamelMongoDbOperationOption": { "index": 25, "kind": "header", 
"displayName": "", "group": "producer update one and options", "label": 
"producer update one and options", "required": false, "javaType": "Object", 
"deprecated": false, "deprecationNote": "", "autowired": false, "secret": 
false, "description": "Options to use. When set, options set in the headers 
will be ignored.", "constantName": 
"org.apache.camel.component.mongodb.MongoDbConstants#OPTIONS" }
diff --git 
a/components/camel-mongodb/src/main/java/org/apache/camel/component/mongodb/MongoAbstractConsumerThread.java
 
b/components/camel-mongodb/src/main/java/org/apache/camel/component/mongodb/MongoAbstractConsumerThread.java
index ea2cca9dc0eb..b1f6cf415a9c 100644
--- 
a/components/camel-mongodb/src/main/java/org/apache/camel/component/mongodb/MongoAbstractConsumerThread.java
+++ 
b/components/camel-mongodb/src/main/java/org/apache/camel/component/mongodb/MongoAbstractConsumerThread.java
@@ -21,6 +21,9 @@ import java.util.concurrent.CountDownLatch;
 import com.mongodb.client.MongoCollection;
 import com.mongodb.client.MongoCursor;
 import org.apache.camel.Consumer;
+import org.apache.camel.Exchange;
+import org.apache.camel.spi.ExceptionHandler;
+import org.apache.camel.support.DefaultConsumer;
 import org.bson.Document;
 import org.slf4j.Logger;
 import org.slf4j.LoggerFactory;
@@ -50,6 +53,37 @@ abstract class MongoAbstractConsumerThread implements 
Runnable {
         this.cursorRegenerationDelayEnabled = !(this.cursorRegenerationDelay 
== 0);
     }
 
+    /**
+     * The exception handler of the consumer this thread feeds, so that a 
failure reaches the configured handler rather
+     * than being discarded on the consumer thread.
+     */
+    protected ExceptionHandler getExceptionHandler() {
+        return ((DefaultConsumer) consumer).getExceptionHandler();
+    }
+
+    /**
+     * Routes the exchange and reports a failure to the consumer's exception 
handler.
+     * <p>
+     * A route that fails does not throw from {@code process}: the failure is 
left on the exchange, so it is checked
+     * there.
+     *
+     * @param  exchange the exchange to route, created with {@code 
createExchange(false)} and released by the caller
+     * @return          {@code true} if the route processed the exchange 
without a failure
+     */
+    protected boolean processExchange(Exchange exchange) {
+        try {
+            consumer.getProcessor().process(exchange);
+        } catch (Exception e) {
+            exchange.setException(e);
+        }
+        Exception cause = exchange.getException();
+        if (cause != null) {
+            getExceptionHandler().handleException("Error processing exchange", 
exchange, cause);
+            return false;
+        }
+        return true;
+    }
+
     protected abstract MongoCursor<Document> initializeCursor();
 
     protected abstract void init() throws Exception;
@@ -70,8 +104,10 @@ abstract class MongoAbstractConsumerThread implements 
Runnable {
                     doRun();
                 } catch (Exception e) {
                     if (keepRunning) {
+                        // this one repeats for as long as the cause persists, 
so it is the branch that
+                        // needs the stack trace, not the one below
                         log.warn("Exception from consuming from MongoDB caused 
by {}. Will try again on next poll.",
-                                e.getMessage());
+                                e.getMessage(), e);
                     } else {
                         log.warn("Exception from consuming from MongoDB caused 
by {}. ConsumerThread will be stopped.",
                                 e.getMessage(), e);
diff --git 
a/components/camel-mongodb/src/main/java/org/apache/camel/component/mongodb/MongoDbChangeStreamsThread.java
 
b/components/camel-mongodb/src/main/java/org/apache/camel/component/mongodb/MongoDbChangeStreamsThread.java
index 1fb8cacae69c..51400d086fe6 100644
--- 
a/components/camel-mongodb/src/main/java/org/apache/camel/component/mongodb/MongoDbChangeStreamsThread.java
+++ 
b/components/camel-mongodb/src/main/java/org/apache/camel/component/mongodb/MongoDbChangeStreamsThread.java
@@ -27,11 +27,16 @@ import org.apache.camel.Exchange;
 import org.apache.camel.Message;
 import org.bson.BsonDocument;
 import org.bson.Document;
+import org.bson.codecs.DecoderContext;
+import org.bson.codecs.DocumentCodec;
 import org.bson.types.ObjectId;
 
 import static org.apache.camel.component.mongodb.MongoDbConstants.MONGO_ID;
 
 class MongoDbChangeStreamsThread extends MongoAbstractConsumerThread {
+
+    private static final DocumentCodec DOCUMENT_CODEC = new DocumentCodec();
+
     private List<BsonDocument> bsonFilter;
     private BsonDocument resumeToken;
     private CommitManager commitManager;
@@ -86,28 +91,38 @@ class MongoDbChangeStreamsThread extends 
MongoAbstractConsumerThread {
                 ChangeStreamDocument<Document> dbObj = 
(ChangeStreamDocument<Document>) cursor.next();
                 Exchange exchange = 
createMongoDbExchange(dbObj.getFullDocument());
 
-                ObjectId documentId = 
dbObj.getDocumentKey().getObjectId(MONGO_ID).getValue();
+                Object documentId = readDocumentId(dbObj.getDocumentKey());
                 OperationType operationType = dbObj.getOperationType();
                 BsonDocument currentResumeToken = dbObj.getResumeToken();
 
-                
exchange.getIn().setHeader(MongoDbConstants.STREAM_OPERATION_TYPE, 
operationType.getValue());
-                exchange.getIn().setHeader(MongoDbConstants.MONGO_ID, 
documentId);
+                if (operationType != null) {
+                    
exchange.getIn().setHeader(MongoDbConstants.STREAM_OPERATION_TYPE, 
operationType.getValue());
+                }
+                if (documentId != null) {
+                    exchange.getIn().setHeader(MongoDbConstants.MONGO_ID, 
documentId);
+                }
                 if (currentResumeToken != null) {
                     exchange.getIn().setHeader(Exchange.OFFSET, 
MongoDbResumable.of(resumeTokenKey, currentResumeToken));
                 }
-                if (operationType == OperationType.DELETE) {
+                if (operationType == OperationType.DELETE && documentId != 
null) {
                     exchange.getIn().setBody(new Document(MONGO_ID, 
documentId));
                 }
 
                 try {
                     if (log.isTraceEnabled()) {
-                        log.trace("Sending exchange: {}, ObjectId: {}", 
exchange, dbObj.getFullDocument().get(MONGO_ID));
+                        log.trace("Sending exchange: {}, id: {}", exchange, 
documentId);
+                    }
+                    // a failed event does not advance the resume token 
itself, but a later event that
+                    // succeeds commits its own, so the failure is reported 
rather than retried
+                    if (processExchange(exchange)) {
+                        this.resumeToken = currentResumeToken;
+                        commitManager.recordResumeToken(currentResumeToken);
+                        commitManager.commit();
                     }
-                    consumer.getProcessor().process(exchange);
-                    this.resumeToken = currentResumeToken;
-                    commitManager.recordResumeToken(currentResumeToken);
-                    commitManager.commit();
-                } catch (Exception ignored) {
+                } catch (Exception e) {
+                    getExceptionHandler().handleException("Error committing 
the change stream resume token", exchange, e);
+                } finally {
+                    consumer.releaseExchange(exchange, false);
                 }
             }
         } catch (MongoException e) {
@@ -122,8 +137,36 @@ class MongoDbChangeStreamsThread extends 
MongoAbstractConsumerThread {
         }
     }
 
+    /**
+     * Reads the {@code _id} out of the change event's document key.
+     * <p>
+     * The key is absent on the events that do not belong to a single document 
({@code invalidate}, {@code drop},
+     * {@code rename}, {@code dropDatabase}), and {@code _id} is only an 
{@link ObjectId} when the collection lets
+     * MongoDB generate it - a document may just as well be keyed by a string, 
a number or a compound value. Reading it
+     * blindly as an {@link ObjectId} threw before the exchange was ever 
created, and since the resume token is only
+     * advanced after a successful exchange, the regenerated cursor kept 
returning the same event.
+     *
+     * @param  documentKey the change event's document key, which may be 
{@code null}
+     * @return             the id as its natural Java type, or {@code null} 
when the event carries no document key
+     */
+    static Object readDocumentId(BsonDocument documentKey) {
+        if (documentKey == null || !documentKey.containsKey(MONGO_ID)) {
+            return null;
+        }
+
+        if (documentKey.get(MONGO_ID).isObjectId()) {
+            return documentKey.getObjectId(MONGO_ID).getValue();
+        }
+
+        // anything else - a string, a number, a compound key - is decoded the 
way the driver decodes a
+        // document, so the header carries the id in its natural Java type
+        Document decoded = DOCUMENT_CODEC.decode(documentKey.asBsonReader(), 
DecoderContext.builder().build());
+        return decoded.get(MONGO_ID);
+    }
+
     private Exchange createMongoDbExchange(Document dbObj) {
-        Exchange exchange = consumer.createExchange(true);
+        // released by doRun once the outcome has been read, as an 
auto-released exchange may already be reset by then
+        Exchange exchange = consumer.createExchange(false);
         Message message = exchange.getIn();
         message.setHeader(MongoDbConstants.DATABASE, endpoint.getDatabase());
         message.setHeader(MongoDbConstants.COLLECTION, 
endpoint.getCollection());
diff --git 
a/components/camel-mongodb/src/main/java/org/apache/camel/component/mongodb/MongoDbConstants.java
 
b/components/camel-mongodb/src/main/java/org/apache/camel/component/mongodb/MongoDbConstants.java
index 37b40d9a712a..cf6c66b6d6b8 100644
--- 
a/components/camel-mongodb/src/main/java/org/apache/camel/component/mongodb/MongoDbConstants.java
+++ 
b/components/camel-mongodb/src/main/java/org/apache/camel/component/mongodb/MongoDbConstants.java
@@ -83,10 +83,12 @@ public final class MongoDbConstants {
     public static final String BULK_ORDERED = "CamelMongoDbBulkOrdered";
     @Metadata(label = "consumer changeStreams",
               description = """
-                      A document that contains the _id of the document created 
or modified by the insert,
-                      replace, delete, update operations (i.e. CRUD 
operations). For sharded collections, also displays the full shard key for
-                      the document. The _id field is not repeated if it is 
already a part of the shard key.""",
-              javaType = "org.bson.types.ObjectId")
+                      The _id of the document created or modified by the 
insert, replace, delete or update operation
+                      (i.e. CRUD operations). It is an org.bson.types.ObjectId 
when MongoDB generated the id, and
+                      otherwise the id in its natural Java type: a String, a 
number, or a Document for a compound key.
+                      The header is absent on the events that do not belong to 
a single document, such as invalidate,
+                      drop, rename and dropDatabase.""",
+              javaType = "Object")
     public static final String MONGO_ID = "_id"; // default id field
     @Metadata(label = "consumer changeStreams", description = """
             The type of operation that occurred. Can
diff --git 
a/components/camel-mongodb/src/main/java/org/apache/camel/component/mongodb/MongoDbTailingThread.java
 
b/components/camel-mongodb/src/main/java/org/apache/camel/component/mongodb/MongoDbTailingThread.java
index e578575846f4..c9249704d3d6 100644
--- 
a/components/camel-mongodb/src/main/java/org/apache/camel/component/mongodb/MongoDbTailingThread.java
+++ 
b/components/camel-mongodb/src/main/java/org/apache/camel/component/mongodb/MongoDbTailingThread.java
@@ -115,11 +115,13 @@ class MongoDbTailingThread extends 
MongoAbstractConsumerThread {
                     if (log.isTraceEnabled()) {
                         log.trace("Sending exchange: {}, ObjectId: {}", 
exchange, dbObj.get(MONGO_ID));
                     }
-                    consumer.getProcessor().process(exchange);
-                } catch (Exception e) {
-                    // do nothing
+                    if (processExchange(exchange)) {
+                        // only once the route accepted it, so a failed record 
does not move the tail position
+                        tailTracking.setLastVal(dbObj);
+                    }
+                } finally {
+                    consumer.releaseExchange(exchange, false);
                 }
-                tailTracking.setLastVal(dbObj);
             }
         } catch (MongoCursorNotFoundException e) {
             // we only log the warning if we are not stopping, otherwise it is
@@ -147,7 +149,8 @@ class MongoDbTailingThread extends 
MongoAbstractConsumerThread {
     }
 
     Exchange createMongoDbExchange(Document dbObj) {
-        Exchange exchange = consumer.createExchange(true);
+        // released by doRun once the outcome has been read, as an 
auto-released exchange may already be reset by then
+        Exchange exchange = consumer.createExchange(false);
         Message message = exchange.getIn();
         message.setHeader(MongoDbConstants.DATABASE, endpoint.getDatabase());
         message.setHeader(MongoDbConstants.COLLECTION, 
endpoint.getCollection());
diff --git 
a/components/camel-mongodb/src/test/java/org/apache/camel/component/mongodb/MongoDbChangeStreamsThreadTest.java
 
b/components/camel-mongodb/src/test/java/org/apache/camel/component/mongodb/MongoDbChangeStreamsThreadTest.java
new file mode 100644
index 000000000000..ed78aa441d75
--- /dev/null
+++ 
b/components/camel-mongodb/src/test/java/org/apache/camel/component/mongodb/MongoDbChangeStreamsThreadTest.java
@@ -0,0 +1,84 @@
+/*
+ * 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.mongodb;
+
+import org.bson.BsonDocument;
+import org.bson.BsonInt32;
+import org.bson.BsonObjectId;
+import org.bson.BsonString;
+import org.bson.Document;
+import org.bson.types.ObjectId;
+import org.junit.jupiter.api.Test;
+
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertInstanceOf;
+import static org.junit.jupiter.api.Assertions.assertNull;
+
+/**
+ * The document key is read before the exchange is created, and the resume 
token only advances after a successful
+ * exchange, so anything thrown here made the consumer re-read the same event 
for as long as the route ran.
+ */
+public class MongoDbChangeStreamsThreadTest {
+
+    @Test
+    public void testGeneratedObjectIdIsReadAsAnObjectId() {
+        ObjectId id = new ObjectId();
+        BsonDocument key = new BsonDocument("_id", new BsonObjectId(id));
+
+        Object read = MongoDbChangeStreamsThread.readDocumentId(key);
+
+        assertInstanceOf(ObjectId.class, read);
+        assertEquals(id, read);
+    }
+
+    @Test
+    public void testStringIdIsReadAsAString() {
+        BsonDocument key = new BsonDocument("_id", new 
BsonString("a-string-key"));
+
+        assertEquals("a-string-key", 
MongoDbChangeStreamsThread.readDocumentId(key));
+    }
+
+    @Test
+    public void testNumericIdIsReadAsANumber() {
+        BsonDocument key = new BsonDocument("_id", new BsonInt32(42));
+
+        assertEquals(42, MongoDbChangeStreamsThread.readDocumentId(key));
+    }
+
+    @Test
+    public void testCompoundIdIsReadAsADocument() {
+        BsonDocument compound = new BsonDocument("tenant", new 
BsonString("acme")).append("seq", new BsonInt32(7));
+        BsonDocument key = new BsonDocument("_id", compound);
+
+        Object read = MongoDbChangeStreamsThread.readDocumentId(key);
+
+        assertInstanceOf(Document.class, read);
+        assertEquals("acme", ((Document) read).get("tenant"));
+        assertEquals(7, ((Document) read).get("seq"));
+    }
+
+    @Test
+    public void testEventWithoutADocumentKeyHasNoId() {
+        // invalidate, drop, rename and dropDatabase carry no document key at 
all
+        assertNull(MongoDbChangeStreamsThread.readDocumentId(null));
+    }
+
+    @Test
+    public void testDocumentKeyWithoutAnIdHasNoId() {
+        assertNull(MongoDbChangeStreamsThread.readDocumentId(new 
BsonDocument()));
+    }
+}
diff --git 
a/components/camel-mongodb/src/test/java/org/apache/camel/component/mongodb/integration/MongoDbChangeStreamsConsumerIT.java
 
b/components/camel-mongodb/src/test/java/org/apache/camel/component/mongodb/integration/MongoDbChangeStreamsConsumerIT.java
index 591ba00e3a94..150d46847cfa 100644
--- 
a/components/camel-mongodb/src/test/java/org/apache/camel/component/mongodb/integration/MongoDbChangeStreamsConsumerIT.java
+++ 
b/components/camel-mongodb/src/test/java/org/apache/camel/component/mongodb/integration/MongoDbChangeStreamsConsumerIT.java
@@ -16,7 +16,9 @@
  */
 package org.apache.camel.component.mongodb.integration;
 
+import java.util.List;
 import java.util.concurrent.CompletableFuture;
+import java.util.concurrent.TimeUnit;
 
 import com.mongodb.client.MongoCollection;
 import com.mongodb.client.model.CreateCollectionOptions;
@@ -26,6 +28,7 @@ import org.apache.camel.builder.RouteBuilder;
 import org.apache.camel.component.mock.MockEndpoint;
 import org.apache.camel.test.infra.core.annotations.RouteFixture;
 import org.apache.camel.test.infra.core.api.ConfigurableRoute;
+import org.awaitility.Awaitility;
 import org.bson.Document;
 import org.bson.types.ObjectId;
 import org.junit.jupiter.api.AfterEach;
@@ -44,6 +47,7 @@ import static org.junit.jupiter.api.Assertions.assertTrue;
 public class MongoDbChangeStreamsConsumerIT extends AbstractMongoDbITSupport 
implements ConfigurableRoute {
 
     private MongoCollection<Document> mongoCollection;
+    private final RecordingExceptionHandler failures = new 
RecordingExceptionHandler("increasing");
 
     /*
      * NOTE: in the case of this test, we *DO* want to recreate everything 
after the test has executed, so that when
@@ -167,6 +171,33 @@ public class MongoDbChangeStreamsConsumerIT extends 
AbstractMongoDbITSupport imp
         context.getRouteController().stopRoute(consumerRouteId);
     }
 
+    @Order(5)
+    @Test
+    public void failedExchangeIsReportedTest() throws Exception {
+        Assumptions.assumeTrue(0 == mongoCollection.countDocuments(), "The 
collection should have no documents");
+        failures.clear();
+        MockEndpoint mock = 
contextExtension.getMockEndpoint("mock:changeStreamFailing");
+        mock.reset();
+        mock.expectedMessageCount(2);
+
+        String consumerRouteId = "failingConsumer";
+        context.getRouteController().startRoute(consumerRouteId);
+
+        // the route fails for the second event, increasing=2
+        CompletableFuture.runAsync(() -> {
+            for (int i = 1; i <= 3; i++) {
+                mongoCollection.insertOne(new Document("increasing", i));
+            }
+        });
+
+        mock.assertIsSatisfied();
+        Awaitility.await().atMost(10, TimeUnit.SECONDS).until(() -> 
!failures.getValues().isEmpty());
+        context.getRouteController().stopRoute(consumerRouteId);
+
+        // the route failure reached the consumer's exception handler
+        assertEquals(List.of(2), failures.getValues());
+    }
+
     private void insertAndDelete(ObjectId objectId) {
         mongoCollection.insertOne(new Document("_id", 
objectId).append("string", "value"));
         mongoCollection.deleteOne(new Document("_id", objectId));
@@ -175,6 +206,7 @@ public class MongoDbChangeStreamsConsumerIT extends 
AbstractMongoDbITSupport imp
     @RouteFixture
     @Override
     public void createRouteBuilder(CamelContext context) throws Exception {
+        context.getRegistry().bind("changeStreamFailures", failures);
         context.addRoutes(new RouteBuilder() {
 
             @Override
@@ -193,6 +225,16 @@ public class MongoDbChangeStreamsConsumerIT extends 
AbstractMongoDbITSupport imp
                         .id("updateWithFullDocumentConsumer")
                         .autoStartup(false)
                         .to("mock:test");
+
+                
from("mongodb:myDb?consumerType=changeStreams&database={{mongodb.testDb}}&collection={{mongodb.testCollection}}&exceptionHandler=#changeStreamFailures")
+                        .id("failingConsumer")
+                        .autoStartup(false)
+                        .process(exchange -> {
+                            if 
(Integer.valueOf(2).equals(exchange.getIn().getBody(Document.class).get("increasing")))
 {
+                                throw new IllegalStateException("Simulated 
route failure");
+                            }
+                        })
+                        .to("mock:changeStreamFailing");
             }
         });
     }
diff --git 
a/components/camel-mongodb/src/test/java/org/apache/camel/component/mongodb/integration/MongoDbTailableCursorConsumerIT.java
 
b/components/camel-mongodb/src/test/java/org/apache/camel/component/mongodb/integration/MongoDbTailableCursorConsumerIT.java
index fa7235096baf..acbe2dc43d68 100644
--- 
a/components/camel-mongodb/src/test/java/org/apache/camel/component/mongodb/integration/MongoDbTailableCursorConsumerIT.java
+++ 
b/components/camel-mongodb/src/test/java/org/apache/camel/component/mongodb/integration/MongoDbTailableCursorConsumerIT.java
@@ -17,6 +17,7 @@
 package org.apache.camel.component.mongodb.integration;
 
 import java.util.Calendar;
+import java.util.List;
 import java.util.concurrent.ExecutorService;
 import java.util.concurrent.Executors;
 import java.util.concurrent.TimeUnit;
@@ -52,6 +53,7 @@ public class MongoDbTailableCursorConsumerIT extends 
AbstractMongoDbITSupport im
     private String cappedTestCollectionName;
     private CreateCollectionOptions createCollectionOptions;
     private ExecutorService executorService = Executors.newCachedThreadPool();
+    private final RecordingExceptionHandler failures = new 
RecordingExceptionHandler("increasing");
 
     @BeforeEach
     void checkDocuments() {
@@ -330,6 +332,31 @@ public class MongoDbTailableCursorConsumerIT extends 
AbstractMongoDbITSupport im
 
     }
 
+    @Test
+    public void testFailedRecordIsReportedAndDoesNotMoveTheTailPosition() 
throws Exception {
+        assertEquals(0, cappedTestCollection.countDocuments());
+        MongoCollection<Document> trackingCol = 
db.getCollection(MongoDbTailTrackingConfig.DEFAULT_COLLECTION, Document.class);
+        trackingCol.deleteMany(eq("persistentId", "failing"));
+        failures.clear();
+
+        MockEndpoint mock = 
contextExtension.getMockEndpoint("mock:tailFailing");
+        mock.reset();
+        mock.expectedMessageCount(2);
+
+        
context.getRouteController().startRoute("tailableCursorConsumerFailing");
+        // the route fails for the last record, increasing=3
+        doQuickInsert(1, 3);
+
+        mock.assertIsSatisfied();
+        Awaitility.await().atMost(10, TimeUnit.SECONDS).until(() -> 
!failures.getValues().isEmpty());
+        
context.getRouteController().stopRoute("tailableCursorConsumerFailing");
+
+        // the route failure reached the consumer's exception handler
+        assertEquals(List.of(3), failures.getValues());
+        // and the failed record did not move the persisted position past the 
last record the route accepted
+        assertEquals(2, trackingCol.find(eq("persistentId", 
"failing")).first().get("lastTrackingValue"));
+    }
+
     public void assertAndResetMockEndpoint(MockEndpoint mock) throws Exception 
{
         mock.assertIsSatisfied();
         mock.reset();
@@ -374,6 +401,7 @@ public class MongoDbTailableCursorConsumerIT extends 
AbstractMongoDbITSupport im
     @RouteFixture
     @Override
     public void createRouteBuilder(CamelContext context) throws Exception {
+        context.getRegistry().bind("tailFailures", failures);
         context.addRoutes(new RouteBuilder() {
 
             @Override
@@ -393,6 +421,15 @@ public class MongoDbTailableCursorConsumerIT extends 
AbstractMongoDbITSupport im
 
                 
from("mongodb:myDb?database={{mongodb.testDb}}&collection={{mongodb.cappedTestCollection}}&tailTrackIncreasingField=increasing")//
 &readPreference=primary")
                         
.id("tailableCursorConsumer1.readPreference").autoStartup(false).to("mock:test");
+                
from("mongodb:myDb?database={{mongodb.testDb}}&collection={{mongodb.cappedTestCollection}}&tailTrackIncreasingField=increasing&"
+                     + 
"persistentTailTracking=true&persistentId=failing&exceptionHandler=#tailFailures")
+                        .id("tailableCursorConsumerFailing").autoStartup(false)
+                        .process(exchange -> {
+                            if 
(Integer.valueOf(3).equals(exchange.getIn().getBody(Document.class).get("increasing")))
 {
+                                throw new IllegalStateException("Simulated 
route failure");
+                            }
+                        })
+                        .to("mock:tailFailing");
 
             }
         });
diff --git 
a/components/camel-mongodb/src/test/java/org/apache/camel/component/mongodb/integration/RecordingExceptionHandler.java
 
b/components/camel-mongodb/src/test/java/org/apache/camel/component/mongodb/integration/RecordingExceptionHandler.java
new file mode 100644
index 000000000000..3ea4c25506ff
--- /dev/null
+++ 
b/components/camel-mongodb/src/test/java/org/apache/camel/component/mongodb/integration/RecordingExceptionHandler.java
@@ -0,0 +1,62 @@
+/*
+ * 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.mongodb.integration;
+
+import java.util.List;
+import java.util.concurrent.CopyOnWriteArrayList;
+
+import org.apache.camel.Exchange;
+import org.apache.camel.spi.ExceptionHandler;
+import org.bson.Document;
+
+/**
+ * An {@link ExceptionHandler} that records one field of the {@link Document} 
body of each failed exchange, so a test
+ * can tell which record the consumer reported.
+ */
+class RecordingExceptionHandler implements ExceptionHandler {
+
+    private final String field;
+    private final List<Object> values = new CopyOnWriteArrayList<>();
+
+    RecordingExceptionHandler(String field) {
+        this.field = field;
+    }
+
+    List<Object> getValues() {
+        return values;
+    }
+
+    void clear() {
+        values.clear();
+    }
+
+    @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) {
+        Document body = exchange != null ? 
exchange.getIn().getBody(Document.class) : null;
+        values.add(body != null ? body.get(field) : null);
+    }
+}
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 473d491af122..0c424094d549 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
@@ -3257,6 +3257,37 @@ healthy while no longer receiving any change event. 
Deployments that use readine
 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-mongodb - a failed exchange is reported, and does not advance the 
consumer's position
+
+Both MongoDB consumers ignored the outcome of the route. A failure left on the 
exchange was never looked at,
+and an exception thrown from the processor was caught and dropped. With the 
default error handler the
+failure was still logged once redelivery was exhausted, but with 
`noErrorHandler`, or an error handler that
+does not log, it left no trace. The failure is now also passed to the 
consumer's `ExceptionHandler`, which
+logs it by default, so expect an additional log line for each failed exchange.
+
+A failed exchange also no longer advances the consumer's position. The 
tailable cursor consumer updates its
+tail tracking value, and the change streams consumer records and commits its 
resume token, only after the
+route processed the event without a failure. Previously both did so for every 
event, whether or not the
+route failed. With `persistentTailTracking=true` the stored position is 
therefore that of the last
+successful record.
+
+This is not a retry. The next event that succeeds moves the position past the 
failed one. A failed event is
+read again only if the consumer reopens its cursor from the stored position 
before a later event has
+succeeded, for example when the cursor is regenerated, or on a restart with 
persistent tail tracking or a
+resume strategy.
+
+=== camel-mongodb - the change stream id header is no longer always an ObjectId
+
+The `_id` header set by the change streams consumer 
(`MongoDbConstants.MONGO_ID`) used to be read as an
+`org.bson.types.ObjectId` unconditionally, which threw for a document whose 
`_id` is a string, a number or
+a compound key, and for the events that carry no document key at all 
(`invalidate`, `drop`, `rename`,
+`dropDatabase`).
+
+It now carries the id in its natural Java type - an `ObjectId` when MongoDB 
generated it, otherwise a
+`String`, a number, or a `Document` for a compound key - and is absent for the 
events without a document
+key. A route that casts this header to `ObjectId` should check the type first, 
or keep working unchanged
+if its collection uses generated ids.
+
 === 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
diff --git 
a/dsl/camel-endpointdsl/src/generated/java/org/apache/camel/builder/endpoint/dsl/MongoDbEndpointBuilderFactory.java
 
b/dsl/camel-endpointdsl/src/generated/java/org/apache/camel/builder/endpoint/dsl/MongoDbEndpointBuilderFactory.java
index 93b98e3d4b50..3e09c50e0d57 100644
--- 
a/dsl/camel-endpointdsl/src/generated/java/org/apache/camel/builder/endpoint/dsl/MongoDbEndpointBuilderFactory.java
+++ 
b/dsl/camel-endpointdsl/src/generated/java/org/apache/camel/builder/endpoint/dsl/MongoDbEndpointBuilderFactory.java
@@ -4464,13 +4464,15 @@ public interface MongoDbEndpointBuilderFactory {
             return "CamelMongoDbBulkOrdered";
         }
         /**
-         * A document that contains the _id of the document created or modified
-         * by the insert, replace, delete, update operations (i.e. CRUD
-         * operations). For sharded collections, also displays the full shard
-         * key for the document. The _id field is not repeated if it is already
-         * a part of the shard key.
+         * The _id of the document created or modified by the insert, replace,
+         * delete or update operation (i.e. CRUD operations). It is an
+         * org.bson.types.ObjectId when MongoDB generated the id, and otherwise
+         * the id in its natural Java type: a String, a number, or a Document
+         * for a compound key. The header is absent on the events that do not
+         * belong to a single document, such as invalidate, drop, rename and
+         * dropDatabase.
          * 
-         * The option is a: {@code org.bson.types.ObjectId} type.
+         * The option is a: {@code Object} type.
          * 
          * Group: consumer changeStreams
          * 

Reply via email to