From a91aa2cd3fa7268f06900688b5f1abfba0bcdd8e Mon Sep 17 00:00:00 2001 From: Artem Bilan Date: Fri, 28 Aug 2020 15:55:08 -0400 Subject: [PATCH] GH-110: Fix filter function for emitting `null` (#111) * GH-110: Fix filter function for emitting `null` Fixes https://github.com/spring-cloud/stream-applications/issues/110 Turns out when functions are composed in the Spring Cloud Stream environment, they are called via reactive wrappers which don't allow to emit `null` from the `map()` operator. * Make `filterFunction` fully reactive to reply on the `Flux.filter()` operator * Add `config-common` for auto-conversion string configuration options into `Expression` instances * Add `spring-boot-starter-json` dependency since it is required by the `SpelExpressionConverterConfiguration` * Add `reactor-test` dependency to test the final solution * Remove redundant dependencies from the `filter-processor` * * Add `proxyBeanMethods = false` into `FilterFunctionConfiguration` * Fix default expression to `true` instead of `payload`, which does not fit to filter logic * Fix JavaDoc for `expression` property * Remove redundant `application.properties` from the `filter-function` --- .../processor/filter-processor/README.adoc | 2 +- .../processor/filter-processor/pom.xml | 17 ----------- functions/function/filter-function/pom.xml | 14 ++++++++++ .../filter/FilterFunctionConfiguration.java | 17 ++++++----- .../fn/filter/FilterFunctionProperties.java | 12 ++++---- .../src/main/resources/application.properties | 1 - .../FilterFunctionApplicationTests.java | 28 ++++++++++++++----- 7 files changed, 49 insertions(+), 42 deletions(-) delete mode 100644 functions/function/filter-function/src/main/resources/application.properties diff --git a/applications/processor/filter-processor/README.adoc b/applications/processor/filter-processor/README.adoc index db61e062..f29bee77 100644 --- a/applications/processor/filter-processor/README.adoc +++ b/applications/processor/filter-processor/README.adoc @@ -16,7 +16,7 @@ If the incoming type is `byte[]` and the content type is set to `text/plain` or == Options //tag::configuration-properties[] -$$filter.function.expression$$:: $$A SpEL expression to apply.$$ *($$String$$, default: `$$$$`)* +$$filter.function.expression$$:: $$Boolean SpEL expression to apply against request message to filter.$$ *($$String$$, default: `$$true$$`)* //end::configuration-properties[] //end::ref-doc[] diff --git a/applications/processor/filter-processor/pom.xml b/applications/processor/filter-processor/pom.xml index 228444a4..ef6879b4 100644 --- a/applications/processor/filter-processor/pom.xml +++ b/applications/processor/filter-processor/pom.xml @@ -20,23 +20,6 @@ org.springframework.cloud.fn filter-function - - - org.springframework.boot - spring-boot-starter-test - test - - - org.junit.vintage - junit-vintage-engine - - - - - org.springframework.boot - spring-boot-starter-json - test - diff --git a/functions/function/filter-function/pom.xml b/functions/function/filter-function/pom.xml index 901a61af..84f586a0 100644 --- a/functions/function/filter-function/pom.xml +++ b/functions/function/filter-function/pom.xml @@ -15,6 +15,11 @@ + + org.springframework.cloud.fn + config-common + ${project.version} + org.springframework.cloud.fn payload-converter-function @@ -24,6 +29,10 @@ org.springframework.boot spring-boot-starter-integration + + org.springframework.boot + spring-boot-starter-json + org.springframework.boot spring-boot-configuration-processor @@ -40,6 +49,11 @@ + + io.projectreactor + reactor-test + test + diff --git a/functions/function/filter-function/src/main/java/org/springframework/cloud/fn/filter/FilterFunctionConfiguration.java b/functions/function/filter-function/src/main/java/org/springframework/cloud/fn/filter/FilterFunctionConfiguration.java index a935ef9a..a9aba2be 100644 --- a/functions/function/filter-function/src/main/java/org/springframework/cloud/fn/filter/FilterFunctionConfiguration.java +++ b/functions/function/filter-function/src/main/java/org/springframework/cloud/fn/filter/FilterFunctionConfiguration.java @@ -16,39 +16,38 @@ package org.springframework.cloud.fn.filter; -import java.util.Optional; import java.util.function.Function; import org.springframework.boot.context.properties.EnableConfigurationProperties; import org.springframework.context.annotation.Bean; import org.springframework.context.annotation.Configuration; -import org.springframework.expression.spel.standard.SpelExpressionParser; import org.springframework.integration.transformer.ExpressionEvaluatingTransformer; import org.springframework.messaging.Message; +import reactor.core.publisher.Flux; + /** * @author Artem Bilan * @author David Turanski */ -@Configuration +@Configuration(proxyBeanMethods = false) @EnableConfigurationProperties(FilterFunctionProperties.class) public class FilterFunctionConfiguration { @Bean - public Function, Message> filterFunction( + public Function>, Flux>> filterFunction( ExpressionEvaluatingTransformer filterExpressionEvaluatingTransformer) { - return message -> Optional.of(message) - .filter(m -> (Boolean) filterExpressionEvaluatingTransformer.transform(m).getPayload()) - .orElse(null); + return flux -> + flux.filter((message) -> + (Boolean) filterExpressionEvaluatingTransformer.transform(message).getPayload()); } @Bean public ExpressionEvaluatingTransformer filterExpressionEvaluatingTransformer( FilterFunctionProperties filterFunctionProperties) { - return new ExpressionEvaluatingTransformer(new SpelExpressionParser() - .parseExpression(filterFunctionProperties.getExpression())); + return new ExpressionEvaluatingTransformer(filterFunctionProperties.getExpression()); } } diff --git a/functions/function/filter-function/src/main/java/org/springframework/cloud/fn/filter/FilterFunctionProperties.java b/functions/function/filter-function/src/main/java/org/springframework/cloud/fn/filter/FilterFunctionProperties.java index bde33f71..ab5798d1 100644 --- a/functions/function/filter-function/src/main/java/org/springframework/cloud/fn/filter/FilterFunctionProperties.java +++ b/functions/function/filter-function/src/main/java/org/springframework/cloud/fn/filter/FilterFunctionProperties.java @@ -18,7 +18,7 @@ package org.springframework.cloud.fn.filter; import org.springframework.boot.context.properties.ConfigurationProperties; import org.springframework.expression.Expression; -import org.springframework.expression.spel.standard.SpelExpressionParser; +import org.springframework.integration.expression.ValueExpression; /** * Configuration properties for the SpEL function. @@ -29,18 +29,16 @@ import org.springframework.expression.spel.standard.SpelExpressionParser; @ConfigurationProperties("filter.function") public class FilterFunctionProperties { - private static final Expression DEFAULT_EXPRESSION = new SpelExpressionParser().parseExpression("payload"); - /** - * A SpEL expression to apply. + * Boolean SpEL expression to apply against request message to filter. */ - private String expression = DEFAULT_EXPRESSION.getExpressionString(); + private Expression expression = new ValueExpression<>(true); - public String getExpression() { + public Expression getExpression() { return this.expression; } - public void setExpression(String expression) { + public void setExpression(Expression expression) { this.expression = expression; } diff --git a/functions/function/filter-function/src/main/resources/application.properties b/functions/function/filter-function/src/main/resources/application.properties deleted file mode 100644 index 86e445ae..00000000 --- a/functions/function/filter-function/src/main/resources/application.properties +++ /dev/null @@ -1 +0,0 @@ -spel.function.expression=true diff --git a/functions/function/filter-function/src/test/java/org/springframework/cloud/fn/filter/FilterFunctionApplicationTests.java b/functions/function/filter-function/src/test/java/org/springframework/cloud/fn/filter/FilterFunctionApplicationTests.java index ed26dfef..0c8fc873 100644 --- a/functions/function/filter-function/src/test/java/org/springframework/cloud/fn/filter/FilterFunctionApplicationTests.java +++ b/functions/function/filter-function/src/test/java/org/springframework/cloud/fn/filter/FilterFunctionApplicationTests.java @@ -30,25 +30,39 @@ import org.springframework.test.annotation.DirtiesContext; import static org.assertj.core.api.Assertions.assertThat; +import reactor.core.publisher.Flux; +import reactor.test.StepVerifier; + +/** + * @author Artem Bilan + * @author David Turanski + */ @SpringBootTest(properties = "filter.function.expression=payload.length() > 5") @DirtiesContext public class FilterFunctionApplicationTests { @Autowired @Qualifier("filterFunction") - Function, Message> filter; + Function>, Flux>> filter; @Test public void testFilter() { - Message filtered = this.filter.apply(new GenericMessage<>("hello")); - assertThat(filtered).isNull(); - filtered = this.filter.apply(new GenericMessage<>("hello world")); - assertThat(filtered).isNotNull() - .extracting(Message::getPayload) - .isEqualTo("hello world"); + Flux> messageFlux = + Flux.just("hello", "hello world") + .map(GenericMessage::new); + Flux> result = this.filter.apply(messageFlux); + result + .map(Message::getPayload) + .cast(String.class) + .as(StepVerifier::create) + .expectNext("hello world") + .expectComplete() + .verify(); } @SpringBootApplication static class FilterFunctionTestApplication { + } + }