Fix ElasticsearchSink index name header (#446)

- Use `Consumer<Message<?>>` on IntegrationFlow
  rather than `Consumer`

Fixes #440
This commit is contained in:
Chris Bono
2023-04-11 15:50:59 -05:00
committed by GitHub
parent 69af9564a3
commit d61bb6db2b
3 changed files with 45 additions and 5 deletions

View File

@@ -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<JsonData> 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

View File

@@ -0,0 +1,13 @@
<configuration>
<appender name="STDOUT" class="ch.qos.logback.core.ConsoleAppender">
<encoder>
<pattern>%d{HH:mm:ss.SSS} [%thread] %-5level %logger - %msg%n</pattern>
</encoder>
</appender>
<root level="WARN">
<appender-ref ref="STDOUT"/>
</root>
<logger name="org.testcontainers" level="ERROR"/>
<logger name="com.github.dockerjava" level="ERROR"/>
<logger name="org.springframework.cloud" level="INFO"/>
</configuration>

View File

@@ -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<Message<?>> {
}
}