This is an automated email from the ASF dual-hosted git repository.
davsclaus pushed a commit to branch camel-4.22.x
in repository https://gitbox.apache.org/repos/asf/camel.git
The following commit(s) were added to refs/heads/camel-4.22.x by this push:
new 289da267e4be [backport camel-4.22.x] CAMEL-24816:
camel-elasticsearch/camel-opensearch - minor robustness fixes (#26722)
289da267e4be is described below
commit 289da267e4be48f511dd30d2e38c875e1584103c
Author: Guillaume Nodet <[email protected]>
AuthorDate: Tue Sep 22 10:43:58 2026 +0200
[backport camel-4.22.x] CAMEL-24816: camel-elasticsearch/camel-opensearch -
minor robustness fixes (#26722)
---
.../camel/component/es/ElasticsearchProducer.java | 27 ++++++-
.../es/ElasticsearchScrollRequestIterator.java | 13 ++--
.../ElasticsearchActionRequestConverter.java | 16 +++-
.../ElasticsearchActionRequestConverterTest.java | 86 ++++++++++++++++++++++
.../component/opensearch/OpensearchProducer.java | 15 +++-
.../OpensearchScrollRequestIterator.java | 13 ++--
.../OpensearchActionRequestConverter.java | 14 +++-
.../OpensearchActionRequestConverterTest.java | 86 ++++++++++++++++++++++
8 files changed, 251 insertions(+), 19 deletions(-)
diff --git
a/components/camel-elasticsearch/src/main/java/org/apache/camel/component/es/ElasticsearchProducer.java
b/components/camel-elasticsearch/src/main/java/org/apache/camel/component/es/ElasticsearchProducer.java
index 5b60d9a543e8..20027487a64a 100644
---
a/components/camel-elasticsearch/src/main/java/org/apache/camel/component/es/ElasticsearchProducer.java
+++
b/components/camel-elasticsearch/src/main/java/org/apache/camel/component/es/ElasticsearchProducer.java
@@ -166,14 +166,18 @@ class ElasticsearchProducer extends DefaultAsyncProducer {
configIndexName = true;
}
+ boolean configSize = false;
Integer size =
message.getHeader(ElasticsearchConstants.PARAM_SIZE, Integer.class);
if (size == null) {
message.setHeader(ElasticsearchConstants.PARAM_SIZE,
configuration.getSize());
+ configSize = true;
}
+ boolean configFrom = false;
Integer from =
message.getHeader(ElasticsearchConstants.PARAM_FROM, Integer.class);
if (from == null) {
message.setHeader(ElasticsearchConstants.PARAM_FROM,
configuration.getFrom());
+ configFrom = true;
}
Boolean enableDocumentOnlyMode =
message.getHeader(ElasticsearchConstants.PARAM_DOCUMENT_MODE, Boolean.class);
@@ -193,7 +197,8 @@ class ElasticsearchProducer extends DefaultAsyncProducer {
documentClass = configuration.getDocumentClass();
}
- ActionContext ctx = new ActionContext(exchange, callback,
transport, configIndexName, configWaitForActiveShards);
+ ActionContext ctx = new ActionContext(
+ exchange, callback, transport, configIndexName,
configWaitForActiveShards, configSize, configFrom);
switch (operation) {
case Index: {
@@ -455,6 +460,12 @@ class ElasticsearchProducer extends DefaultAsyncProducer {
if (ctx.isConfigWaitForActiveShards()) {
message.removeHeader(ElasticsearchConstants.PARAM_WAIT_FOR_ACTIVE_SHARDS);
}
+ if (ctx.isConfigSize()) {
+ message.removeHeader(ElasticsearchConstants.PARAM_SIZE);
+ }
+ if (ctx.isConfigFrom()) {
+ message.removeHeader(ElasticsearchConstants.PARAM_FROM);
+ }
if (configuration.isDisconnect()) {
IOHelper.close(ctx.getTransport());
if (configuration.isEnableSniffer()) {
@@ -586,14 +597,18 @@ class ElasticsearchProducer extends DefaultAsyncProducer {
private final ElasticsearchTransport transport;
private final boolean configIndexName;
private final boolean configWaitForActiveShards;
+ private final boolean configSize;
+ private final boolean configFrom;
ActionContext(Exchange exchange, AsyncCallback callback,
ElasticsearchTransport transport, boolean configIndexName,
- boolean configWaitForActiveShards) {
+ boolean configWaitForActiveShards, boolean configSize,
boolean configFrom) {
this.exchange = exchange;
this.callback = callback;
this.transport = transport;
this.configIndexName = configIndexName;
this.configWaitForActiveShards = configWaitForActiveShards;
+ this.configSize = configSize;
+ this.configFrom = configFrom;
}
ElasticsearchTransport getTransport() {
@@ -612,6 +627,14 @@ class ElasticsearchProducer extends DefaultAsyncProducer {
return configWaitForActiveShards;
}
+ boolean isConfigSize() {
+ return configSize;
+ }
+
+ boolean isConfigFrom() {
+ return configFrom;
+ }
+
Exchange getExchange() {
return exchange;
}
diff --git
a/components/camel-elasticsearch/src/main/java/org/apache/camel/component/es/ElasticsearchScrollRequestIterator.java
b/components/camel-elasticsearch/src/main/java/org/apache/camel/component/es/ElasticsearchScrollRequestIterator.java
index b4ebf0d51910..b336998609e8 100644
---
a/components/camel-elasticsearch/src/main/java/org/apache/camel/component/es/ElasticsearchScrollRequestIterator.java
+++
b/components/camel-elasticsearch/src/main/java/org/apache/camel/component/es/ElasticsearchScrollRequestIterator.java
@@ -125,11 +125,14 @@ public class
ElasticsearchScrollRequestIterator<TDocument> implements Iterator<H
public void close() {
if (!closed) {
try {
- ClearScrollRequest clearScrollRequest = new
ClearScrollRequest.Builder()
- .scrollId(List.of(scrollId))
- .build();
-
- esClient.clearScroll(clearScrollRequest);
+ // scrollId can be null if the initial search returned no
scroll id; List.of(null) would NPE
+ if (scrollId != null) {
+ ClearScrollRequest clearScrollRequest = new
ClearScrollRequest.Builder()
+ .scrollId(List.of(scrollId))
+ .build();
+
+ esClient.clearScroll(clearScrollRequest);
+ }
closed = true;
exchange.setProperty(ElasticsearchConstants.PROPERTY_SCROLL_ES_QUERY_COUNT,
requestCount);
} catch (IOException e) {
diff --git
a/components/camel-elasticsearch/src/main/java/org/apache/camel/component/es/converter/ElasticsearchActionRequestConverter.java
b/components/camel-elasticsearch/src/main/java/org/apache/camel/component/es/converter/ElasticsearchActionRequestConverter.java
index 90b2fc8ff76e..7e18c719dd8e 100644
---
a/components/camel-elasticsearch/src/main/java/org/apache/camel/component/es/converter/ElasticsearchActionRequestConverter.java
+++
b/components/camel-elasticsearch/src/main/java/org/apache/camel/component/es/converter/ElasticsearchActionRequestConverter.java
@@ -85,7 +85,12 @@ public final class ElasticsearchActionRequestConverter {
@Converter
public static IndexRequest.Builder<?> toIndexRequestBuilder(Object
document, Exchange exchange) throws IOException {
if (document instanceof IndexRequest.Builder<?> indexReqBuilder) {
- return
indexReqBuilder.id(exchange.getIn().getHeader(ElasticsearchConstants.PARAM_INDEX_ID,
String.class));
+ // only override the id when the header is present, otherwise a
caller-supplied id would be cleared
+ String id =
exchange.getIn().getHeader(ElasticsearchConstants.PARAM_INDEX_ID, String.class);
+ if (id != null) {
+ indexReqBuilder.id(id);
+ }
+ return indexReqBuilder;
}
IndexRequest.Builder<Object> builder = new IndexRequest.Builder<>();
if (document instanceof byte[] byteArray) {
@@ -115,12 +120,17 @@ public final class ElasticsearchActionRequestConverter {
@Converter
public static UpdateRequest.Builder<?, ?> toUpdateRequestBuilder(Object
document, Exchange exchange) throws IOException {
if (document instanceof UpdateRequest.Builder<?, ?> updateReqBuilder) {
- return
updateReqBuilder.id(exchange.getIn().getHeader(ElasticsearchConstants.PARAM_INDEX_ID,
String.class));
+ // only override the id when the header is present, otherwise a
caller-supplied id would be cleared
+ String id =
exchange.getIn().getHeader(ElasticsearchConstants.PARAM_INDEX_ID, String.class);
+ if (id != null) {
+ updateReqBuilder.id(id);
+ }
+ return updateReqBuilder;
}
UpdateRequest.Builder<?, Object> builder = new
UpdateRequest.Builder<>();
Boolean enableDocumentOnlyMode
=
exchange.getIn().getHeader(ElasticsearchConstants.PARAM_DOCUMENT_MODE,
Boolean.FALSE, Boolean.class);
- Mode mode = enableDocumentOnlyMode == Boolean.TRUE ?
Mode.DOCUMENT_ONLY : Mode.DEFAULT;
+ Mode mode = Boolean.TRUE.equals(enableDocumentOnlyMode) ?
Mode.DOCUMENT_ONLY : Mode.DEFAULT;
if (document instanceof byte[] byteArray) {
mode.addDocToUpdateRequestBuilder(builder, new
ByteArrayInputStream(byteArray));
} else if (document instanceof InputStream inputStream) {
diff --git
a/components/camel-elasticsearch/src/test/java/org/apache/camel/component/es/converter/ElasticsearchActionRequestConverterTest.java
b/components/camel-elasticsearch/src/test/java/org/apache/camel/component/es/converter/ElasticsearchActionRequestConverterTest.java
new file mode 100644
index 000000000000..796e859980b3
--- /dev/null
+++
b/components/camel-elasticsearch/src/test/java/org/apache/camel/component/es/converter/ElasticsearchActionRequestConverterTest.java
@@ -0,0 +1,86 @@
+/*
+ * 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.converter;
+
+import java.util.Map;
+
+import co.elastic.clients.elasticsearch.core.IndexRequest;
+import co.elastic.clients.elasticsearch.core.UpdateRequest;
+import org.apache.camel.CamelContext;
+import org.apache.camel.Exchange;
+import org.apache.camel.component.es.ElasticsearchConstants;
+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;
+
+public class ElasticsearchActionRequestConverterTest {
+
+ private CamelContext context;
+
+ @BeforeEach
+ void setUp() {
+ context = new DefaultCamelContext();
+ }
+
+ @AfterEach
+ void tearDown() {
+ context.stop();
+ }
+
+ private IndexRequest.Builder<Object> preBuiltIndexBuilder() {
+ return new
IndexRequest.Builder<>().index("idx").id("original").document(Map.of("k", "v"));
+ }
+
+ @Test
+ void preBuiltIndexBuilderKeepsItsIdWhenHeaderAbsent() throws Exception {
+ Exchange exchange = new DefaultExchange(context);
+
+ IndexRequest.Builder<?> result
+ =
ElasticsearchActionRequestConverter.toIndexRequestBuilder(preBuiltIndexBuilder(),
exchange);
+
+ // no CamelIndexId header -> the caller's id must be preserved, not
overwritten with null
+ assertEquals("original", result.build().id());
+ }
+
+ @Test
+ void preBuiltIndexBuilderHeaderOverridesId() throws Exception {
+ Exchange exchange = new DefaultExchange(context);
+ exchange.getIn().setHeader(ElasticsearchConstants.PARAM_INDEX_ID,
"fromHeader");
+
+ IndexRequest.Builder<?> result
+ =
ElasticsearchActionRequestConverter.toIndexRequestBuilder(preBuiltIndexBuilder(),
exchange);
+
+ assertEquals("fromHeader", result.build().id());
+ }
+
+ @Test
+ void preBuiltUpdateBuilderKeepsItsIdWhenHeaderAbsent() throws Exception {
+ Exchange exchange = new DefaultExchange(context);
+ UpdateRequest.Builder<Object, Object> preBuilt
+ = new
UpdateRequest.Builder<>().index("idx").id("original").doc(Map.of("k", "v"));
+
+ UpdateRequest.Builder<?, ?> result
+ =
ElasticsearchActionRequestConverter.toUpdateRequestBuilder(preBuilt, exchange);
+
+ // no CamelIndexId header -> the caller's id must be preserved, not
overwritten with null
+ assertEquals("original", result.build().id());
+ }
+}
diff --git
a/components/camel-opensearch/src/main/java/org/apache/camel/component/opensearch/OpensearchProducer.java
b/components/camel-opensearch/src/main/java/org/apache/camel/component/opensearch/OpensearchProducer.java
index f90f6314a1ca..2d96335235bf 100644
---
a/components/camel-opensearch/src/main/java/org/apache/camel/component/opensearch/OpensearchProducer.java
+++
b/components/camel-opensearch/src/main/java/org/apache/camel/component/opensearch/OpensearchProducer.java
@@ -189,14 +189,18 @@ class OpensearchProducer extends DefaultAsyncProducer {
configIndexName = true;
}
+ boolean configSize = false;
Integer size = message.getHeader(OpensearchConstants.PARAM_SIZE,
Integer.class);
if (size == null) {
message.setHeader(OpensearchConstants.PARAM_SIZE,
configuration.getSize());
+ configSize = true;
}
+ boolean configFrom = false;
Integer from = message.getHeader(OpensearchConstants.PARAM_FROM,
Integer.class);
if (from == null) {
message.setHeader(OpensearchConstants.PARAM_FROM,
configuration.getFrom());
+ configFrom = true;
}
boolean configWaitForActiveShards = false;
@@ -211,7 +215,8 @@ class OpensearchProducer extends DefaultAsyncProducer {
documentClass = configuration.getDocumentClass();
}
- ActionContext ctx = new ActionContext(exchange, callback,
transport, configIndexName, configWaitForActiveShards);
+ ActionContext ctx = new ActionContext(
+ exchange, callback, transport, configIndexName,
configWaitForActiveShards, configSize, configFrom);
switch (operation) {
case Index -> processIndexAsync(ctx);
@@ -441,6 +446,12 @@ class OpensearchProducer extends DefaultAsyncProducer {
if (ctx.configWaitForActiveShards()) {
message.removeHeader(OpensearchConstants.PARAM_WAIT_FOR_ACTIVE_SHARDS);
}
+ if (ctx.configSize()) {
+ message.removeHeader(OpensearchConstants.PARAM_SIZE);
+ }
+ if (ctx.configFrom()) {
+ message.removeHeader(OpensearchConstants.PARAM_FROM);
+ }
if (configuration.isDisconnect() && openSearchClient == null) {
IOHelper.close(ctx.transport());
if (configuration.isEnableSniffer()) {
@@ -604,7 +615,7 @@ class OpensearchProducer extends DefaultAsyncProducer {
* An inner class providing all the information that an asynchronous
action could need.
*/
private record ActionContext(Exchange exchange, AsyncCallback callback,
OpenSearchTransport transport,
- boolean configIndexName, boolean configWaitForActiveShards) {
+ boolean configIndexName, boolean configWaitForActiveShards,
boolean configSize, boolean configFrom) {
OpenSearchAsyncClient getClient() {
return new OpenSearchAsyncClient(transport);
diff --git
a/components/camel-opensearch/src/main/java/org/apache/camel/component/opensearch/OpensearchScrollRequestIterator.java
b/components/camel-opensearch/src/main/java/org/apache/camel/component/opensearch/OpensearchScrollRequestIterator.java
index 70e19c9c0880..5c70ac93c61c 100644
---
a/components/camel-opensearch/src/main/java/org/apache/camel/component/opensearch/OpensearchScrollRequestIterator.java
+++
b/components/camel-opensearch/src/main/java/org/apache/camel/component/opensearch/OpensearchScrollRequestIterator.java
@@ -125,11 +125,14 @@ public class OpensearchScrollRequestIterator<TDocument>
implements Iterator<Hit<
public void close() {
if (!closed) {
try {
- ClearScrollRequest clearScrollRequest = new
ClearScrollRequest.Builder()
- .scrollId(List.of(scrollId))
- .build();
-
- esClient.clearScroll(clearScrollRequest);
+ // scrollId can be null if the initial search returned no
scroll id; List.of(null) would NPE
+ if (scrollId != null) {
+ ClearScrollRequest clearScrollRequest = new
ClearScrollRequest.Builder()
+ .scrollId(List.of(scrollId))
+ .build();
+
+ esClient.clearScroll(clearScrollRequest);
+ }
closed = true;
exchange.setProperty(OpensearchConstants.PROPERTY_SCROLL_OPENSEARCH_QUERY_COUNT,
requestCount);
} catch (IOException e) {
diff --git
a/components/camel-opensearch/src/main/java/org/apache/camel/component/opensearch/converter/OpensearchActionRequestConverter.java
b/components/camel-opensearch/src/main/java/org/apache/camel/component/opensearch/converter/OpensearchActionRequestConverter.java
index 512c2aa118e8..bdb6c711ec9c 100644
---
a/components/camel-opensearch/src/main/java/org/apache/camel/component/opensearch/converter/OpensearchActionRequestConverter.java
+++
b/components/camel-opensearch/src/main/java/org/apache/camel/component/opensearch/converter/OpensearchActionRequestConverter.java
@@ -88,7 +88,12 @@ public final class OpensearchActionRequestConverter {
@Converter
public static IndexRequest.Builder<?> toIndexRequestBuilder(Object
document, Exchange exchange) throws IOException {
if (document instanceof IndexRequest.Builder<?> builder) {
- return
builder.id(exchange.getIn().getHeader(OpensearchConstants.PARAM_INDEX_ID,
String.class));
+ // only override the id when the header is present, otherwise a
caller-supplied id would be cleared
+ String id =
exchange.getIn().getHeader(OpensearchConstants.PARAM_INDEX_ID, String.class);
+ if (id != null) {
+ builder.id(id);
+ }
+ return builder;
}
JacksonJsonpMapper mapper = createMapper();
IndexRequest.Builder<Object> builder = new IndexRequest.Builder<>();
@@ -116,7 +121,12 @@ public final class OpensearchActionRequestConverter {
@Converter
public static UpdateRequest.Builder<?, ?> toUpdateRequestBuilder(Object
document, Exchange exchange) throws IOException {
if (document instanceof UpdateRequest.Builder<?, ?> builder) {
- return
builder.id(exchange.getIn().getHeader(OpensearchConstants.PARAM_INDEX_ID,
String.class));
+ // only override the id when the header is present, otherwise a
caller-supplied id would be cleared
+ String id =
exchange.getIn().getHeader(OpensearchConstants.PARAM_INDEX_ID, String.class);
+ if (id != null) {
+ builder.id(id);
+ }
+ return builder;
}
JacksonJsonpMapper mapper = createMapper();
UpdateRequest.Builder<?, Object> builder = new
UpdateRequest.Builder<>();
diff --git
a/components/camel-opensearch/src/test/java/org/apache/camel/component/opensearch/converter/OpensearchActionRequestConverterTest.java
b/components/camel-opensearch/src/test/java/org/apache/camel/component/opensearch/converter/OpensearchActionRequestConverterTest.java
new file mode 100644
index 000000000000..87dede478c7d
--- /dev/null
+++
b/components/camel-opensearch/src/test/java/org/apache/camel/component/opensearch/converter/OpensearchActionRequestConverterTest.java
@@ -0,0 +1,86 @@
+/*
+ * 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.converter;
+
+import java.util.Map;
+
+import org.apache.camel.CamelContext;
+import org.apache.camel.Exchange;
+import org.apache.camel.component.opensearch.OpensearchConstants;
+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.IndexRequest;
+import org.opensearch.client.opensearch.core.UpdateRequest;
+
+import static org.junit.jupiter.api.Assertions.assertEquals;
+
+public class OpensearchActionRequestConverterTest {
+
+ private CamelContext context;
+
+ @BeforeEach
+ void setUp() {
+ context = new DefaultCamelContext();
+ }
+
+ @AfterEach
+ void tearDown() {
+ context.stop();
+ }
+
+ private IndexRequest.Builder<Object> preBuiltIndexBuilder() {
+ return new
IndexRequest.Builder<>().index("idx").id("original").document(Map.of("k", "v"));
+ }
+
+ @Test
+ void preBuiltIndexBuilderKeepsItsIdWhenHeaderAbsent() throws Exception {
+ Exchange exchange = new DefaultExchange(context);
+
+ IndexRequest.Builder<?> result
+ =
OpensearchActionRequestConverter.toIndexRequestBuilder(preBuiltIndexBuilder(),
exchange);
+
+ // no CamelOpensearchIndexId header -> the caller's id must be
preserved, not overwritten with null
+ assertEquals("original", result.build().id());
+ }
+
+ @Test
+ void preBuiltIndexBuilderHeaderOverridesId() throws Exception {
+ Exchange exchange = new DefaultExchange(context);
+ exchange.getIn().setHeader(OpensearchConstants.PARAM_INDEX_ID,
"fromHeader");
+
+ IndexRequest.Builder<?> result
+ =
OpensearchActionRequestConverter.toIndexRequestBuilder(preBuiltIndexBuilder(),
exchange);
+
+ assertEquals("fromHeader", result.build().id());
+ }
+
+ @Test
+ void preBuiltUpdateBuilderKeepsItsIdWhenHeaderAbsent() throws Exception {
+ Exchange exchange = new DefaultExchange(context);
+ UpdateRequest.Builder<Object, Object> preBuilt
+ = new
UpdateRequest.Builder<>().index("idx").id("original").doc(Map.of("k", "v"));
+
+ UpdateRequest.Builder<?, ?> result
+ =
OpensearchActionRequestConverter.toUpdateRequestBuilder(preBuilt, exchange);
+
+ // no CamelOpensearchIndexId header -> the caller's id must be
preserved, not overwritten with null
+ assertEquals("original", result.build().id());
+ }
+}