[SB3] Update Elasticsearch Client. Update Boot and Spring versions. Improve Rabbit listener (#416)

This commit is contained in:
Corneil du Plessis
2022-11-25 11:51:53 +02:00
committed by GitHub
parent f44dc56f50
commit 1459bfabd9
3 changed files with 266 additions and 261 deletions

View File

@@ -10,7 +10,7 @@
</parent>
<properties>
<elasticsearch.version>7.15.2</elasticsearch.version>
<elasticsearch.version>8.5.0</elasticsearch.version>
</properties>
<artifactId>elasticsearch-consumer</artifactId>
@@ -19,13 +19,8 @@
<dependencies>
<dependency>
<groupId>org.elasticsearch</groupId>
<artifactId>elasticsearch</artifactId>
<version>${elasticsearch.version}</version>
</dependency>
<dependency>
<groupId>org.elasticsearch.client</groupId>
<artifactId>elasticsearch-rest-high-level-client</artifactId>
<groupId>co.elastic.clients</groupId>
<artifactId>elasticsearch-java</artifactId>
<version>${elasticsearch.version}</version>
</dependency>
<dependency>

View File

@@ -17,27 +17,25 @@
package org.springframework.cloud.fn.consumer.elasticsearch;
import java.io.IOException;
import java.io.StringReader;
import java.util.ArrayList;
import java.util.Collection;
import java.util.List;
import java.util.Map;
import java.util.concurrent.CompletableFuture;
import java.util.function.Consumer;
import java.util.stream.StreamSupport;
import co.elastic.clients.elasticsearch.ElasticsearchAsyncClient;
import co.elastic.clients.elasticsearch.ElasticsearchClient;
import co.elastic.clients.elasticsearch._types.Time;
import co.elastic.clients.elasticsearch.core.BulkRequest;
import co.elastic.clients.elasticsearch.core.BulkResponse;
import co.elastic.clients.elasticsearch.core.IndexRequest;
import co.elastic.clients.elasticsearch.core.IndexResponse;
import co.elastic.clients.elasticsearch.core.bulk.BulkResponseItem;
import org.apache.commons.logging.Log;
import org.apache.commons.logging.LogFactory;
import org.elasticsearch.action.ActionListener;
import org.elasticsearch.action.DocWriteResponse;
import org.elasticsearch.action.bulk.BulkItemResponse;
import org.elasticsearch.action.bulk.BulkRequest;
import org.elasticsearch.action.bulk.BulkResponse;
import org.elasticsearch.action.index.IndexRequest;
import org.elasticsearch.action.index.IndexResponse;
import org.elasticsearch.client.RequestOptions;
import org.elasticsearch.client.RestHighLevelClient;
import org.elasticsearch.common.xcontent.XContentBuilder;
import org.elasticsearch.common.xcontent.XContentType;
import org.elasticsearch.core.TimeValue;
import org.springframework.beans.factory.FactoryBean;
import org.springframework.beans.factory.annotation.Qualifier;
@@ -120,11 +118,14 @@ public class ElasticsearchConsumerConfiguration {
}
@Bean
IntegrationFlow elasticsearchConsumerFlow(@Qualifier("aggregator") MessageHandler aggregator, ElasticsearchConsumerProperties properties,
@Qualifier("indexingHandler") MessageHandler indexingHandler) {
IntegrationFlow elasticsearchConsumerFlow(
@Qualifier("aggregator") MessageHandler aggregator,
ElasticsearchConsumerProperties properties,
@Qualifier("indexingHandler") MessageHandler indexingHandler
) {
final IntegrationFlowBuilder builder =
IntegrationFlows.from(Consumer.class, gateway -> gateway.beanName("elasticsearchConsumer"));
IntegrationFlows.from(Consumer.class, gateway -> gateway.beanName("elasticsearchConsumer"));
if (properties.getBatchSize() > 1) {
builder.handle(aggregator);
}
@@ -132,34 +133,38 @@ public class ElasticsearchConsumerConfiguration {
}
@Bean
public MessageHandler indexingHandler(RestHighLevelClient restHighLevelClient,
ElasticsearchConsumerProperties consumerProperties) {
public MessageHandler indexingHandler(
ElasticsearchClient elasticsearchClient,
ElasticsearchConsumerProperties consumerProperties
) {
return message -> {
if (message.getPayload() instanceof Iterable) {
BulkRequest bulkRequest = new BulkRequest();
BulkRequest.Builder builder = new BulkRequest.Builder();
StreamSupport.stream(((Iterable<?>) message.getPayload()).spliterator(), false)
.filter(MessageWrapper.class::isInstance)
.map(itemPayload -> ((MessageWrapper) itemPayload).getMessage())
.map(m -> buildIndexRequest(m, consumerProperties))
.forEach(bulkRequest::add);
.filter(MessageWrapper.class::isInstance)
.map(itemPayload -> ((MessageWrapper) itemPayload).getMessage())
.map(m -> buildIndexRequest(m, consumerProperties))
.forEach(indexRequest ->
builder.operations(builder1 -> builder1.index(idx -> idx.index(indexRequest.index()).id(indexRequest.id()).document(indexRequest.document())))
);
index(restHighLevelClient, bulkRequest, consumerProperties.isAsync());
index(elasticsearchClient, builder.build(), consumerProperties.isAsync());
}
else {
IndexRequest request = buildIndexRequest(message, consumerProperties);
index(restHighLevelClient, request, consumerProperties.isAsync());
index(elasticsearchClient, request, consumerProperties.isAsync());
}
};
}
private IndexRequest buildIndexRequest(Message<?> message, ElasticsearchConsumerProperties consumerProperties) {
IndexRequest request = new IndexRequest();
IndexRequest.Builder requestBuilder = new IndexRequest.Builder();
String index = consumerProperties.getIndex();
if (message.getHeaders().containsKey(INDEX_NAME_HEADER)) {
index = (String) message.getHeaders().get(INDEX_NAME_HEADER);
}
request.index(index);
requestBuilder.index(index);
String id = "";
if (message.getHeaders().containsKey(INDEX_ID_HEADER)) {
@@ -168,45 +173,42 @@ public class ElasticsearchConsumerConfiguration {
else if (consumerProperties.getId() != null) {
id = consumerProperties.getId().getValue(message, String.class);
}
request.id(id);
requestBuilder.id(id);
if (message.getPayload() instanceof String) {
request.source((String) message.getPayload(), XContentType.JSON);
requestBuilder.withJson(new StringReader((String) message.getPayload()));
}
else if (message.getPayload() instanceof Map) {
request.source((Map<String, ?>) message.getPayload(), XContentType.JSON);
}
else if (message.getPayload() instanceof XContentBuilder) {
request.source((XContentBuilder) message.getPayload());
requestBuilder.document((Map<String, ?>) message.getPayload());
}
if (!StringUtils.isEmpty(consumerProperties.getRouting())) {
request.routing(consumerProperties.getRouting());
if (StringUtils.hasText(consumerProperties.getRouting())) {
requestBuilder.routing(consumerProperties.getRouting());
}
if (consumerProperties.getTimeoutSeconds() > 0) {
request.timeout(TimeValue.timeValueSeconds(consumerProperties.getTimeoutSeconds()));
requestBuilder.timeout(new Time.Builder().time(consumerProperties.getTimeoutSeconds() + "s").build());
}
return request;
return requestBuilder.build();
}
private void index(RestHighLevelClient restHighLevelClient, BulkRequest request, boolean isAsync) {
private void index(ElasticsearchClient elasticsearchClient, BulkRequest request, boolean isAsync) {
if (isAsync) {
restHighLevelClient.bulkAsync(request, RequestOptions.DEFAULT, new ActionListener<BulkResponse>() {
@Override
public void onResponse(BulkResponse bulkResponse) {
ElasticsearchAsyncClient elasticsearchAsyncClient = new ElasticsearchAsyncClient(elasticsearchClient._transport());
CompletableFuture<BulkResponse> responseCompletableFuture = elasticsearchAsyncClient.bulk(request);
responseCompletableFuture.whenComplete((bulkResponse, x) -> {
if (x != null) {
throw new IllegalStateException("Error occurred while performing bulk index operation: " + x.getMessage(), x);
}
else {
handleBulkResponse(bulkResponse);
}
@Override
public void onFailure(Exception e) {
throw new IllegalStateException("Error occurred while performing bulk index operation: " + e.getMessage(), e);
}
});
}
else {
try {
BulkResponse bulkResponse = restHighLevelClient.bulk(request, RequestOptions.DEFAULT);
BulkResponse bulkResponse = elasticsearchClient.bulk(request);
handleBulkResponse(bulkResponse);
}
catch (IOException e) {
@@ -215,23 +217,22 @@ public class ElasticsearchConsumerConfiguration {
}
}
private void index(RestHighLevelClient restHighLevelClient, IndexRequest request, boolean isAsync) {
private void index(ElasticsearchClient elasticsearchClient, IndexRequest request, boolean isAsync) {
if (isAsync) {
restHighLevelClient.indexAsync(request, RequestOptions.DEFAULT, new ActionListener<IndexResponse>() {
@Override
public void onResponse(IndexResponse indexResponse) {
handleResponse(indexResponse);
ElasticsearchAsyncClient elasticsearchAsyncClient = new ElasticsearchAsyncClient(elasticsearchClient._transport());
CompletableFuture<IndexResponse> responseCompletableFuture = elasticsearchAsyncClient.index(request);
responseCompletableFuture.whenComplete((indexResponse, x) -> {
if (x != null) {
throw new IllegalStateException("Error occurred while indexing document: " + x.getMessage(), x);
}
@Override
public void onFailure(Exception e) {
throw new IllegalStateException("Error occurred while indexing document: " + e.getMessage(), e);
else {
handleResponse(indexResponse);
}
});
}
else {
try {
IndexResponse response = restHighLevelClient.index(request, RequestOptions.DEFAULT);
IndexResponse response = elasticsearchClient.index(request);
handleResponse(response);
}
catch (IOException e) {
@@ -241,30 +242,46 @@ public class ElasticsearchConsumerConfiguration {
}
private void handleBulkResponse(BulkResponse response) {
if (logger.isDebugEnabled() || response.hasFailures()) {
for (BulkItemResponse itemResponse : response) {
if (itemResponse.isFailed()) {
logger.error(String.format("Index operation [i=%d, id=%s, index=%s] failed: %s",
itemResponse.getItemId(), itemResponse.getId(), itemResponse.getIndex(), itemResponse.getFailureMessage())
if (logger.isDebugEnabled() || response.errors()) {
for (BulkResponseItem itemResponse : response.items()) {
if (itemResponse.error() != null) {
if (logger.isDebugEnabled()) {
logger.debug("itemResponse.error=" + itemResponse.error());
}
logger.error(String.format("Index operation [id=%s, index=%s] failed: %s",
itemResponse.id(), itemResponse.index(), itemResponse.error().toString())
);
}
else {
DocWriteResponse r = itemResponse.getResponse();
logger.debug(String.format("Index operation [i=%d, id=%s, index=%s] succeeded: document [id=%s, version=%d] was written on shard %s.",
itemResponse.getItemId(), itemResponse.getId(), itemResponse.getIndex(), r.getId(), r.getVersion(), r.getShardId())
);
var r = itemResponse.get();
if (r != null) {
if (logger.isDebugEnabled()) {
logger.debug("itemResponse:" + r);
}
logger.debug(String.format("Index operation [id=%s, index=%s] succeeded: document [id=%s, version=%s] was written on shard %s.",
itemResponse.id(), itemResponse.index(), r.source().get("id"), r.source().get("version"), r.source().get("shardId"))
);
}
else {
logger.debug(String.format("Index operation [id=%s, index=%s] succeeded", itemResponse.id(), itemResponse.index()));
}
}
}
}
if (response.hasFailures()) {
throw new IllegalStateException("Bulk indexing operation completed with failures: " + response.buildFailureMessage());
if (response.errors()) {
String error = response.items()
.stream()
.map(bulkResponseItem -> bulkResponseItem.error() != null ? bulkResponseItem.error().toString() : "")
.reduce((errorCause, errorCause2) -> errorCause != null ? errorCause + " : " + errorCause2 : errorCause2)
.orElseGet(() -> response.toString());
throw new IllegalStateException("Bulk indexing operation completed with failures: " + error);
}
}
private void handleResponse(IndexResponse response) {
logger.debug(String.format("Index operation [index=%s] succeeded: document [id=%s, version=%d] was written on shard %s.",
response.getIndex(), response.getId(), response.getVersion(), response.getShardId())
response.index(), response.id(), response.version(), response.shards().toString())
);
}

View File

@@ -22,25 +22,26 @@ import java.util.Map;
import java.util.UUID;
import java.util.function.Consumer;
import co.elastic.clients.elasticsearch.ElasticsearchClient;
import co.elastic.clients.elasticsearch._types.ElasticsearchException;
import co.elastic.clients.elasticsearch.core.GetRequest;
import co.elastic.clients.elasticsearch.core.GetResponse;
import co.elastic.clients.json.JsonData;
import org.awaitility.Awaitility;
import org.elasticsearch.ElasticsearchStatusException;
import org.elasticsearch.action.get.GetRequest;
import org.elasticsearch.action.get.GetResponse;
import org.elasticsearch.client.RequestOptions;
import org.elasticsearch.client.RestHighLevelClient;
import org.elasticsearch.common.Strings;
import org.elasticsearch.common.xcontent.XContentBuilder;
import org.elasticsearch.common.xcontent.XContentFactory;
import org.junit.jupiter.api.Disabled;
import org.junit.jupiter.api.Tag;
import org.junit.jupiter.api.Test;
import org.testcontainers.elasticsearch.ElasticsearchContainer;
import org.testcontainers.junit.jupiter.Container;
import org.testcontainers.junit.jupiter.Testcontainers;
import org.testcontainers.utility.DockerImageName;
import org.springframework.boot.autoconfigure.SpringBootApplication;
import org.springframework.boot.test.context.runner.ApplicationContextRunner;
import org.springframework.cloud.fn.common.config.SpelExpressionConverterConfiguration;
import org.springframework.context.annotation.Configuration;
import org.springframework.data.elasticsearch.client.ClientConfiguration;
import org.springframework.data.elasticsearch.client.elc.ElasticsearchConfiguration;
import org.springframework.lang.NonNull;
import org.springframework.messaging.Message;
import org.springframework.messaging.support.MessageBuilder;
@@ -53,254 +54,246 @@ import static org.assertj.core.api.Assertions.assertThatIllegalStateException;
* @author Andrea Montemaggio
*/
@Tag("integration")
@Disabled
@Testcontainers(disabledWithoutDocker = true)
public class ElasticsearchConsumerApplicationTests {
@Container
static final ElasticsearchContainer elasticsearch = new ElasticsearchContainer().withStartupAttempts(5)
.withStartupTimeout(Duration.ofMinutes(10));
static final ElasticsearchContainer elasticsearch = new ElasticsearchContainer(
DockerImageName.parse("docker.elastic.co/elasticsearch/elasticsearch")
.withTag("7.17.7")
).withStartupAttempts(5).withStartupTimeout(Duration.ofMinutes(10));
private final ApplicationContextRunner contextRunner = new ApplicationContextRunner()
.withUserConfiguration(ElasticsearchConsumerTestApplication.class, SpelExpressionConverterConfiguration.class);
.withUserConfiguration(ElasticsearchConsumerTestApplication.class, SpelExpressionConverterConfiguration.class);
@Test
public void testBasicJsonString() {
this.contextRunner
.withPropertyValues("elasticsearch.consumer.index=foo", "elasticsearch.consumer.id=1",
"spring.elasticsearch.rest.uris=http://" + elasticsearch.getHttpHostAddress())
.run(context -> {
Consumer<Message<?>> elasticsearchConsumer = context.getBean("elasticsearchConsumer", Consumer.class);
.withPropertyValues("elasticsearch.consumer.index=foo", "elasticsearch.consumer.id=1",
"spring.elasticsearch.rest.uris=http://" + elasticsearch.getHttpHostAddress())
.run(context -> {
final Consumer<Message<?>> elasticsearchConsumer = context.getBean("elasticsearchConsumer", Consumer.class);
String jsonObject = "{\"age\":10,\"dateOfBirth\":1471466076564,"
+ "\"fullName\":\"John Doe\"}";
final Message<String> message = MessageBuilder.withPayload(jsonObject).build();
final String jsonObject = "{\"age\":10,\"dateOfBirth\":1471466076564,"
+ "\"fullName\":\"John Doe\"}";
final Message<String> message = MessageBuilder.withPayload(jsonObject).build();
elasticsearchConsumer.accept(message);
elasticsearchConsumer.accept(message);
final ElasticsearchClient elasticsearchClient = context.getBean(ElasticsearchClient.class);
final GetRequest getRequest = new GetRequest.Builder().index("foo").id("1").build();
final GetResponse<JsonData> response = elasticsearchClient.get(getRequest, JsonData.class);
RestHighLevelClient restHighLevelClient = context.getBean(RestHighLevelClient.class);
GetRequest getRequest = new GetRequest("foo").id("1");
final GetResponse response = restHighLevelClient.get(getRequest, RequestOptions.DEFAULT);
assertThat(response.isExists()).isTrue();
assertThat(response.getSourceAsString()).isEqualTo(jsonObject);
});
assertThat(response.found()).isTrue();
assertThat(response.source()).isNotNull();
assertThat(response.source().toJson()).isEqualTo(JsonData.fromJson(jsonObject).toJson());
});
}
@Test
public void testIdPassedAsMessageHeader() {
this.contextRunner
.withPropertyValues("elasticsearch.consumer.index=foo",
"spring.elasticsearch.rest.uris=http://" + elasticsearch.getHttpHostAddress())
.run(context -> {
Consumer<Message<?>> elasticsearchConsumer = context.getBean("elasticsearchConsumer", Consumer.class);
.withPropertyValues("elasticsearch.consumer.index=foo",
"spring.elasticsearch.rest.uris=http://" + elasticsearch.getHttpHostAddress())
.run(context -> {
final Consumer<Message<?>> elasticsearchConsumer = context.getBean("elasticsearchConsumer", Consumer.class);
String jsonObject = "{\"age\":10,\"dateOfBirth\":1471466076564,"
+ "\"fullName\":\"John Doe\"}";
final Message<String> message = MessageBuilder.withPayload(jsonObject)
.setHeader(ElasticsearchConsumerConfiguration.INDEX_ID_HEADER, "2").build();
final String jsonObject = "{\"age\":10,\"dateOfBirth\":1471466076564,"
+ "\"fullName\":\"John Doe\"}";
final Message<String> message = MessageBuilder.withPayload(jsonObject)
.setHeader(ElasticsearchConsumerConfiguration.INDEX_ID_HEADER, "2").build();
elasticsearchConsumer.accept(message);
elasticsearchConsumer.accept(message);
RestHighLevelClient restHighLevelClient = context.getBean(RestHighLevelClient.class);
GetRequest getRequest = new GetRequest("foo").id("2");
final GetResponse response = restHighLevelClient.get(getRequest, RequestOptions.DEFAULT);
assertThat(response.isExists()).isTrue();
assertThat(response.getSourceAsString()).isEqualTo(jsonObject);
assertThat(response.getId()).isEqualTo("2");
});
final ElasticsearchClient elasticsearchClient = context.getBean(ElasticsearchClient.class);
final GetRequest getRequest = new GetRequest.Builder().index("foo").id("2").build();
final GetResponse<JsonData> response = elasticsearchClient.get(getRequest, JsonData.class);
assertThat(response.found()).isTrue();
assertThat(response.source()).isNotNull();
assertThat(response.source().toJson()).isEqualTo(JsonData.fromJson(jsonObject).toJson());
assertThat(response.id()).isEqualTo("2");
});
}
@Test
public void testJsonAsMap() {
this.contextRunner
.withPropertyValues("elasticsearch.consumer.index=foo", "elasticsearch.consumer.id=3",
"spring.elasticsearch.rest.uris=http://" + elasticsearch.getHttpHostAddress())
.run(context -> {
Consumer<Message<?>> elasticsearchConsumer = context.getBean("elasticsearchConsumer", Consumer.class);
.withPropertyValues("elasticsearch.consumer.index=foo", "elasticsearch.consumer.id=3",
"spring.elasticsearch.rest.uris=http://" + elasticsearch.getHttpHostAddress())
.run(context -> {
final Consumer<Message<?>> elasticsearchConsumer = context.getBean("elasticsearchConsumer", Consumer.class);
Map<String, Object> jsonMap = new HashMap<>();
jsonMap.put("age", 10);
jsonMap.put("dateOfBirth", 1471466076564L);
jsonMap.put("fullName", "John Doe");
final Message<Map<String, Object>> message = MessageBuilder.withPayload(jsonMap).build();
final Map<String, Object> jsonMap = new HashMap<>();
jsonMap.put("age", 10);
jsonMap.put("dateOfBirth", 1471466076564L);
jsonMap.put("fullName", "John Doe");
final Message<Map<String, Object>> message = MessageBuilder.withPayload(jsonMap).build();
elasticsearchConsumer.accept(message);
elasticsearchConsumer.accept(message);
final ElasticsearchClient elasticsearchClient = context.getBean(ElasticsearchClient.class);
final GetRequest getRequest = new GetRequest.Builder().index("foo").id("3").build();
final GetResponse<HashMap> response = elasticsearchClient.get(getRequest, HashMap.class);
RestHighLevelClient restHighLevelClient = context.getBean(RestHighLevelClient.class);
GetRequest getRequest = new GetRequest("foo").id("3");
assertThat(response.found()).isTrue();
HashMap map = response.source();
final GetResponse response = restHighLevelClient.get(getRequest, RequestOptions.DEFAULT);
assertThat(response.isExists()).isTrue();
assertThat(response.getSource()).containsAllEntriesOf(jsonMap);
assertThat(response.getId()).isEqualTo("3");
});
}
@Test
public void testXContentBuilder() {
this.contextRunner
.withPropertyValues("elasticsearch.consumer.index=foo", "elasticsearch.consumer.id=4",
"spring.elasticsearch.rest.uris=http://" + elasticsearch.getHttpHostAddress())
.run(context -> {
Consumer<Message<?>> elasticsearchConsumer = context.getBean("elasticsearchConsumer", Consumer.class);
XContentBuilder builder = XContentFactory.jsonBuilder();
builder.startObject();
builder.field("user", "kimchy");
builder.timeField("postDate", 1471466076564L);
builder.field("message", "trying out Elasticsearch");
builder.endObject();
final Message<XContentBuilder> message = MessageBuilder.withPayload(builder).build();
elasticsearchConsumer.accept(message);
RestHighLevelClient restHighLevelClient = context.getBean(RestHighLevelClient.class);
GetRequest getRequest = new GetRequest("foo").id("4");
final GetResponse response = restHighLevelClient.get(getRequest, RequestOptions.DEFAULT);
assertThat(response.isExists()).isTrue();
assertThat(response.getSourceAsString()).isEqualTo(Strings.toString(builder));
jsonMap.entrySet().forEach(entry -> {
Object value = map.get(entry.getKey());
assertThat(value).isNotNull();
assertThat(value).isEqualTo(entry.getValue());
});
assertThat(response.id()).isEqualTo("3");
});
}
@Test
public void testAsyncIndexing() {
this.contextRunner
.withPropertyValues("elasticsearch.consumer.index=foo", "elasticsearch.consumer.async=true",
"elasticsearch.consumer.id=5",
"spring.elasticsearch.rest.uris=http://" + elasticsearch.getHttpHostAddress())
.run(context -> {
Consumer<Message<?>> elasticsearchConsumer = context.getBean("elasticsearchConsumer", Consumer.class);
.withPropertyValues("elasticsearch.consumer.index=foo", "elasticsearch.consumer.async=true",
"elasticsearch.consumer.id=5",
"spring.elasticsearch.rest.uris=http://" + elasticsearch.getHttpHostAddress())
.run(context -> {
final Consumer<Message<?>> elasticsearchConsumer = context.getBean("elasticsearchConsumer", Consumer.class);
String jsonObject = "{\"age\":10,\"dateOfBirth\":1471466076564,"
+ "\"fullName\":\"John Doe\"}";
final Message<String> message = MessageBuilder.withPayload(jsonObject).build();
final String jsonObject = "{\"age\":10,\"dateOfBirth\":1471466076564,"
+ "\"fullName\":\"John Doe\"}";
final Message<String> message = MessageBuilder.withPayload(jsonObject).build();
elasticsearchConsumer.accept(message);
elasticsearchConsumer.accept(message);
RestHighLevelClient restHighLevelClient = context.getBean(RestHighLevelClient.class);
GetRequest getRequest = new GetRequest("foo").id("5");
final ElasticsearchClient elasticsearchClient = context.getBean(ElasticsearchClient.class);
final GetRequest getRequest = new GetRequest.Builder().index("foo").id("5").build();
Awaitility.given()
.ignoreException(ElasticsearchStatusException.class)
.await()
.until(() -> restHighLevelClient.get(getRequest, RequestOptions.DEFAULT).isExists());
});
Awaitility.given()
.ignoreException(ElasticsearchException.class)
.await()
.until(() -> elasticsearchClient.get(getRequest, JsonData.class).found());
});
}
@Test
public void testBulkIndexingWithIdFromHeader() {
this.contextRunner
.withPropertyValues("elasticsearch.consumer.index=foo_" + UUID.randomUUID(), "elasticsearch.consumer.batch-size=10",
"spring.elasticsearch.rest.uris=http://" + elasticsearch.getHttpHostAddress())
.run(context -> {
Consumer<Message<?>> elasticsearchConsumer = context.getBean("elasticsearchConsumer", Consumer.class);
ElasticsearchConsumerProperties properties = context.getBean(ElasticsearchConsumerProperties.class);
RestHighLevelClient restHighLevelClient = context.getBean(RestHighLevelClient.class);
.withPropertyValues("elasticsearch.consumer.index=foo_" + UUID.randomUUID(), "elasticsearch.consumer.batch-size=10",
"spring.elasticsearch.rest.uris=http://" + elasticsearch.getHttpHostAddress())
.run(context -> {
final Consumer<Message<?>> elasticsearchConsumer = context.getBean("elasticsearchConsumer", Consumer.class);
final ElasticsearchConsumerProperties properties = context.getBean(ElasticsearchConsumerProperties.class);
final ElasticsearchClient elasticsearchClient = context.getBean(ElasticsearchClient.class);
for (int i = 0; i < properties.getBatchSize(); i++) {
final GetRequest getRequest = new GetRequest(properties.getIndex()).id(Integer.toString(i));
assertThatExceptionOfType(ElasticsearchStatusException.class)
.isThrownBy(() -> restHighLevelClient.get(getRequest, RequestOptions.DEFAULT))
.withFailMessage("Expected index not found exception for message %d")
.withMessageContaining("index_not_found_exception");
final Message<String> message = MessageBuilder
.withPayload("{\"seq\":" + i + ",\"age\":10,\"dateOfBirth\":1471466076564,"
+ "\"fullName\":\"John Doe\"}")
.setHeader(ElasticsearchConsumerConfiguration.INDEX_ID_HEADER, Integer.toString(i))
.build();
for (int i = 0; i < properties.getBatchSize(); i++) {
final GetRequest getRequest = new GetRequest.Builder().index(properties.getIndex()).id(Integer.toString(i)).build();
assertThatExceptionOfType(ElasticsearchException.class)
.isThrownBy(() -> elasticsearchClient.get(getRequest, JsonData.class))
.withFailMessage("Expected index not found exception for message %d")
.withMessageContaining("index_not_found_exception");
elasticsearchConsumer.accept(message);
}
final Message<String> message = MessageBuilder
.withPayload("{\"seq\":" + i + ",\"age\":10,\"dateOfBirth\":1471466076564,"
+ "\"fullName\":\"John Doe\"}")
.setHeader(ElasticsearchConsumerConfiguration.INDEX_ID_HEADER, Integer.toString(i))
.build();
for (int i = 0; i < properties.getBatchSize(); i++) {
GetRequest getRequest = new GetRequest(properties.getIndex()).id(Integer.toString(i));
GetResponse response = restHighLevelClient.get(getRequest, RequestOptions.DEFAULT);
elasticsearchConsumer.accept(message);
}
assertThat(response.isExists())
.withFailMessage("Document with id=%d cannot be found.", i)
.isTrue();
assertThat(response.getSource().get("seq")).isEqualTo(i);
}
});
for (int i = 0; i < properties.getBatchSize(); i++) {
final GetRequest getRequest = new GetRequest.Builder().index(properties.getIndex()).id(Integer.toString(i)).build();
final GetResponse<JsonData> response = elasticsearchClient.get(getRequest, JsonData.class);
assertThat(response.found())
.withFailMessage("Document with id=%d cannot be found.", i)
.isTrue();
assertThat(response.source().toJson().asJsonObject().get("seq").toString()).isEqualTo(Integer.toString(i));
}
});
}
@Test
public void testBulkIndexingItemFailure() {
this.contextRunner
.withPropertyValues("elasticsearch.consumer.index=foo_" + UUID.randomUUID(), "elasticsearch.consumer.batch-size=10",
"spring.elasticsearch.rest.uris=http://" + elasticsearch.getHttpHostAddress())
.run(context -> {
Consumer<Message<?>> elasticsearchConsumer = context.getBean("elasticsearchConsumer", Consumer.class);
ElasticsearchConsumerProperties properties = context.getBean(ElasticsearchConsumerProperties.class);
RestHighLevelClient restHighLevelClient = context.getBean(RestHighLevelClient.class);
.withPropertyValues("elasticsearch.consumer.index=foo_" + UUID.randomUUID(), "elasticsearch.consumer.batch-size=10",
"spring.elasticsearch.rest.uris=http://" + elasticsearch.getHttpHostAddress())
.run(context -> {
final Consumer<Message<?>> elasticsearchConsumer = context.getBean("elasticsearchConsumer", Consumer.class);
final ElasticsearchConsumerProperties properties = context.getBean(ElasticsearchConsumerProperties.class);
final ElasticsearchClient elasticsearchClient = context.getBean(ElasticsearchClient.class);
for (int i = 0; i < properties.getBatchSize(); i++) {
final GetRequest getRequest = new GetRequest(properties.getIndex()).id(Integer.toString(i));
assertThatExceptionOfType(ElasticsearchStatusException.class)
.isThrownBy(() -> restHighLevelClient.get(getRequest, RequestOptions.DEFAULT))
.withFailMessage("Expected index not found exception for message %d")
.withMessageContaining("index_not_found_exception");
MessageBuilder<String> builder = MessageBuilder
.withPayload("{\"seq\":" + i + ",\"age\":10,\"dateOfBirth\":1471466076564,"
+ "\"fullName\":\"John Doe\"}")
.setHeader(ElasticsearchConsumerConfiguration.INDEX_ID_HEADER, Integer.toString(i));
for (int i = 0; i < properties.getBatchSize(); i++) {
final GetRequest getRequest = new GetRequest.Builder().index(properties.getIndex()).id(Integer.toString(i)).build();
assertThatExceptionOfType(ElasticsearchException.class)
.isThrownBy(() -> elasticsearchClient.get(getRequest, JsonData.class))
.withFailMessage("Expected index not found exception for message %d")
.withMessageContaining("index_not_found_exception");
if (i == 0) {
// set an invalid index name to make the first request fail
builder.setHeader(ElasticsearchConsumerConfiguration.INDEX_NAME_HEADER, "_" + properties.getIndex());
}
MessageBuilder<String> builder = MessageBuilder
.withPayload("{\"seq\":" + i + ",\"age\":10,\"dateOfBirth\":1471466076564,"
+ "\"fullName\":\"John Doe\"}")
.setHeader(ElasticsearchConsumerConfiguration.INDEX_ID_HEADER, Integer.toString(i));
final Message<String> message = builder.build();
if (i < properties.getBatchSize() - 1) {
elasticsearchConsumer.accept(message);
}
else {
// last invocation
assertThatIllegalStateException()
.isThrownBy(() -> elasticsearchConsumer.accept(message))
.withMessageContaining("Bulk indexing operation completed with failures");
}
if (i == 0) {
// set an invalid index name to make the first request fail
builder.setHeader(ElasticsearchConsumerConfiguration.INDEX_NAME_HEADER, "_" + properties.getIndex());
}
});
final Message<String> message = builder.build();
if (i < properties.getBatchSize() - 1) {
elasticsearchConsumer.accept(message);
}
else {
// last invocation
assertThatIllegalStateException()
.isThrownBy(() -> elasticsearchConsumer.accept(message))
.withMessageContaining("Bulk indexing operation completed with failures");
}
}
});
}
@Test
public void testIndexFromMessageHeader() {
this.contextRunner
.withPropertyValues("elasticsearch.consumer.index=foo",
"spring.elasticsearch.rest.uris=http://" + elasticsearch.getHttpHostAddress())
.run(context -> {
Consumer<Message<?>> elasticsearchConsumer = context.getBean("elasticsearchConsumer", Consumer.class);
ElasticsearchConsumerProperties properties = context.getBean(ElasticsearchConsumerProperties.class);
.withPropertyValues("elasticsearch.consumer.index=foo",
"spring.elasticsearch.rest.uris=http://" + elasticsearch.getHttpHostAddress())
.run(context -> {
final Consumer<Message<?>> elasticsearchConsumer = context.getBean("elasticsearchConsumer", Consumer.class);
final ElasticsearchConsumerProperties properties = context.getBean(ElasticsearchConsumerProperties.class);
final String dynamicIndex = properties.getIndex() + "-2";
final String dynamicIndex = properties.getIndex() + "-2";
String jsonObject = "{\"age\":10,\"dateOfBirth\":1471466076564,"
+ "\"fullName\":\"John Doe\"}";
final Message<String> message = MessageBuilder.withPayload(jsonObject)
.setHeader(ElasticsearchConsumerConfiguration.INDEX_ID_HEADER, "2")
.setHeader(ElasticsearchConsumerConfiguration.INDEX_NAME_HEADER, dynamicIndex)
.build();
final String jsonObject = "{\"age\":10,\"dateOfBirth\":1471466076564,"
+ "\"fullName\":\"John Doe\"}";
final Message<String> message = MessageBuilder.withPayload(jsonObject)
.setHeader(ElasticsearchConsumerConfiguration.INDEX_ID_HEADER, "2")
.setHeader(ElasticsearchConsumerConfiguration.INDEX_NAME_HEADER, dynamicIndex)
.build();
elasticsearchConsumer.accept(message);
elasticsearchConsumer.accept(message);
final ElasticsearchClient elasticsearchClient = context.getBean(ElasticsearchClient.class);
RestHighLevelClient restHighLevelClient = context.getBean(RestHighLevelClient.class);
GetRequest getRequest = new GetRequest(dynamicIndex).id("2");
final GetResponse response = restHighLevelClient.get(getRequest, RequestOptions.DEFAULT);
assertThat(response.isExists()).isTrue();
assertThat(response.getSourceAsString()).isEqualTo(jsonObject);
assertThat(response.getId()).isEqualTo("2");
});
GetRequest getRequest = new GetRequest.Builder().index(dynamicIndex).id("2").build();
final GetResponse<JsonData> response = elasticsearchClient.get(getRequest, JsonData.class);
assertThat(response.found()).isTrue();
assertThat(response.source()).isNotNull();
assertThat(response.source().toJson()).isEqualTo(JsonData.fromJson(jsonObject).toJson());
assertThat(response.id()).isEqualTo("2");
});
}
@SpringBootApplication
static class ElasticsearchConsumerTestApplication {
}
@Configuration
static class Config extends ElasticsearchConfiguration {
@NonNull
@Override
public ClientConfiguration clientConfiguration() {
return ClientConfiguration.builder()
.connectedTo(elasticsearch.getHttpHostAddress())
.build();
}
}
}