GH-430: Change Filter function to non-reactive

Fixes https://github.com/spring-cloud/stream-applications/issues/430

Fix greenmail exclusion.
Added -X for github.debug
This commit is contained in:
Corneil du Plessis
2023-02-14 18:23:31 +02:00
committed by GitHub
parent d29546c9bd
commit 8e4c20dc77
4 changed files with 30 additions and 23 deletions

View File

@@ -32,9 +32,11 @@
<exclusions>
<exclusion>
<groupId>com.sun.mail</groupId>
<artifactId>jakarta.mail</artifactId>
</exclusion>
<exclusion>
<groupId>jakarta.activation</groupId>
<artifactId>jakarta.activation-api</artifactId>
</exclusion>
</exclusions>
</dependency>

View File

@@ -18,8 +18,6 @@ package org.springframework.cloud.fn.filter;
import java.util.function.Function;
import reactor.core.publisher.Flux;
import org.springframework.boot.autoconfigure.AutoConfiguration;
import org.springframework.boot.context.properties.EnableConfigurationProperties;
import org.springframework.context.annotation.Bean;
@@ -35,17 +33,22 @@ import org.springframework.messaging.Message;
public class FilterFunctionConfiguration {
@Bean
public Function<Flux<Message<?>>, Flux<Message<?>>> filterFunction(
ExpressionEvaluatingTransformer filterExpressionEvaluatingTransformer) {
public Function<Message<?>, Message<?>> filterFunction(
ExpressionEvaluatingTransformer filterExpressionEvaluatingTransformer) {
return flux ->
flux.filter((message) ->
(Boolean) filterExpressionEvaluatingTransformer.transform(message).getPayload());
return message -> {
if ((Boolean) filterExpressionEvaluatingTransformer.transform(message).getPayload()) {
return message;
}
else {
return null;
}
};
}
@Bean
public ExpressionEvaluatingTransformer filterExpressionEvaluatingTransformer(
FilterFunctionProperties filterFunctionProperties) {
FilterFunctionProperties filterFunctionProperties) {
return new ExpressionEvaluatingTransformer(filterFunctionProperties.getExpression());
}

View File

@@ -16,11 +16,15 @@
package org.springframework.cloud.fn.filter;
import java.util.List;
import java.util.function.Function;
import java.util.stream.Collectors;
import java.util.stream.Stream;
import org.junit.jupiter.api.Test;
import reactor.core.publisher.Flux;
import reactor.test.StepVerifier;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.beans.factory.annotation.Qualifier;
@@ -30,6 +34,8 @@ import org.springframework.messaging.Message;
import org.springframework.messaging.support.GenericMessage;
import org.springframework.test.annotation.DirtiesContext;
import static org.assertj.core.api.Assertions.assertThat;
/**
* @author Artem Bilan
* @author David Turanski
@@ -40,21 +46,17 @@ public class FilterFunctionApplicationTests {
@Autowired
@Qualifier("filterFunction")
Function<Flux<Message<?>>, Flux<Message<?>>> filter;
Function<Message<?>, Message<?>> filter;
@Test
public void testFilter() {
Flux<Message<?>> messageFlux =
Flux.just("hello", "hello world")
.map(GenericMessage::new);
Flux<Message<?>> result = this.filter.apply(messageFlux);
result
.map(Message::getPayload)
.cast(String.class)
.as(StepVerifier::create)
.expectNext("hello world")
.expectComplete()
.verify();
Stream<Message<?>> messages = List.of("hello", "hello world")
.stream()
.map(GenericMessage::new);
List<Message<?>> result = messages.filter(message -> this.filter.apply(message) != null).collect(Collectors.toList());
assertThat(result.size()).isEqualTo(1);
assertThat(result.get(0).getPayload()).isNotNull();
assertThat(result.get(0).getPayload()).isEqualTo("hello world");
}
@SpringBootApplication

View File

@@ -2,7 +2,7 @@
This module provides an HTTP request function that can be reused and composed in other applications.
The `Function` uses the reactive `WebClient` from `Spring WebFlux` and is implemented as a `java.util.function.Function`.
This function gives you a reactive stream of `ResponseEntity` given a stream of request messages as the function a signature of `Function<Flux<Message<?>,Flux<ResponseEntity>>`.
This function gives you a reactive stream of `ResponseEntity` given a stream of request messages as the function a signature of `Function<Message<?>,ResponseEntity>`.
Users have to subscribe to the returned `Flux` to receive the data.
## Beans for injection