This is an automated email from the ASF dual-hosted git repository.
gnodet pushed a commit to branch camel-4.18.x
in repository https://gitbox.apache.org/repos/asf/camel.git
The following commit(s) were added to refs/heads/camel-4.18.x by this push:
new da0dfd76a8ce [backport camel-4.18.x] CAMEL-24814:
camel-elasticsearch/camel-opensearch - BulkRequestAggregationStrategy must
return the new exchange (#26724)
da0dfd76a8ce is described below
commit da0dfd76a8ce94145afeca719f39b93b0860b557
Author: Guillaume Nodet <[email protected]>
AuthorDate: Tue Sep 22 18:01:28 2026 +0200
[backport camel-4.18.x] CAMEL-24814: camel-elasticsearch/camel-opensearch -
BulkRequestAggregationStrategy must return the new exchange (#26724)
---
...lasticsearchBulkRequestAggregationStrategy.java | 7 +-
...icsearchBulkRequestAggregationStrategyTest.java | 89 ++++++++++++++++++++++
.../OpensearchBulkRequestAggregationStrategy.java | 7 +-
...ensearchBulkRequestAggregationStrategyTest.java | 89 ++++++++++++++++++++++
4 files changed, 188 insertions(+), 4 deletions(-)
diff --git
a/components/camel-elasticsearch/src/main/java/org/apache/camel/component/es/aggregation/ElasticsearchBulkRequestAggregationStrategy.java
b/components/camel-elasticsearch/src/main/java/org/apache/camel/component/es/aggregation/ElasticsearchBulkRequestAggregationStrategy.java
index eb7d2e91b65a..fa56af0bb4bd 100644
---
a/components/camel-elasticsearch/src/main/java/org/apache/camel/component/es/aggregation/ElasticsearchBulkRequestAggregationStrategy.java
+++
b/components/camel-elasticsearch/src/main/java/org/apache/camel/component/es/aggregation/ElasticsearchBulkRequestAggregationStrategy.java
@@ -42,12 +42,15 @@ public class ElasticsearchBulkRequestAggregationStrategy
implements AggregationS
BulkOperation[] newBody = (BulkOperation[]) objBody;
BulkRequest.Builder builder = new BulkRequest.Builder();
- builder.operations(List.of(newBody));
if (oldExchange != null) {
+ // add the already-aggregated operations first so the merged
request keeps insertion order
BulkRequest request =
oldExchange.getIn().getBody(BulkRequest.class);
builder.operations(request.operations());
}
+ builder.operations(List.of(newBody));
+ // the merged BulkRequest is stored on the new exchange, so the new
exchange must be the one returned
+ // (returning oldExchange would return null on the first aggregation
call, which is not allowed)
newExchange.getIn().setBody(builder.build());
- return oldExchange;
+ return newExchange;
}
}
diff --git
a/components/camel-elasticsearch/src/test/java/org/apache/camel/component/es/aggregation/ElasticsearchBulkRequestAggregationStrategyTest.java
b/components/camel-elasticsearch/src/test/java/org/apache/camel/component/es/aggregation/ElasticsearchBulkRequestAggregationStrategyTest.java
new file mode 100644
index 000000000000..cbf4ea2e387c
--- /dev/null
+++
b/components/camel-elasticsearch/src/test/java/org/apache/camel/component/es/aggregation/ElasticsearchBulkRequestAggregationStrategyTest.java
@@ -0,0 +1,89 @@
+/*
+ * 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.es.aggregation;
+
+import java.util.List;
+import java.util.Map;
+
+import co.elastic.clients.elasticsearch.core.BulkRequest;
+import co.elastic.clients.elasticsearch.core.bulk.BulkOperation;
+import co.elastic.clients.elasticsearch.core.bulk.IndexOperation;
+import org.apache.camel.CamelContext;
+import org.apache.camel.Exchange;
+import org.apache.camel.impl.DefaultCamelContext;
+import org.apache.camel.support.DefaultExchange;
+import org.junit.jupiter.api.AfterEach;
+import org.junit.jupiter.api.BeforeEach;
+import org.junit.jupiter.api.Test;
+
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertNotNull;
+import static org.junit.jupiter.api.Assertions.assertSame;
+
+public class ElasticsearchBulkRequestAggregationStrategyTest {
+
+ private final ElasticsearchBulkRequestAggregationStrategy strategy = new
ElasticsearchBulkRequestAggregationStrategy();
+ private CamelContext context;
+
+ @BeforeEach
+ void setUp() {
+ context = new DefaultCamelContext();
+ }
+
+ @AfterEach
+ void tearDown() {
+ context.stop();
+ }
+
+ private Exchange exchangeWith(String id) {
+ Exchange exchange = new DefaultExchange(context);
+ BulkOperation operation = new BulkOperation.Builder()
+ .index(new
IndexOperation.Builder<>().index("idx").id(id).document(Map.of("id",
id)).build())
+ .build();
+ exchange.getIn().setBody(new BulkOperation[] { operation });
+ return exchange;
+ }
+
+ private static List<String> ids(BulkRequest request) {
+ return request.operations().stream().map(op ->
op.index().id()).toList();
+ }
+
+ @Test
+ void firstAggregationMustReturnNewExchangeCarryingTheRequest() {
+ Exchange newExchange = exchangeWith("1");
+
+ Exchange result = strategy.aggregate(null, newExchange);
+
+ // On the first call oldExchange is null; returning it (the previous
bug) would make the
+ // AggregationStrategy return null, which AggregateProcessor rejects.
+ assertSame(newExchange, result);
+ BulkRequest request = result.getIn().getBody(BulkRequest.class);
+ assertNotNull(request);
+ assertEquals(List.of("1"), ids(request));
+ }
+
+ @Test
+ void subsequentAggregationMergesAllOperationsInInsertionOrder() {
+ Exchange step1 = strategy.aggregate(null, exchangeWith("1"));
+ Exchange step2 = strategy.aggregate(step1, exchangeWith("2"));
+ Exchange step3 = strategy.aggregate(step2, exchangeWith("3"));
+
+ // a 3-step chain proves the merged request preserves insertion order
cumulatively
+ assertEquals(List.of("1", "2"),
ids(step2.getIn().getBody(BulkRequest.class)));
+ assertEquals(List.of("1", "2", "3"),
ids(step3.getIn().getBody(BulkRequest.class)));
+ }
+}
diff --git
a/components/camel-opensearch/src/main/java/org/apache/camel/component/opensearch/aggregation/OpensearchBulkRequestAggregationStrategy.java
b/components/camel-opensearch/src/main/java/org/apache/camel/component/opensearch/aggregation/OpensearchBulkRequestAggregationStrategy.java
index a4e15fd73296..914c0251027e 100644
---
a/components/camel-opensearch/src/main/java/org/apache/camel/component/opensearch/aggregation/OpensearchBulkRequestAggregationStrategy.java
+++
b/components/camel-opensearch/src/main/java/org/apache/camel/component/opensearch/aggregation/OpensearchBulkRequestAggregationStrategy.java
@@ -46,12 +46,15 @@ public class OpensearchBulkRequestAggregationStrategy
implements AggregationStra
BulkOperation[] newBody = (BulkOperation[]) objBody;
BulkRequest.Builder builder = new BulkRequest.Builder();
- builder.operations(List.of(newBody));
if (ObjectHelper.isNotEmpty(oldExchange)) {
+ // add the already-aggregated operations first so the merged
request keeps insertion order
BulkRequest request =
oldExchange.getIn().getBody(BulkRequest.class);
builder.operations(request.operations());
}
+ builder.operations(List.of(newBody));
+ // the merged BulkRequest is stored on the new exchange, so the new
exchange must be the one returned
+ // (returning oldExchange would return null on the first aggregation
call, which is not allowed)
newExchange.getIn().setBody(builder.build());
- return oldExchange;
+ return newExchange;
}
}
diff --git
a/components/camel-opensearch/src/test/java/org/apache/camel/component/opensearch/aggregation/OpensearchBulkRequestAggregationStrategyTest.java
b/components/camel-opensearch/src/test/java/org/apache/camel/component/opensearch/aggregation/OpensearchBulkRequestAggregationStrategyTest.java
new file mode 100644
index 000000000000..9c4b6d247c18
--- /dev/null
+++
b/components/camel-opensearch/src/test/java/org/apache/camel/component/opensearch/aggregation/OpensearchBulkRequestAggregationStrategyTest.java
@@ -0,0 +1,89 @@
+/*
+ * 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.opensearch.aggregation;
+
+import java.util.List;
+import java.util.Map;
+
+import org.apache.camel.CamelContext;
+import org.apache.camel.Exchange;
+import org.apache.camel.impl.DefaultCamelContext;
+import org.apache.camel.support.DefaultExchange;
+import org.junit.jupiter.api.AfterEach;
+import org.junit.jupiter.api.BeforeEach;
+import org.junit.jupiter.api.Test;
+import org.opensearch.client.opensearch.core.BulkRequest;
+import org.opensearch.client.opensearch.core.bulk.BulkOperation;
+import org.opensearch.client.opensearch.core.bulk.IndexOperation;
+
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertNotNull;
+import static org.junit.jupiter.api.Assertions.assertSame;
+
+public class OpensearchBulkRequestAggregationStrategyTest {
+
+ private final OpensearchBulkRequestAggregationStrategy strategy = new
OpensearchBulkRequestAggregationStrategy();
+ private CamelContext context;
+
+ @BeforeEach
+ void setUp() {
+ context = new DefaultCamelContext();
+ }
+
+ @AfterEach
+ void tearDown() {
+ context.stop();
+ }
+
+ private Exchange exchangeWith(String id) {
+ Exchange exchange = new DefaultExchange(context);
+ BulkOperation operation = new BulkOperation.Builder()
+ .index(new
IndexOperation.Builder<>().index("idx").id(id).document(Map.of("id",
id)).build())
+ .build();
+ exchange.getIn().setBody(new BulkOperation[] { operation });
+ return exchange;
+ }
+
+ private static List<String> ids(BulkRequest request) {
+ return request.operations().stream().map(op ->
op.index().id()).toList();
+ }
+
+ @Test
+ void firstAggregationMustReturnNewExchangeCarryingTheRequest() {
+ Exchange newExchange = exchangeWith("1");
+
+ Exchange result = strategy.aggregate(null, newExchange);
+
+ // On the first call oldExchange is null; returning it (the previous
bug) would make the
+ // AggregationStrategy return null, which AggregateProcessor rejects.
+ assertSame(newExchange, result);
+ BulkRequest request = result.getIn().getBody(BulkRequest.class);
+ assertNotNull(request);
+ assertEquals(List.of("1"), ids(request));
+ }
+
+ @Test
+ void subsequentAggregationMergesAllOperationsInInsertionOrder() {
+ Exchange step1 = strategy.aggregate(null, exchangeWith("1"));
+ Exchange step2 = strategy.aggregate(step1, exchangeWith("2"));
+ Exchange step3 = strategy.aggregate(step2, exchangeWith("3"));
+
+ // a 3-step chain proves the merged request preserves insertion order
cumulatively
+ assertEquals(List.of("1", "2"),
ids(step2.getIn().getBody(BulkRequest.class)));
+ assertEquals(List.of("1", "2", "3"),
ids(step3.getIn().getBody(BulkRequest.class)));
+ }
+}