diff --git a/docs/src/main/asciidoc/spring-cloud-stream.adoc b/docs/src/main/asciidoc/spring-cloud-stream.adoc index 14244456b..d9aa5e03a 100644 --- a/docs/src/main/asciidoc/spring-cloud-stream.adoc +++ b/docs/src/main/asciidoc/spring-cloud-stream.adoc @@ -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 even() { + return value -> { + System.out.println("EVEN: " + value); + }; + } + + @Bean + public Consumer 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"); diff --git a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/config/BinderFactoryAutoConfiguration.java b/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/config/BinderFactoryAutoConfiguration.java index e3bc395ee..ea3e196ba 100644 --- a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/config/BinderFactoryAutoConfiguration.java +++ b/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/config/BinderFactoryAutoConfiguration.java @@ -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) { diff --git a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/function/FunctionConfiguration.java b/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/function/FunctionConfiguration.java index a7ed72d94..62cab49ac 100644 --- a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/function/FunctionConfiguration.java +++ b/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/function/FunctionConfiguration.java @@ -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> triggerRef = new AtomicReference<>(); - Publisher 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> triggerRef = new AtomicReference<>(); + Publisher 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> { + private static class FunctionWrapper implements Function, Object> { private final Function function; FunctionWrapper(Function function) { @@ -338,8 +335,11 @@ public class FunctionConfiguration { @SuppressWarnings("unchecked") @Override public Message apply(Message t) { - Message resultMessage = (Message) 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) result; } } } diff --git a/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/function/RoutingFunctionTests.java b/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/function/RoutingFunctionTests.java index 88690b146..05acfa37e 100644 --- a/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/function/RoutingFunctionTests.java +++ b/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/function/RoutingFunctionTests.java @@ -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 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 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 inputMessage = MessageBuilder + .withPayload("Hello".getBytes()) + .setHeader(MessageHeaders.CONTENT_TYPE, MimeTypeUtils.TEXT_PLAIN) + .build(); + inputDestination.send(inputMessage); + + Message 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 inputMessage = MessageBuilder + .withPayload("Hello".getBytes()) + .setHeader("spring.cloud.function.routing-expression", "'echo'") + .setHeader(MessageHeaders.CONTENT_TYPE, MimeTypeUtils.TEXT_PLAIN) + .build(); + inputDestination.send(inputMessage); + + Message 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 inputMessage = MessageBuilder + .withPayload("Hello".getBytes()) + .setHeader("spring.cloud.function.definition", "echo|uppercase") + .setHeader(MessageHeaders.CONTENT_TYPE, MimeTypeUtils.TEXT_PLAIN) + .build(); + inputDestination.send(inputMessage); + + Message 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 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 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 inputMessage = MessageBuilder - .withPayload("Hello".getBytes()) - .setHeader("function.name", "echo") - .setHeader(MessageHeaders.CONTENT_TYPE, MimeTypeUtils.TEXT_PLAIN) - .build(); - inputDestination.send(inputMessage); - - Message 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 inputMessage = MessageBuilder - .withPayload("Hello".getBytes()) - .setHeader("function.name", "echo|uppercase") - .setHeader(MessageHeaders.CONTENT_TYPE, MimeTypeUtils.TEXT_PLAIN) - .build(); - inputDestination.send(inputMessage); - - Message 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 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> 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(); }; }