Consolidated all function wrappers around WrappedFunction
This commit is contained in:
@@ -36,6 +36,7 @@ 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.FluxToMonoFunction;
|
||||||
import org.springframework.cloud.function.core.FluxedConsumer;
|
import org.springframework.cloud.function.core.FluxedConsumer;
|
||||||
|
import org.springframework.cloud.function.core.FluxedFunction;
|
||||||
import org.springframework.cloud.function.core.MonoToFluxFunction;
|
import org.springframework.cloud.function.core.MonoToFluxFunction;
|
||||||
import org.springframework.util.Assert;
|
import org.springframework.util.Assert;
|
||||||
import org.springframework.util.CollectionUtils;
|
import org.springframework.util.CollectionUtils;
|
||||||
@@ -171,7 +172,7 @@ public class FunctionRegistration<T> implements BeanNameAware {
|
|||||||
target = (S) new FluxedConsumer((Consumer<?>) target);
|
target = (S) new FluxedConsumer((Consumer<?>) target);
|
||||||
}
|
}
|
||||||
else if (target instanceof Function) {
|
else if (target instanceof Function) {
|
||||||
// target = (S) new FluxedFunction((Function<?, ?>) target);
|
target = (S) new FluxedFunction((Function<?, ?>) target);
|
||||||
}
|
}
|
||||||
|
|
||||||
result = result.target(target).names(this.names)
|
result = result.target(target).names(this.names)
|
||||||
|
|||||||
@@ -22,7 +22,12 @@ import reactor.core.publisher.Flux;
|
|||||||
import reactor.core.publisher.Mono;
|
import reactor.core.publisher.Mono;
|
||||||
|
|
||||||
/**
|
/**
|
||||||
* Marker wrapper for target {@code Function<Flux<?>, Mono<?>>}.
|
* Wrapper to mark function {@code Function<Flux<?>, Mono<?>>}.
|
||||||
|
*
|
||||||
|
* While it may look similar to {@link FluxedConsumer}
|
||||||
|
* the fundamental difference is that this class represents a function that
|
||||||
|
* returns {@link Mono} of type {@code <O>}, while {@link FluxedConsumer} is
|
||||||
|
* a consumer that has been decorated as {@code Function<Flux<?>, Mono<Void>>}.
|
||||||
*
|
*
|
||||||
* @param <I> type of {@link Flux} input of the target function
|
* @param <I> type of {@link Flux} input of the target function
|
||||||
* @param <O> type of {@link Mono} output of the target function
|
* @param <O> type of {@link Mono} output of the target function
|
||||||
@@ -30,25 +35,15 @@ import reactor.core.publisher.Mono;
|
|||||||
* @since 2.0
|
* @since 2.0
|
||||||
*/
|
*/
|
||||||
public class FluxToMonoFunction<I, O>
|
public class FluxToMonoFunction<I, O>
|
||||||
implements Function<Flux<I>, Mono<O>>, FluxWrapper<Function<Flux<I>, Mono<O>>> {
|
extends WrappedFunction<I, O, Flux<I>, Mono<O>, Function<Flux<I>, Mono<O>>> {
|
||||||
|
|
||||||
private final Function<Flux<I>, Mono<O>> function;
|
public FluxToMonoFunction(Function<Flux<I>, Mono<O>> target) {
|
||||||
|
super(target);
|
||||||
/**
|
|
||||||
* @param function target function
|
|
||||||
*/
|
|
||||||
public FluxToMonoFunction(Function<Flux<I>, Mono<O>> function) {
|
|
||||||
this.function = function;
|
|
||||||
}
|
|
||||||
|
|
||||||
@Override
|
|
||||||
public Function<Flux<I>, Mono<O>> getTarget() {
|
|
||||||
return this.function;
|
|
||||||
}
|
}
|
||||||
|
|
||||||
@Override
|
@Override
|
||||||
public Mono<O> apply(Flux<I> input) {
|
public Mono<O> apply(Flux<I> input) {
|
||||||
return this.function.apply(input);
|
return this.getTarget().apply(input);
|
||||||
}
|
}
|
||||||
|
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -30,25 +30,15 @@ import reactor.core.publisher.Mono;
|
|||||||
* @since 2.0
|
* @since 2.0
|
||||||
*/
|
*/
|
||||||
public class MonoToFluxFunction<I, O>
|
public class MonoToFluxFunction<I, O>
|
||||||
implements Function<Mono<I>, Flux<O>>, FluxWrapper<Function<Mono<I>, Flux<O>>> {
|
extends WrappedFunction<I, O, Mono<I>, Flux<O>, Function<Mono<I>, Flux<O>>> {
|
||||||
|
|
||||||
private final Function<Mono<I>, Flux<O>> function;
|
public MonoToFluxFunction(Function<Mono<I>, Flux<O>> target) {
|
||||||
|
super(target);
|
||||||
/**
|
|
||||||
* @param function target function
|
|
||||||
*/
|
|
||||||
public MonoToFluxFunction(Function<Mono<I>, Flux<O>> function) {
|
|
||||||
this.function = function;
|
|
||||||
}
|
|
||||||
|
|
||||||
@Override
|
|
||||||
public Function<Mono<I>, Flux<O>> getTarget() {
|
|
||||||
return this.function;
|
|
||||||
}
|
}
|
||||||
|
|
||||||
@Override
|
@Override
|
||||||
public Flux<O> apply(Mono<I> input) {
|
public Flux<O> apply(Mono<I> input) {
|
||||||
return this.function.apply(input);
|
return this.getTarget().apply(input);
|
||||||
}
|
}
|
||||||
|
|
||||||
}
|
}
|
||||||
|
|||||||
Reference in New Issue
Block a user