This is an automated email from the ASF dual-hosted git repository.
davsclaus 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 b4aab1b04bd8 CAMEL-24814: camel-elasticsearch/camel-opensearch -
BulkRequestAggregationStrategy must return the new exchange
b4aab1b04bd8 is described below
commit b4aab1b04bd86b7d37aac285f6878144f0938f87
Author: Andrea Cosentino <[email protected]>
AuthorDate: Tue Sep 22 09:12:28 2026 +0200
CAMEL-24814: camel-elasticsearch/camel-opensearch -
BulkRequestAggregationStrategy must return the new exchange
ElasticsearchBulkRequestAggregationStrategy and
OpensearchBulkRequestAggregationStrategy
merge BulkOperation[] messages into a single BulkRequest for use with the
Aggregate EIP.
Both stored the merged request on newExchange but returned oldExchange,
which is null on
the first call of every aggregation group. AggregateProcessor rejects a
null return, so
the first exchange of every group failed and later returns never carried
the merged
request, leaving the strategy unusable. Neither component had a test for it.
Return newExchange from both strategies, add the already-aggregated
operations before
the new ones so the merged request preserves insertion order, and add unit
tests
covering the first call and a 3-step aggregation chain in both components.
Closes #26588
Co-authored-by: Claude <[email protected]>
---
...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 cd5ca8f91dcb..f1b7bca87c44 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)));
+ }
+}