diff --git a/applications/sink/elasticsearch-sink/src/test/java/org/springframework/cloud/stream/app/sink/elasticsearch/ElasticsearchSinkTests.java b/applications/sink/elasticsearch-sink/src/test/java/org/springframework/cloud/stream/app/sink/elasticsearch/ElasticsearchSinkTests.java index d5d69497..d3bf671b 100644 --- a/applications/sink/elasticsearch-sink/src/test/java/org/springframework/cloud/stream/app/sink/elasticsearch/ElasticsearchSinkTests.java +++ b/applications/sink/elasticsearch-sink/src/test/java/org/springframework/cloud/stream/app/sink/elasticsearch/ElasticsearchSinkTests.java @@ -17,6 +17,7 @@ package org.springframework.cloud.stream.app.sink.elasticsearch; import java.time.Duration; +import java.util.Map; import co.elastic.clients.elasticsearch.ElasticsearchClient; import co.elastic.clients.elasticsearch.core.GetRequest; @@ -60,7 +61,7 @@ public class ElasticsearchSinkTests { @Test - public void tesElasticSearchSink() { + void elasticSearchSinkWithIndexNameProperty() { this.contextRunner .withPropertyValues("spring.cloud.function.definition=elasticsearchConsumer", "elasticsearch.consumer.index=foo", "elasticsearch.consumer.id=1", @@ -82,12 +83,34 @@ public class ElasticsearchSinkTests { }); } + @Test + void elasticSearchSinkWithIndexNameFromHeader() { + this.contextRunner + .withPropertyValues("spring.cloud.function.definition=elasticsearchConsumer", "elasticsearch.consumer.id=1", + "spring.elasticsearch.rest.uris=http://" + elasticsearch.getHttpHostAddress()) + .run(context -> { + + final InputDestination inputDestination = context.getBean(InputDestination.class); + final String jsonObject = "{\"age\":10,\"dateOfBirth\":1471466076564," + + "\"fullName\":\"John Doe\"}"; + + inputDestination.send(new GenericMessage<>(jsonObject, Map.of("INDEX_NAME", "foo"))); + + 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); + assertThat(response.found()).isTrue(); + assertThat(response.source()).isNotNull(); + assertThat(response.source().toJson()).isEqualTo(JsonData.fromJson(jsonObject).toJson()); + }); + } + @SpringBootApplication @Import(ElasticsearchConsumerConfiguration.class) static class ElasticsearchSinkTestApplication { - } - @Configuration + + @Configuration(proxyBeanMethods = false) static class Config extends ElasticsearchConfiguration { @NonNull @Override diff --git a/applications/sink/elasticsearch-sink/src/test/resources/logback-test.xml b/applications/sink/elasticsearch-sink/src/test/resources/logback-test.xml new file mode 100644 index 00000000..d155fd18 --- /dev/null +++ b/applications/sink/elasticsearch-sink/src/test/resources/logback-test.xml @@ -0,0 +1,13 @@ + + + + %d{HH:mm:ss.SSS} [%thread] %-5level %logger - %msg%n + + + + + + + + + diff --git a/functions/consumer/elasticsearch-consumer/src/main/java/org/springframework/cloud/fn/consumer/elasticsearch/ElasticsearchConsumerConfiguration.java b/functions/consumer/elasticsearch-consumer/src/main/java/org/springframework/cloud/fn/consumer/elasticsearch/ElasticsearchConsumerConfiguration.java index b60ba662..636a22a6 100644 --- a/functions/consumer/elasticsearch-consumer/src/main/java/org/springframework/cloud/fn/consumer/elasticsearch/ElasticsearchConsumerConfiguration.java +++ b/functions/consumer/elasticsearch-consumer/src/main/java/org/springframework/cloud/fn/consumer/elasticsearch/ElasticsearchConsumerConfiguration.java @@ -47,7 +47,6 @@ import org.springframework.integration.aggregator.MessageCountReleaseStrategy; import org.springframework.integration.config.AggregatorFactoryBean; import org.springframework.integration.dsl.IntegrationFlow; import org.springframework.integration.dsl.IntegrationFlowBuilder; -import org.springframework.integration.dsl.IntegrationFlows; import org.springframework.integration.expression.ValueExpression; import org.springframework.integration.store.MessageGroup; import org.springframework.integration.store.MessageGroupStore; @@ -126,7 +125,7 @@ public class ElasticsearchConsumerConfiguration { ) { final IntegrationFlowBuilder builder = - IntegrationFlows.from(Consumer.class, gateway -> gateway.beanName("elasticsearchConsumer")); + IntegrationFlow.from(MessageConsumer.class, gateway -> gateway.beanName("elasticsearchConsumer")); if (properties.getBatchSize() > 1) { builder.handle(aggregator); } @@ -297,4 +296,9 @@ public class ElasticsearchConsumerConfiguration { return message; } } + + private interface MessageConsumer extends Consumer> { + + } + }