GH-1748 - Enhance routing function with SpEL and application properties

Resolves #1748
This commit is contained in:
Oleg Zhurakousky
2019-09-10 14:40:32 +02:00
parent aa495ac4ae
commit f3b696f35e
4 changed files with 184 additions and 116 deletions

View File

@@ -635,7 +635,6 @@ Here is the example of a Processor application defined as `java.util.function.Fu
[source,java]
----
@SpringBootApplication
@EnableBinding(Processor.class)
public static class ProcessorFromFunction {
public static void main(String[] args) {
SpringApplication.run(ProcessorFromFunction.class, "--spring.cloud.stream.function.definition=toUpperCase");
@@ -651,7 +650,6 @@ Here is the example of a Sink application defined as `java.util.function.Consume
[source,java]
----
@EnableAutoConfiguration
@EnableBinding(Sink.class)
public static class SinkFromConsumer {
public static void main(String[] args) {
SpringApplication.run(SinkFromConsumer.class, "--spring.cloud.stream.function.definition=sink");
@@ -662,6 +660,57 @@ public static class SinkFromConsumer {
}
}
----
===== Content-based routing with functions
Routing with functions can achieved by relying on `RoutingFunction` available in Spring Cloud Function 3.0. All you need to do is enable it via
`--spring.cloud.stream.function.routing.enabled=true` application property. Once enabled `RoutingFunction` will be bound to input destination
receiving all the messages and route them to other functions based on the provided instruction.
Instruction could be provided with individual messages as well as application properties.
Here are couple of samples:
***Using message headers***
[source,java]
----
@SpringBootApplication
public class SampleApplication {
public static void main(String[] args) {
SpringApplication.run(SampleApplication.class,
"--spring.cloud.stream.function.routing.enabled=true");
}
@Bean
public Consumer<String> even() {
return value -> {
System.out.println("EVEN: " + value);
};
}
@Bean
public Consumer<String> odd() {
return value -> {
System.out.println("ODD: " + value);
};
}
}
----
By default `RoutingFunction` will look for `spring.cloud.function.definition` header and if it is found its value will be treated as routing instruction.
So in the above case the value of such header should be either `odd` or `even` (the name of the function beans) to route request to available functions.
You can also use SpEL for more dynamic scenarios via `spring.cloud.function.routing-expression` header.
For example,
setting `spring.cloud.function.routing-expression` header to value `T(java.lang.System).currentTimeMillis() % 2 == 0 ? 'even' : 'odd'` will end up semi-randomly routing request to either `odd` or `even` functions.
Also, for SpEL, the _root object_ of the evaluation context is `Message` so you can do evaluation on individual headers (or message) as well `....routing-expression=headers['type']`
***Using application properties***
The `spring.cloud.function.routing-expression` and/or `spring.cloud.function.definition`
can be passed as application properties (e.g., `spring.cloud.function.routing-expression=headers['type']`.
Passing instructions via application properties is especially important for reactive functions since given that fact that reactive
function is only invoked once to pass the Publisher, so access to the individual items is limited.
===== Reactive Functions support
Since _Spring Cloud Function_ is build on top of https://projectreactor.io/[Project Reactor] there isn't much you need to do
@@ -672,7 +721,6 @@ For example:
[source,java]
----
@EnableAutoConfiguration
@EnableBinding(Processor.class)
public static class SinkFromConsumer {
public static void main(String[] args) {
SpringApplication.run(SinkFromConsumer.class, "--spring.cloud.stream.function.definition=reactiveUpperCase");

View File

@@ -318,7 +318,7 @@ public class BinderFactoryAutoConfiguration {
name = environment.getProperty("spring.cloud.function.definition");
}
if (!StringUtils.hasText(name) && Boolean.parseBoolean(
environment.getProperty("spring.cloud.function.routing.enabled", "false"))) {
environment.getProperty("spring.cloud.stream.function.routing.enabled", "false"))) {
name = RoutingFunction.FUNCTION_NAME;
}
if (!StringUtils.hasText(name) && catalog.size() >= 1 && catalog.size() <= 2) {

View File

@@ -83,7 +83,7 @@ import org.springframework.util.ReflectionUtils;
@EnableConfigurationProperties(StreamFunctionProperties.class)
@Import(BinderFactoryAutoConfiguration.class)
@AutoConfigureBefore(BindingServiceConfiguration.class)
public class FunctionConfiguration {
class FunctionConfiguration {
@Bean
public InitializingBean functionChannelBindingInitializer(FunctionCatalog functionCatalog, FunctionInspector functionInspector,
@@ -95,44 +95,41 @@ public class FunctionConfiguration {
@Bean
public IntegrationFlow standAloneSupplierFlow(FunctionCatalog functionCatalog, FunctionInspector functionInspector,
StreamFunctionProperties functionProperties, GenericApplicationContext context) {
FunctionInvocationWrapper functionWrapper = functionCatalog.lookup(functionProperties.getDefinition());
IntegrationFlow integrationFlow = null;
if (functionCatalog != null && ObjectUtils.isEmpty(context.getBeanNamesForAnnotation(EnableBinding.class))) {
FunctionInvocationWrapper functionWrapper = functionCatalog.lookup(functionProperties.getDefinition());
if (functionWrapper != null /*&& functionWrapper.getTarget() instanceof Supplier*/) {
AtomicReference<MonoSink<Object>> triggerRef = new AtomicReference<>();
Publisher<Object> beginPublishingTrigger = Mono.create(emmiter -> {
triggerRef.set(emmiter);
});
context.addApplicationListener(event -> {
if (event instanceof BindingCreatedEvent) {
if (triggerRef.get() != null) {
triggerRef.get().success();
}
}
});
RootBeanDefinition bd = (RootBeanDefinition) context.getBeanDefinition(functionProperties.getParsedDefinition()[0]);
Method factoryMethod = bd.getResolvedFactoryMethod();
if (factoryMethod == null) {
Object source = bd.getSource();
if (source instanceof MethodMetadata) {
Class<?> factory = ClassUtils.resolveClassName(((MethodMetadata) source).getDeclaringClassName(), null);
Class<?>[] params = FunctionContextUtils.getParamTypesFromBeanDefinitionFactory(factory, bd);
factoryMethod = ReflectionUtils.findMethod(factory, ((MethodMetadata) source).getMethodName(), params);
if (ObjectUtils.isEmpty(context.getBeanNamesForAnnotation(EnableBinding.class)) && functionWrapper != null && functionWrapper.isSupplier()) {
AtomicReference<MonoSink<Object>> triggerRef = new AtomicReference<>();
Publisher<Object> beginPublishingTrigger = Mono.create(emmiter -> {
triggerRef.set(emmiter);
});
context.addApplicationListener(event -> {
if (event instanceof BindingCreatedEvent) {
if (triggerRef.get() != null) {
triggerRef.get().success();
}
}
Assert.notNull(factoryMethod, "Failed to introspect factory method since it was not discovered for function '"
+ functionProperties.getDefinition() + "'");
PollableSupplier pollable = factoryMethod.getReturnType().isAssignableFrom(Supplier.class)
? AnnotationUtils.findAnnotation(factoryMethod, PollableSupplier.class)
: null;
});
if (!functionProperties.isComposeFrom() && !functionProperties.isComposeTo() && functionWrapper.isSupplier()) {
integrationFlow = this.integrationFlowFromProvidedSupplier(functionWrapper, functionInspector, beginPublishingTrigger, pollable)
.channel("output").get();
RootBeanDefinition bd = (RootBeanDefinition) context.getBeanDefinition(functionProperties.getParsedDefinition()[0]);
Method factoryMethod = bd.getResolvedFactoryMethod();
if (factoryMethod == null) {
Object source = bd.getSource();
if (source instanceof MethodMetadata) {
Class<?> factory = ClassUtils.resolveClassName(((MethodMetadata) source).getDeclaringClassName(), null);
Class<?>[] params = FunctionContextUtils.getParamTypesFromBeanDefinitionFactory(factory, bd);
factoryMethod = ReflectionUtils.findMethod(factory, ((MethodMetadata) source).getMethodName(), params);
}
}
Assert.notNull(factoryMethod, "Failed to introspect factory method since it was not discovered for function '"
+ functionProperties.getDefinition() + "'");
PollableSupplier pollable = factoryMethod.getReturnType().isAssignableFrom(Supplier.class)
? AnnotationUtils.findAnnotation(factoryMethod, PollableSupplier.class)
: null;
if (!functionProperties.isComposeFrom() && !functionProperties.isComposeTo()) {
integrationFlow = this.integrationFlowFromProvidedSupplier(functionWrapper, functionInspector, beginPublishingTrigger, pollable)
.channel("output").get();
}
}
return integrationFlow;
@@ -253,7 +250,7 @@ public class FunctionConfiguration {
}
else {
FunctionInvocationWrapper function = functionCatalog.lookup(functionProperties.getDefinition(), "application/json");
if (!function.isSupplier() && "input".equals(channelName)) {
if (/*!function.isSupplier() && */"input".equals(channelName)) {
this.postProcessForStandAloneFunction(function, messageChannel);
}
}
@@ -329,7 +326,7 @@ public class FunctionConfiguration {
*
*/
@SuppressWarnings("rawtypes")
private static class FunctionWrapper implements Function<Message<byte[]>, Message<byte[]>> {
private static class FunctionWrapper implements Function<Message<byte[]>, Object> {
private final Function function;
FunctionWrapper(Function function) {
@@ -338,8 +335,11 @@ public class FunctionConfiguration {
@SuppressWarnings("unchecked")
@Override
public Message<byte[]> apply(Message<byte[]> t) {
Message<byte[]> resultMessage = (Message<byte[]>) function.apply(t);
return resultMessage;
Object result = function.apply(t);
if (result instanceof Publisher) {
throw new IllegalStateException("Routing to functions that return Publisher is not supported in the context of Spring Cloud Stream.");
}
return (Message<byte[]>) result;
}
}
}

View File

@@ -20,15 +20,17 @@ import java.util.function.Function;
import org.junit.After;
import org.junit.Before;
import org.junit.Ignore;
import org.junit.Test;
import reactor.core.publisher.Flux;
import org.springframework.boot.WebApplicationType;
import org.springframework.boot.autoconfigure.EnableAutoConfiguration;
import org.springframework.boot.builder.SpringApplicationBuilder;
import org.springframework.cloud.function.context.FunctionProperties;
import org.springframework.cloud.function.context.config.RoutingFunction;
import org.springframework.cloud.stream.binder.test.InputDestination;
import org.springframework.cloud.stream.binder.test.OutputDestination;
import org.springframework.cloud.stream.binder.test.TestChannelBinder;
import org.springframework.cloud.stream.binder.test.TestChannelBinderConfiguration;
import org.springframework.context.ConfigurableApplicationContext;
import org.springframework.context.annotation.Bean;
@@ -44,7 +46,6 @@ import static org.assertj.core.api.Assertions.assertThat;
* @author Oleg Zhurakousky
* @since 2.2.1
*/
@Ignore
public class RoutingFunctionTests {
@After
@@ -60,13 +61,13 @@ public class RoutingFunctionTests {
}
@Test
public void testDefaultRoutingFunctionBinding() {
public void testRoutingViaExplicitEnablingAndDefinitionHeader() {
try (ConfigurableApplicationContext context = new SpringApplicationBuilder(
TestChannelBinderConfiguration.getCompleteConfiguration(
RoutingFunctionConfiguration.class))
.web(WebApplicationType.NONE)
.run("--spring.jmx.enabled=false",
"--spring.cloud.function.routing.enabled=true")) {
"--spring.cloud.stream.function.routing.enabled=true")) {
InputDestination inputDestination = context.getBean(InputDestination.class);
OutputDestination outputDestination = context
@@ -74,14 +75,88 @@ public class RoutingFunctionTests {
Message<byte[]> inputMessage = MessageBuilder
.withPayload("Hello".getBytes())
.setHeader("function.name", "echo")
.setHeader(FunctionProperties.PREFIX + ".definition", "echo")
.setHeader(MessageHeaders.CONTENT_TYPE, MimeTypeUtils.TEXT_PLAIN)
.build();
inputDestination.send(inputMessage);
Message<byte[]> outputMessage = outputDestination.receive();
assertThat(outputMessage.getPayload()).isEqualTo("Hello".getBytes());
}
}
@Test
public void testRoutingViaExplicitEnablingAndRoutingExpressionProperty() {
try (ConfigurableApplicationContext context = new SpringApplicationBuilder(
TestChannelBinderConfiguration.getCompleteConfiguration(
RoutingFunctionConfiguration.class))
.web(WebApplicationType.NONE)
.run("--spring.jmx.enabled=false",
"--spring.cloud.function.routing-expression=headers.contentType.toString().equals('text/plain') ? 'echo' : null",
"--spring.cloud.stream.function.routing.enabled=true")) {
InputDestination inputDestination = context.getBean(InputDestination.class);
OutputDestination outputDestination = context
.getBean(OutputDestination.class);
Message<byte[]> inputMessage = MessageBuilder
.withPayload("Hello".getBytes())
.setHeader(MessageHeaders.CONTENT_TYPE, MimeTypeUtils.TEXT_PLAIN)
.build();
inputDestination.send(inputMessage);
Message<byte[]> outputMessage = outputDestination.receive();
assertThat(outputMessage.getPayload()).isEqualTo("Hello".getBytes());
}
}
@Test
public void testRoutingViaExplicitEnablingAndRoutingExpressionHeader() {
try (ConfigurableApplicationContext context = new SpringApplicationBuilder(
TestChannelBinderConfiguration.getCompleteConfiguration(
RoutingFunctionConfiguration.class))
.web(WebApplicationType.NONE)
.run("--spring.jmx.enabled=false",
"--spring.cloud.stream.function.routing.enabled=true")) {
InputDestination inputDestination = context.getBean(InputDestination.class);
OutputDestination outputDestination = context
.getBean(OutputDestination.class);
Message<byte[]> inputMessage = MessageBuilder
.withPayload("Hello".getBytes())
.setHeader("spring.cloud.function.routing-expression", "'echo'")
.setHeader(MessageHeaders.CONTENT_TYPE, MimeTypeUtils.TEXT_PLAIN)
.build();
inputDestination.send(inputMessage);
Message<byte[]> outputMessage = outputDestination.receive();
assertThat(outputMessage.getPayload()).isEqualTo("Hello".getBytes());
}
}
@Test
public void testRoutingViaExplicitDefinitionAndDefinitionHeader() {
try (ConfigurableApplicationContext context = new SpringApplicationBuilder(
TestChannelBinderConfiguration.getCompleteConfiguration(
RoutingFunctionConfiguration.class))
.web(WebApplicationType.NONE)
.run("--spring.jmx.enabled=false",
"--spring.cloud.function.definition=" + RoutingFunction.FUNCTION_NAME)) {
InputDestination inputDestination = context.getBean(InputDestination.class);
OutputDestination outputDestination = context
.getBean(OutputDestination.class);
Message<byte[]> inputMessage = MessageBuilder
.withPayload("Hello".getBytes())
.setHeader("spring.cloud.function.definition", "echo|uppercase")
.setHeader(MessageHeaders.CONTENT_TYPE, MimeTypeUtils.TEXT_PLAIN)
.build();
inputDestination.send(inputMessage);
Message<byte[]> outputMessage = outputDestination.receive();
assertThat(outputMessage.getPayload()).isEqualTo("HELLO".getBytes());
}
}
@@ -92,76 +167,21 @@ public class RoutingFunctionTests {
RoutingFunctionConfiguration.class))
.web(WebApplicationType.NONE)
.run("--spring.jmx.enabled=false",
"--spring.cloud.function.routing.enabled=true")) {
"--spring.cloud.stream.function.routing.enabled=true")) {
InputDestination inputDestination = context.getBean(InputDestination.class);
OutputDestination outputDestination = context
.getBean(OutputDestination.class);
Message<byte[]> inputMessage = MessageBuilder
.withPayload("Hello".getBytes())
.setHeader("function.name", "echoFlux")
.setHeader("spring.cloud.function.definition", "echoFlux")
.setHeader(MessageHeaders.CONTENT_TYPE, MimeTypeUtils.TEXT_PLAIN)
.build();
inputDestination.send(inputMessage);
Message<byte[]> outputMessage = outputDestination.receive();
assertThat(outputMessage.getPayload()).isEqualTo("Hello".getBytes());
}
}
@Test // see RoutingFunctionEnvironmentPostProcessor
public void testExplicitRoutingFunctionBinding() {
try (ConfigurableApplicationContext context = new SpringApplicationBuilder(
TestChannelBinderConfiguration.getCompleteConfiguration(
RoutingFunctionConfiguration.class))
.web(WebApplicationType.NONE)
.run("--spring.jmx.enabled=false",
"--spring.cloud.stream.function.definition=router")) {
InputDestination inputDestination = context.getBean(InputDestination.class);
OutputDestination outputDestination = context
.getBean(OutputDestination.class);
Message<byte[]> inputMessage = MessageBuilder
.withPayload("Hello".getBytes())
.setHeader("function.name", "echo")
.setHeader(MessageHeaders.CONTENT_TYPE, MimeTypeUtils.TEXT_PLAIN)
.build();
inputDestination.send(inputMessage);
Message<byte[]> outputMessage = outputDestination.receive();
assertThat(outputMessage.getPayload()).isEqualTo("Hello".getBytes());
}
}
@Test
public void testCompositionViaFunctionName() {
try (ConfigurableApplicationContext context = new SpringApplicationBuilder(
TestChannelBinderConfiguration.getCompleteConfiguration(
RoutingFunctionConfiguration.class))
.web(WebApplicationType.NONE)
.run("--spring.jmx.enabled=false",
"--spring.cloud.stream.function.definition=router")) {
InputDestination inputDestination = context.getBean(InputDestination.class);
OutputDestination outputDestination = context
.getBean(OutputDestination.class);
Message<byte[]> inputMessage = MessageBuilder
.withPayload("Hello".getBytes())
.setHeader("function.name", "echo|uppercase")
.setHeader(MessageHeaders.CONTENT_TYPE, MimeTypeUtils.TEXT_PLAIN)
.build();
inputDestination.send(inputMessage);
Message<byte[]> outputMessage = outputDestination.receive();
assertThat(outputMessage.getPayload()).isEqualTo("HELLO".getBytes());
TestChannelBinder binder = context.getBean(TestChannelBinder.class);
Throwable ex = ((Exception) binder.getLastError().getPayload()).getCause();
assertThat(ex).isInstanceOf(IllegalStateException.class);
assertThat(ex.getMessage()).isEqualTo("Routing to functions that return Publisher is not supported in the context of Spring Cloud Stream.");
}
}
@@ -173,7 +193,7 @@ public class RoutingFunctionTests {
RoutingFunctionConfiguration.class))
.web(WebApplicationType.NONE)
.run("--spring.jmx.enabled=false",
"--spring.cloud.stream.function.definition=router")) {
"--spring.cloud.stream.function.definition=" + RoutingFunction.FUNCTION_NAME)) {
InputDestination inputDestination = context.getBean(InputDestination.class);
OutputDestination outputDestination = context
@@ -181,7 +201,7 @@ public class RoutingFunctionTests {
Message<byte[]> inputMessage = MessageBuilder
.withPayload("{\"name\":\"bob\"}".getBytes())
.setHeader("function.name", "pojoecho")
.setHeader("spring.cloud.function.definition", "pojoecho")
.build();
inputDestination.send(inputMessage);
@@ -198,8 +218,8 @@ public class RoutingFunctionTests {
RoutingFunctionConfiguration.class))
.web(WebApplicationType.NONE)
.run("--spring.jmx.enabled=false",
"--spring.cloud.stream.function.definition=enrich|router",
"--spring.cloud.function.routing.enabled=true")) {
"--spring.cloud.stream.function.definition=enrich|" + RoutingFunction.FUNCTION_NAME,
"--spring.cloud.stream.function.routing.enabled=true")) {
InputDestination inputDestination = context.getBean(InputDestination.class);
OutputDestination outputDestination = context
@@ -221,7 +241,7 @@ public class RoutingFunctionTests {
RoutingFunctionConfiguration.class))
.web(WebApplicationType.NONE)
.run("--spring.jmx.enabled=false",
"--spring.cloud.stream.function.definition=enrich|router")) {
"--spring.cloud.stream.function.definition=enrich|" + RoutingFunction.FUNCTION_NAME)) {
InputDestination inputDestination = context.getBean(InputDestination.class);
OutputDestination outputDestination = context
@@ -243,8 +263,8 @@ public class RoutingFunctionTests {
RoutingFunctionConfiguration.class))
.web(WebApplicationType.NONE)
.run("--spring.jmx.enabled=false",
"--spring.cloud.stream.function.definition=enrich|router|reverse",
"--spring.cloud.function.routing.enabled=true")) {
"--spring.cloud.stream.function.definition=enrich|" + RoutingFunction.FUNCTION_NAME + "|reverse",
"--spring.cloud.stream.function.routing.enabled=true")) {
InputDestination inputDestination = context.getBean(InputDestination.class);
OutputDestination outputDestination = context
@@ -266,7 +286,7 @@ public class RoutingFunctionTests {
RoutingFunctionConfiguration.class))
.web(WebApplicationType.NONE)
.run("--spring.jmx.enabled=false",
"--spring.cloud.stream.function.definition=enrich|router|reverse")) {
"--spring.cloud.stream.function.definition=enrich|" + RoutingFunction.FUNCTION_NAME + "|reverse")) {
InputDestination inputDestination = context.getBean(InputDestination.class);
OutputDestination outputDestination = context
@@ -314,7 +334,7 @@ public class RoutingFunctionTests {
public Function<Message<String>, Message<String>> enrich() {
return x -> {
System.out.println("===> enrich");
return MessageBuilder.withPayload(x.getPayload()).setHeader("function.name", "uppercase").build();
return MessageBuilder.withPayload(x.getPayload()).setHeader("spring.cloud.function.definition", "uppercase").build();
};
}