GH-56 Added support for Function<Flux, Mono> and <Mono, Flux>
Resolves #56 Resolves #218
This commit is contained in:
@@ -52,9 +52,11 @@ import org.springframework.cloud.function.context.catalog.FunctionUnregistration
|
|||||||
import org.springframework.cloud.function.core.FluxConsumer;
|
import org.springframework.cloud.function.core.FluxConsumer;
|
||||||
import org.springframework.cloud.function.core.FluxFunction;
|
import org.springframework.cloud.function.core.FluxFunction;
|
||||||
import org.springframework.cloud.function.core.FluxSupplier;
|
import org.springframework.cloud.function.core.FluxSupplier;
|
||||||
|
import org.springframework.cloud.function.core.FluxToMonoFunction;
|
||||||
import org.springframework.cloud.function.core.IsolatedConsumer;
|
import org.springframework.cloud.function.core.IsolatedConsumer;
|
||||||
import org.springframework.cloud.function.core.IsolatedFunction;
|
import org.springframework.cloud.function.core.IsolatedFunction;
|
||||||
import org.springframework.cloud.function.core.IsolatedSupplier;
|
import org.springframework.cloud.function.core.IsolatedSupplier;
|
||||||
|
import org.springframework.cloud.function.core.MonoToFluxFunction;
|
||||||
import org.springframework.cloud.function.json.GsonMapper;
|
import org.springframework.cloud.function.json.GsonMapper;
|
||||||
import org.springframework.cloud.function.json.JacksonMapper;
|
import org.springframework.cloud.function.json.JacksonMapper;
|
||||||
import org.springframework.context.ApplicationEventPublisher;
|
import org.springframework.context.ApplicationEventPublisher;
|
||||||
@@ -349,7 +351,22 @@ public class ContextFunctionCatalogAutoConfiguration {
|
|||||||
else if (a instanceof Function && b instanceof Function) {
|
else if (a instanceof Function && b instanceof Function) {
|
||||||
Function<Object, Object> function1 = (Function<Object, Object>) a;
|
Function<Object, Object> function1 = (Function<Object, Object>) a;
|
||||||
Function<Object, Object> function2 = (Function<Object, Object>) b;
|
Function<Object, Object> function2 = (Function<Object, Object>) b;
|
||||||
return function1.andThen(function2);
|
if (function1 instanceof FluxToMonoFunction) {
|
||||||
|
if (function2 instanceof MonoToFluxFunction) {
|
||||||
|
return function1.andThen(function2);
|
||||||
|
}
|
||||||
|
else {
|
||||||
|
throw new IllegalStateException("The provided function is finite (i.e., returns Mono<?>) "
|
||||||
|
+ "therefore it can *only* be composed with compatible function (i.e., Function<Mono, Flux>");
|
||||||
|
}
|
||||||
|
}
|
||||||
|
else if (function2 instanceof FluxToMonoFunction) {
|
||||||
|
return new FluxToMonoFunction<Object, Object>(((Function<Flux<Object>, Flux<Object>>)a)
|
||||||
|
.andThen(((FluxToMonoFunction<Object,Object>) b).getTarget()));
|
||||||
|
}
|
||||||
|
else {
|
||||||
|
return function1.andThen(function2);
|
||||||
|
}
|
||||||
}
|
}
|
||||||
else if (a instanceof Function && b instanceof Consumer) {
|
else if (a instanceof Function && b instanceof Consumer) {
|
||||||
Function<Object, Object> function = (Function<Object, Object>) a;
|
Function<Object, Object> function = (Function<Object, Object>) a;
|
||||||
@@ -490,6 +507,12 @@ public class ContextFunctionCatalogAutoConfiguration {
|
|||||||
}
|
}
|
||||||
registration.target(target);
|
registration.target(target);
|
||||||
}
|
}
|
||||||
|
if (Mono.class.isAssignableFrom(type.getOutputWrapper())) {
|
||||||
|
registration.target(new FluxToMonoFunction<>((Function) target));
|
||||||
|
}
|
||||||
|
else if (Mono.class.isAssignableFrom(type.getInputWrapper())) {
|
||||||
|
registration.target(new MonoToFluxFunction<>((Function) target));
|
||||||
|
}
|
||||||
return registration;
|
return registration;
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
@@ -120,6 +120,33 @@ public class BeanFactoryFunctionCatalogTests {
|
|||||||
assertThat(foos.apply(Flux.just(2)).blockFirst()).isEqualTo("Hello 4");
|
assertThat(foos.apply(Flux.just(2)).blockFirst()).isEqualTo("Hello 4");
|
||||||
}
|
}
|
||||||
|
|
||||||
|
@Test
|
||||||
|
public void composeWithFiniteFunction() {
|
||||||
|
Function<String, String> func1 = x -> x.toUpperCase();
|
||||||
|
processor.register(new FunctionRegistration<>(func1, "func1"));
|
||||||
|
processor.register(new FunctionRegistration<>(new FluxThenMonoFunction(), "func2"));
|
||||||
|
Function<Flux<String>, Mono<Long>> foos = processor.lookup(Function.class, "func1,func2");
|
||||||
|
assertThat(foos.apply(Flux.fromArray(new String[] {"a", "b", "c"})).block()).isEqualTo(3);
|
||||||
|
}
|
||||||
|
|
||||||
|
@Test
|
||||||
|
public void composeWithFiniteFunctionAndContinueWithCompatible() {
|
||||||
|
Function<String, String> func1 = x -> x.toUpperCase();
|
||||||
|
processor.register(new FunctionRegistration<>(func1, "func1"));
|
||||||
|
processor.register(new FunctionRegistration<>(new FluxThenMonoFunction(), "func2"));
|
||||||
|
processor.register(new FunctionRegistration<>(new MonoThenFluxFunction(), "func3"));
|
||||||
|
Function<Flux<String>, Flux<Integer>> foos = processor.lookup(Function.class, "func1,func2,func3");
|
||||||
|
assertThat(foos.apply(Flux.fromArray(new String[] {"a", "b", "c"})).collectList().block().size()).isEqualTo(3);
|
||||||
|
}
|
||||||
|
|
||||||
|
@Test(expected=IllegalStateException.class)
|
||||||
|
public void composeIncompatibleFunctions() {
|
||||||
|
Function<String, String> func1 = x -> x.toUpperCase();
|
||||||
|
processor.register(new FunctionRegistration<>(func1, "func1"));
|
||||||
|
processor.register(new FunctionRegistration<>(new FluxThenMonoFunction(), "func2"));
|
||||||
|
processor.lookup(Function.class, "func2,func1");
|
||||||
|
}
|
||||||
|
|
||||||
@Test
|
@Test
|
||||||
public void composeSupplier() {
|
public void composeSupplier() {
|
||||||
processor.register(new FunctionRegistration<>(new Source(), "numbers"));
|
processor.register(new FunctionRegistration<>(new Source(), "numbers"));
|
||||||
@@ -231,4 +258,20 @@ public class BeanFactoryFunctionCatalogTests {
|
|||||||
|
|
||||||
}
|
}
|
||||||
|
|
||||||
|
protected static class FluxThenMonoFunction implements Function<Flux<String>, Mono<Long>> {
|
||||||
|
|
||||||
|
@Override
|
||||||
|
public Mono<Long> apply(Flux<String> t) {
|
||||||
|
return t.count();
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
protected static class MonoThenFluxFunction implements Function<Mono<Long>, Flux<Integer>> {
|
||||||
|
|
||||||
|
@Override
|
||||||
|
public Flux<Integer> apply(Mono<Long> t) {
|
||||||
|
return Flux.range(0, Integer.parseInt(t.block().toString()));
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -0,0 +1,53 @@
|
|||||||
|
/*
|
||||||
|
* Copyright 2018 the original author or authors.
|
||||||
|
*
|
||||||
|
* Licensed under the Apache License, Version 2.0 (the "License");
|
||||||
|
* you may not use this file except in compliance with the License.
|
||||||
|
* You may obtain a copy of the License at
|
||||||
|
*
|
||||||
|
* http://www.apache.org/licenses/LICENSE-2.0
|
||||||
|
*
|
||||||
|
* Unless required by applicable law or agreed to in writing, software
|
||||||
|
* distributed under the License is distributed on an "AS IS" BASIS,
|
||||||
|
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
|
||||||
|
* See the License for the specific language governing permissions and
|
||||||
|
* limitations under the License.
|
||||||
|
*/
|
||||||
|
|
||||||
|
package org.springframework.cloud.function.core;
|
||||||
|
|
||||||
|
import java.util.function.Function;
|
||||||
|
|
||||||
|
import reactor.core.publisher.Flux;
|
||||||
|
import reactor.core.publisher.Mono;
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Marker wrapper for target Function<Flux, Mono>
|
||||||
|
*
|
||||||
|
* @author Oleg Zhurakousky
|
||||||
|
* @since 2.0
|
||||||
|
*
|
||||||
|
* @param <T> type of {@link Flux} input of the target function
|
||||||
|
* @param <R> type of {@link Mono} output of the target function
|
||||||
|
*/
|
||||||
|
public class FluxToMonoFunction<I,O> implements Function<Flux<I>, Mono<O>>, FluxWrapper<Function<Flux<I>, Mono<O>>> {
|
||||||
|
|
||||||
|
private final Function<Flux<I>, Mono<O>> function;
|
||||||
|
|
||||||
|
/**
|
||||||
|
* @param function target function
|
||||||
|
*/
|
||||||
|
public FluxToMonoFunction(Function<Flux<I>, Mono<O>> function) {
|
||||||
|
this.function = function;
|
||||||
|
}
|
||||||
|
|
||||||
|
@Override
|
||||||
|
public Function<Flux<I>, Mono<O>> getTarget() {
|
||||||
|
return function;
|
||||||
|
}
|
||||||
|
|
||||||
|
@Override
|
||||||
|
public Mono<O> apply(Flux<I> input) {
|
||||||
|
return function.apply(input);
|
||||||
|
}
|
||||||
|
}
|
||||||
@@ -0,0 +1,54 @@
|
|||||||
|
/*
|
||||||
|
* Copyright 2018 the original author or authors.
|
||||||
|
*
|
||||||
|
* Licensed under the Apache License, Version 2.0 (the "License");
|
||||||
|
* you may not use this file except in compliance with the License.
|
||||||
|
* You may obtain a copy of the License at
|
||||||
|
*
|
||||||
|
* http://www.apache.org/licenses/LICENSE-2.0
|
||||||
|
*
|
||||||
|
* Unless required by applicable law or agreed to in writing, software
|
||||||
|
* distributed under the License is distributed on an "AS IS" BASIS,
|
||||||
|
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
|
||||||
|
* See the License for the specific language governing permissions and
|
||||||
|
* limitations under the License.
|
||||||
|
*/
|
||||||
|
|
||||||
|
package org.springframework.cloud.function.core;
|
||||||
|
|
||||||
|
import java.util.function.Function;
|
||||||
|
|
||||||
|
import reactor.core.publisher.Flux;
|
||||||
|
import reactor.core.publisher.Mono;
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Marker wrapper for target Function<Mono, Flux>
|
||||||
|
*
|
||||||
|
* @author Oleg Zhurakousky
|
||||||
|
* @since 2.0
|
||||||
|
*
|
||||||
|
* @param <T> type of {@link Mono} input of the target function
|
||||||
|
* @param <R> type of {@link Flux} output of the target function
|
||||||
|
*/
|
||||||
|
public class MonoToFluxFunction<I,O> implements Function<Mono<I>, Flux<O>>, FluxWrapper<Function<Mono<I>, Flux<O>>> {
|
||||||
|
|
||||||
|
private final Function<Mono<I>, Flux<O>> function;
|
||||||
|
|
||||||
|
/**
|
||||||
|
* @param function target function
|
||||||
|
* @param names name(s) of the target function (optional)
|
||||||
|
*/
|
||||||
|
public MonoToFluxFunction(Function<Mono<I>, Flux<O>> function) {
|
||||||
|
this.function = function;
|
||||||
|
}
|
||||||
|
|
||||||
|
@Override
|
||||||
|
public Function<Mono<I>, Flux<O>> getTarget() {
|
||||||
|
return function;
|
||||||
|
}
|
||||||
|
|
||||||
|
@Override
|
||||||
|
public Flux<O> apply(Mono<I> input) {
|
||||||
|
return function.apply(input);
|
||||||
|
}
|
||||||
|
}
|
||||||
Reference in New Issue
Block a user