support function composition for web and stream

This commit is contained in:
markfisher
2016-10-01 16:52:59 -04:00
parent b78024e7ea
commit 8e5d631db9
3 changed files with 28 additions and 5 deletions

View File

@@ -16,4 +16,14 @@ Run a REST Microservice using that Function:
./web.sh -p /words -f uppercase
```
To compose Functions:
(assuming the `uppercase` function was already registered as above)
```
./register.sh -n pluralize -f "f->f.map(s->s+\"S\")"
./web.sh -p /words -f uppercase,pluralize
```
(more docs soon)

View File

@@ -16,6 +16,8 @@
package org.springframework.cloud.function.stream;
import java.util.function.Function;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.boot.context.properties.EnableConfigurationProperties;
import org.springframework.cloud.function.invoker.AbstractFunctionInvoker;
@@ -24,6 +26,9 @@ import org.springframework.cloud.function.registry.FunctionRegistry;
import org.springframework.cloud.stream.annotation.EnableBinding;
import org.springframework.cloud.stream.messaging.Processor;
import org.springframework.context.annotation.Bean;
import org.springframework.util.StringUtils;
import reactor.core.publisher.Flux;
/**
* @author Mark Fisher
@@ -41,7 +46,11 @@ public class StreamConfiguration {
}
@Bean
public AbstractFunctionInvoker<?,?> invoker() {
return new StreamListeningFunctionInvoker(registry().lookup(properties.getName()));
public AbstractFunctionInvoker<?,?> invoker(FunctionRegistry registry) {
String name = properties.getName();
Function<Flux<Object>, Flux<Object>> function = (name.indexOf(',') == -1)
? registry.lookup(name)
: registry.compose(StringUtils.commaDelimitedListToStringArray(name));
return new StreamListeningFunctionInvoker(function);
}
}

View File

@@ -19,24 +19,25 @@ package org.springframework.cloud.function.web;
import static org.springframework.http.MediaType.TEXT_PLAIN;
import static org.springframework.web.reactive.function.BodyExtractors.toFlux;
import static org.springframework.web.reactive.function.BodyInserters.fromPublisher;
import static org.springframework.web.reactive.function.RequestPredicates.contentType;
import static org.springframework.web.reactive.function.RequestPredicates.POST;
import static org.springframework.web.reactive.function.RequestPredicates.contentType;
import java.util.concurrent.Executors;
import java.util.function.Function;
import org.reactivestreams.Publisher;
import org.springframework.beans.factory.DisposableBean;
import org.springframework.beans.factory.InitializingBean;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.boot.context.properties.EnableConfigurationProperties;
import org.springframework.cloud.function.registry.FileSystemFunctionRegistry;
import org.springframework.cloud.function.registry.FunctionRegistry;
import org.springframework.cloud.function.registry.InMemoryFunctionRegistry;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import org.springframework.http.server.reactive.HttpHandler;
import org.springframework.http.server.reactive.ReactorHttpHandlerAdapter;
import org.springframework.util.StringUtils;
import org.springframework.web.reactive.function.Request;
import org.springframework.web.reactive.function.Response;
import org.springframework.web.reactive.function.RouterFunction;
@@ -65,7 +66,10 @@ public class RestConfiguration {
@Bean
public HttpHandler httpHandler(FunctionRegistry registry) {
Function<Flux<String>, Flux<String>> function = registry.lookup(functionProperties.getName());
String name = functionProperties.getName();
Function<Flux<String>, Flux<String>> function = (name.indexOf(',') == -1)
? registry.lookup(name)
: registry.compose(StringUtils.commaDelimitedListToStringArray(name));
FunctionInvokingHandler handler = new FunctionInvokingHandler(function);
RouterFunction<Publisher<String>> route = RouterFunctions.route(
POST(webProperties.getPath()).and(contentType(TEXT_PLAIN)), handler::handleText);