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 c3632125413b CAMEL-25159, CAMEL-25160: camel-mongodb - stop running
queries twice, and harden tail tracking (#27116)
c3632125413b is described below
commit c3632125413bec6fe6de557e0898b592b3ff684b
Author: Andrea Cosentino <[email protected]>
AuthorDate: Wed Sep 30 11:46:35 2026 +0200
CAMEL-25159, CAMEL-25160: camel-mongodb - stop running queries twice, and
harden tail tracking (#27116)
CAMEL-25159: findAll, distinct and aggregate read their results with
ret.iterator().forEachRemaining(...) and then, in a finally block meant to
release
the cursor, called ret.iterator().close(). MongoIterable.iterator()
executes the
operation, so the finally opened and closed a second cursor (a second round
trip,
and for aggregate a second run of the pipeline) and left the cursor that
was read
unclosed. Hold the cursor with try-with-resources, so the query runs once
and the
cursor that is read is the one that is closed.
CAMEL-25160: persistToStore stored the full document returned by
findOneAndUpdate
back into trackingObj, and from then on its filter also matched the previous
value, so another writer of that field (two routes sharing a persistentId)
left
the update matching nothing and the next persist throwing. Filter on the id
only.
recoverFromStore now handles a missing tracking document by warning and
starting
from the beginning.
Co-Authored-By: Claude Opus 5 <[email protected]>
Signed-off-by: Andrea Cosentino <[email protected]>
---
components/camel-mongodb/pom.xml | 6 +
.../camel/component/mongodb/MongoDbProducer.java | 26 ++--
.../mongodb/MongoDbTailTrackingManager.java | 15 ++-
.../mongodb/MongoDbProducerCursorTest.java | 140 +++++++++++++++++++++
4 files changed, 168 insertions(+), 19 deletions(-)
diff --git a/components/camel-mongodb/pom.xml b/components/camel-mongodb/pom.xml
index faec97739249..64002f9b2a1c 100644
--- a/components/camel-mongodb/pom.xml
+++ b/components/camel-mongodb/pom.xml
@@ -92,6 +92,12 @@
<scope>test</scope>
</dependency>
+ <dependency>
+ <groupId>org.mockito</groupId>
+ <artifactId>mockito-core</artifactId>
+ <version>${mockito-version}</version>
+ <scope>test</scope>
+ </dependency>
<dependency>
<groupId>org.apache.commons</groupId>
<artifactId>commons-lang3</artifactId>
diff --git
a/components/camel-mongodb/src/main/java/org/apache/camel/component/mongodb/MongoDbProducer.java
b/components/camel-mongodb/src/main/java/org/apache/camel/component/mongodb/MongoDbProducer.java
index 9383d8e3ae00..0cc7af5dc8ea 100644
---
a/components/camel-mongodb/src/main/java/org/apache/camel/component/mongodb/MongoDbProducer.java
+++
b/components/camel-mongodb/src/main/java/org/apache/camel/component/mongodb/MongoDbProducer.java
@@ -30,6 +30,7 @@ import com.mongodb.client.AggregateIterable;
import com.mongodb.client.DistinctIterable;
import com.mongodb.client.FindIterable;
import com.mongodb.client.MongoCollection;
+import com.mongodb.client.MongoCursor;
import com.mongodb.client.MongoDatabase;
import com.mongodb.client.model.BulkWriteOptions;
import com.mongodb.client.model.Filters;
@@ -362,11 +363,10 @@ public class MongoDbProducer extends DefaultProducer {
} else {
ret = dbCol.distinct(distinctFieldName, String.class);
}
- try {
- ret.iterator().forEachRemaining(result::add);
+ // hold the cursor: every call to iterator() executes the query
again
+ try (MongoCursor<String> cursor = ret.iterator()) {
+ cursor.forEachRemaining(result::add);
exchange.getMessage().setHeader(MongoDbConstants.RESULT_PAGE_SIZE,
result.size());
- } finally {
- ret.iterator().close();
}
return result;
};
@@ -420,12 +420,11 @@ public class MongoDbProducer extends DefaultProducer {
ret.allowDiskUse(exchange.getIn().getHeader(MongoDbConstants.ALLOW_DISK_USE,
Boolean.class));
if
(!MongoDbOutputType.MongoIterable.equals(endpoint.getOutputType())) {
- try {
- result = new ArrayList<>();
- ret.iterator().forEachRemaining(((List<Document>)
result)::add);
+ result = new ArrayList<>();
+ // hold the cursor: every call to iterator() executes the
query again
+ try (MongoCursor<Document> cursor = ret.iterator()) {
+ cursor.forEachRemaining(((List<Document>) result)::add);
exchange.getMessage().setHeader(RESULT_PAGE_SIZE,
((List<Document>) result).size());
- } finally {
- ret.iterator().close();
}
} else {
result = ret;
@@ -572,12 +571,11 @@ public class MongoDbProducer extends DefaultProducer {
Iterable<Document> result;
if
(!MongoDbOutputType.MongoIterable.equals(endpoint.getOutputType())) {
- try {
- result = new ArrayList<>();
-
aggregationResult.iterator().forEachRemaining(((List<Document>) result)::add);
+ result = new ArrayList<>();
+ // hold the cursor: every call to iterator() executes the
pipeline again
+ try (MongoCursor<Document> cursor =
aggregationResult.iterator()) {
+ cursor.forEachRemaining(((List<Document>)
result)::add);
exchange.getMessage().setHeader(MongoDbConstants.RESULT_PAGE_SIZE,
((List<Document>) result).size());
- } finally {
- aggregationResult.iterator().close();
}
} else {
result = aggregationResult;
diff --git
a/components/camel-mongodb/src/main/java/org/apache/camel/component/mongodb/MongoDbTailTrackingManager.java
b/components/camel-mongodb/src/main/java/org/apache/camel/component/mongodb/MongoDbTailTrackingManager.java
index 653d32b03ea3..074f7602b56d 100644
---
a/components/camel-mongodb/src/main/java/org/apache/camel/component/mongodb/MongoDbTailTrackingManager.java
+++
b/components/camel-mongodb/src/main/java/org/apache/camel/component/mongodb/MongoDbTailTrackingManager.java
@@ -21,8 +21,6 @@ import java.util.concurrent.locks.ReentrantLock;
import com.mongodb.client.MongoClient;
import com.mongodb.client.MongoCollection;
-import com.mongodb.client.model.FindOneAndUpdateOptions;
-import com.mongodb.client.model.ReturnDocument;
import com.mongodb.client.model.Updates;
import org.bson.Document;
import org.bson.conversions.Bson;
@@ -76,8 +74,9 @@ public class MongoDbTailTrackingManager {
}
Bson updateObj = Updates.set(config.field, lastVal);
- FindOneAndUpdateOptions options = new
FindOneAndUpdateOptions().returnDocument(ReturnDocument.AFTER);
- trackingObj = dbCol.findOneAndUpdate(trackingObj, updateObj,
options);
+ // filter by the id only: storing the returned document back would
make every later update
+ // also match on the previous value, and a filter that stops
matching leaves trackingObj null
+ dbCol.updateOne(trackingObj, updateObj);
} finally {
lock.unlock();
}
@@ -90,7 +89,13 @@ public class MongoDbTailTrackingManager {
return null;
}
- lastVal = dbCol.find(trackingObj).first().get(config.field);
+ Document tracked = dbCol.find(trackingObj).first();
+ if (tracked == null) {
+ LOG.warn("Tail tracking document {} is gone from collection
{}, starting from the beginning",
+ trackingObj, config.collection);
+ return null;
+ }
+ lastVal = tracked.get(config.field);
if (LOG.isDebugEnabled()) {
LOG.debug("Recovered lastVal={} from store, collection: {}",
lastVal, config.collection);
diff --git
a/components/camel-mongodb/src/test/java/org/apache/camel/component/mongodb/MongoDbProducerCursorTest.java
b/components/camel-mongodb/src/test/java/org/apache/camel/component/mongodb/MongoDbProducerCursorTest.java
new file mode 100644
index 000000000000..b10caf79ca5a
--- /dev/null
+++
b/components/camel-mongodb/src/test/java/org/apache/camel/component/mongodb/MongoDbProducerCursorTest.java
@@ -0,0 +1,140 @@
+/*
+ * 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 com.mongodb.client.AggregateIterable;
+import com.mongodb.client.DistinctIterable;
+import com.mongodb.client.FindIterable;
+import com.mongodb.client.MongoClient;
+import com.mongodb.client.MongoCollection;
+import com.mongodb.client.MongoCursor;
+import com.mongodb.client.MongoDatabase;
+import org.apache.camel.CamelContext;
+import org.apache.camel.RoutesBuilder;
+import org.apache.camel.builder.RouteBuilder;
+import org.apache.camel.test.junit6.CamelTestSupport;
+import org.bson.Document;
+import org.junit.jupiter.api.Test;
+
+import static org.mockito.ArgumentMatchers.any;
+import static org.mockito.ArgumentMatchers.anyString;
+import static org.mockito.ArgumentMatchers.eq;
+import static org.mockito.Mockito.mock;
+import static org.mockito.Mockito.times;
+import static org.mockito.Mockito.verify;
+import static org.mockito.Mockito.when;
+
+/**
+ * Every call to MongoIterable.iterator() executes the operation again - it is
not an accessor for an already open
+ * cursor - so closing "the" cursor in a finally block by calling iterator() a
second time ran the whole query twice and
+ * left the cursor that had actually been read unclosed.
+ */
+class MongoDbProducerCursorTest extends CamelTestSupport {
+
+ private final MongoCollection<Document> collection =
mock(MongoCollection.class);
+ private final DistinctIterable<String> distinctIterable =
mock(DistinctIterable.class);
+ private final FindIterable<Document> findIterable =
mock(FindIterable.class);
+ private final AggregateIterable<Document> aggregateIterable =
mock(AggregateIterable.class);
+
+ @Test
+ void testDistinctOpensOneCursor() {
+ // findDistinct needs the field name, otherwise the operation has
nothing to ask for
+ template.requestBodyAndHeader("direct:distinct", null,
MongoDbConstants.DISTINCT_QUERY_FIELD, "name");
+
+ verify(distinctIterable, times(1)).iterator();
+ }
+
+ @Test
+ void testFindAllOpensOneCursor() {
+ template.requestBody("direct:findAll", (Object) null);
+
+ verify(findIterable, times(1)).iterator();
+ }
+
+ @Test
+ void testAggregateOpensOneCursor() {
+ template.requestBody("direct:aggregate", "[{$match: {}}]");
+
+ verify(aggregateIterable, times(1)).iterator();
+ }
+
+ @Override
+ protected CamelContext createCamelContext() throws Exception {
+ final CamelContext context = super.createCamelContext();
+
+ when(collection.withWriteConcern(any())).thenReturn(collection);
+ when(collection.distinct(anyString(),
eq(String.class))).thenReturn(distinctIterable);
+ when(collection.distinct(anyString(), any(),
eq(String.class))).thenReturn(distinctIterable);
+
when(collection.find(any(org.bson.conversions.Bson.class))).thenReturn(findIterable);
+ when(collection.find()).thenReturn(findIterable);
+ when(collection.aggregate(any())).thenReturn(aggregateIterable);
+
+ stubCursor(distinctIterable);
+ stubFind(findIterable);
+ stubAggregate(aggregateIterable);
+
+ final MongoDatabase database = mock(MongoDatabase.class);
+ when(database.getCollection(anyString(),
eq(Document.class))).thenReturn(collection);
+
+ final MongoClient client = mock(MongoClient.class);
+ when(client.getDatabase(anyString())).thenReturn(database);
+
when(client.listDatabaseNames()).thenReturn(mock(com.mongodb.client.MongoIterable.class));
+
+ context.getRegistry().bind("mongoClient", client);
+ return context;
+ }
+
+ private void stubCursor(DistinctIterable<String> iterable) {
+ final MongoCursor<String> cursor = mock(MongoCursor.class);
+ when(cursor.hasNext()).thenReturn(false);
+ when(iterable.iterator()).thenReturn(cursor);
+ }
+
+ private void stubFind(FindIterable<Document> iterable) {
+ final MongoCursor<Document> cursor = mock(MongoCursor.class);
+ when(cursor.hasNext()).thenReturn(false);
+ when(iterable.iterator()).thenReturn(cursor);
+ when(iterable.projection(any())).thenReturn(iterable);
+ when(iterable.sort(any())).thenReturn(iterable);
+
when(iterable.skip(org.mockito.ArgumentMatchers.anyInt())).thenReturn(iterable);
+
when(iterable.limit(org.mockito.ArgumentMatchers.anyInt())).thenReturn(iterable);
+ when(iterable.allowDiskUse(any())).thenReturn(iterable);
+ }
+
+ private void stubAggregate(AggregateIterable<Document> iterable) {
+ final MongoCursor<Document> cursor = mock(MongoCursor.class);
+ when(cursor.hasNext()).thenReturn(false);
+ when(iterable.iterator()).thenReturn(cursor);
+ when(iterable.allowDiskUse(any())).thenReturn(iterable);
+
when(iterable.batchSize(org.mockito.ArgumentMatchers.anyInt())).thenReturn(iterable);
+ }
+
+ @Override
+ protected RoutesBuilder createRouteBuilder() {
+ return new RouteBuilder() {
+ @Override
+ public void configure() {
+ from("direct:distinct")
+
.to("mongodb:mongoClient?database=test&collection=c&operation=findDistinct&dynamicity=false");
+ from("direct:findAll")
+
.to("mongodb:mongoClient?database=test&collection=c&operation=findAll&dynamicity=false");
+ from("direct:aggregate")
+
.to("mongodb:mongoClient?database=test&collection=c&operation=aggregate&dynamicity=false");
+ }
+ };
+ }
+}