From 1459bfabd94da16b3d783346028f7d3a29682a1c Mon Sep 17 00:00:00 2001 From: Corneil du Plessis Date: Fri, 25 Nov 2022 11:51:53 +0200 Subject: [PATCH] [SB3] Update Elasticsearch Client. Update Boot and Spring versions. Improve Rabbit listener (#416) --- consumer/elasticsearch-consumer/pom.xml | 11 +- .../ElasticsearchConsumerConfiguration.java | 153 ++++---- ...ElasticsearchConsumerApplicationTests.java | 363 +++++++++--------- 3 files changed, 266 insertions(+), 261 deletions(-) diff --git a/consumer/elasticsearch-consumer/pom.xml b/consumer/elasticsearch-consumer/pom.xml index 60859ec1..a38b0c40 100644 --- a/consumer/elasticsearch-consumer/pom.xml +++ b/consumer/elasticsearch-consumer/pom.xml @@ -10,7 +10,7 @@ - 7.15.2 + 8.5.0 elasticsearch-consumer @@ -19,13 +19,8 @@ - org.elasticsearch - elasticsearch - ${elasticsearch.version} - - - org.elasticsearch.client - elasticsearch-rest-high-level-client + co.elastic.clients + elasticsearch-java ${elasticsearch.version} diff --git a/consumer/elasticsearch-consumer/src/main/java/org/springframework/cloud/fn/consumer/elasticsearch/ElasticsearchConsumerConfiguration.java b/consumer/elasticsearch-consumer/src/main/java/org/springframework/cloud/fn/consumer/elasticsearch/ElasticsearchConsumerConfiguration.java index 30db60a9..ea382fc3 100644 --- a/consumer/elasticsearch-consumer/src/main/java/org/springframework/cloud/fn/consumer/elasticsearch/ElasticsearchConsumerConfiguration.java +++ b/consumer/elasticsearch-consumer/src/main/java/org/springframework/cloud/fn/consumer/elasticsearch/ElasticsearchConsumerConfiguration.java @@ -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) message.getPayload(), XContentType.JSON); - } - else if (message.getPayload() instanceof XContentBuilder) { - request.source((XContentBuilder) message.getPayload()); + requestBuilder.document((Map) 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() { - @Override - public void onResponse(BulkResponse bulkResponse) { + ElasticsearchAsyncClient elasticsearchAsyncClient = new ElasticsearchAsyncClient(elasticsearchClient._transport()); + CompletableFuture 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() { - @Override - public void onResponse(IndexResponse indexResponse) { - handleResponse(indexResponse); + ElasticsearchAsyncClient elasticsearchAsyncClient = new ElasticsearchAsyncClient(elasticsearchClient._transport()); + CompletableFuture 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()) ); } diff --git a/consumer/elasticsearch-consumer/src/test/java/org/springframework/cloud/fn/consumer/elasticsearch/ElasticsearchConsumerApplicationTests.java b/consumer/elasticsearch-consumer/src/test/java/org/springframework/cloud/fn/consumer/elasticsearch/ElasticsearchConsumerApplicationTests.java index ed28fd24..d6c61568 100644 --- a/consumer/elasticsearch-consumer/src/test/java/org/springframework/cloud/fn/consumer/elasticsearch/ElasticsearchConsumerApplicationTests.java +++ b/consumer/elasticsearch-consumer/src/test/java/org/springframework/cloud/fn/consumer/elasticsearch/ElasticsearchConsumerApplicationTests.java @@ -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> 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> elasticsearchConsumer = context.getBean("elasticsearchConsumer", Consumer.class); - String jsonObject = "{\"age\":10,\"dateOfBirth\":1471466076564," - + "\"fullName\":\"John Doe\"}"; - final Message message = MessageBuilder.withPayload(jsonObject).build(); + final String jsonObject = "{\"age\":10,\"dateOfBirth\":1471466076564," + + "\"fullName\":\"John Doe\"}"; + final Message 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 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> elasticsearchConsumer = context.getBean("elasticsearchConsumer", Consumer.class); + .withPropertyValues("elasticsearch.consumer.index=foo", + "spring.elasticsearch.rest.uris=http://" + elasticsearch.getHttpHostAddress()) + .run(context -> { + final Consumer> elasticsearchConsumer = context.getBean("elasticsearchConsumer", Consumer.class); - String jsonObject = "{\"age\":10,\"dateOfBirth\":1471466076564," - + "\"fullName\":\"John Doe\"}"; - final Message message = MessageBuilder.withPayload(jsonObject) - .setHeader(ElasticsearchConsumerConfiguration.INDEX_ID_HEADER, "2").build(); + final String jsonObject = "{\"age\":10,\"dateOfBirth\":1471466076564," + + "\"fullName\":\"John Doe\"}"; + final Message 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 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> 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> elasticsearchConsumer = context.getBean("elasticsearchConsumer", Consumer.class); - Map jsonMap = new HashMap<>(); - jsonMap.put("age", 10); - jsonMap.put("dateOfBirth", 1471466076564L); - jsonMap.put("fullName", "John Doe"); - final Message> message = MessageBuilder.withPayload(jsonMap).build(); + final Map jsonMap = new HashMap<>(); + jsonMap.put("age", 10); + jsonMap.put("dateOfBirth", 1471466076564L); + jsonMap.put("fullName", "John Doe"); + final Message> 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 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> 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 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> 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> elasticsearchConsumer = context.getBean("elasticsearchConsumer", Consumer.class); - String jsonObject = "{\"age\":10,\"dateOfBirth\":1471466076564," - + "\"fullName\":\"John Doe\"}"; - final Message message = MessageBuilder.withPayload(jsonObject).build(); + final String jsonObject = "{\"age\":10,\"dateOfBirth\":1471466076564," + + "\"fullName\":\"John Doe\"}"; + final Message 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> 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> 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 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 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 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> 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> 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 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 builder = MessageBuilder + .withPayload("{\"seq\":" + i + ",\"age\":10,\"dateOfBirth\":1471466076564," + + "\"fullName\":\"John Doe\"}") + .setHeader(ElasticsearchConsumerConfiguration.INDEX_ID_HEADER, Integer.toString(i)); - final Message 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 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> 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> 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 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 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 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(); + } + } }