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
*