GH-1611 Removed FluxedConsumerWrapper and FunctionCatalogWrapper
- Upgraded to spring-cloud-function 2.0.2.BUILD-SNAPSHOT - General cleanup and polishing Resolves #1611
This commit is contained in:
2
pom.xml
2
pom.xml
@@ -28,7 +28,7 @@
|
||||
<reactor.version>Californium-RELEASE</reactor.version>
|
||||
<kryo-shaded.version>3.0.3</kryo-shaded.version>
|
||||
<objenesis.version>2.1</objenesis.version>
|
||||
<spring-cloud-function.version>2.0.0.BUILD-SNAPSHOT
|
||||
<spring-cloud-function.version>2.0.2.BUILD-SNAPSHOT
|
||||
</spring-cloud-function.version>
|
||||
|
||||
<maven-checkstyle-plugin.failsOnError>true</maven-checkstyle-plugin.failsOnError>
|
||||
|
||||
@@ -1,52 +0,0 @@
|
||||
/*
|
||||
* Copyright 2018-2019 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.stream.function;
|
||||
|
||||
import java.util.function.Consumer;
|
||||
import java.util.function.Function;
|
||||
|
||||
import reactor.core.publisher.Flux;
|
||||
import reactor.core.publisher.Mono;
|
||||
|
||||
import org.springframework.cloud.function.core.FluxWrapper;
|
||||
|
||||
/**
|
||||
* @author Oleg Zhurakousky
|
||||
* @since 2.1
|
||||
*
|
||||
* Will most likely be moved to SCF
|
||||
*/
|
||||
class FluxedConsumerWrapper<T>
|
||||
implements Function<Flux<T>, Mono<Void>>, FluxWrapper<Consumer<Flux<T>>> {
|
||||
|
||||
private final Consumer<Flux<T>> consumer;
|
||||
|
||||
FluxedConsumerWrapper(Consumer<Flux<T>> consumer) {
|
||||
this.consumer = consumer;
|
||||
}
|
||||
|
||||
@Override
|
||||
public Consumer<Flux<T>> getTarget() {
|
||||
return this.consumer;
|
||||
}
|
||||
|
||||
@Override
|
||||
public Mono<Void> apply(Flux<T> t) {
|
||||
return Mono.fromRunnable(() -> this.consumer.accept(t));
|
||||
}
|
||||
|
||||
}
|
||||
@@ -1,54 +0,0 @@
|
||||
/*
|
||||
* Copyright 2018-2019 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.stream.function;
|
||||
|
||||
import org.springframework.cloud.function.context.FunctionCatalog;
|
||||
import org.springframework.util.Assert;
|
||||
|
||||
/**
|
||||
* @author David Turanski
|
||||
* @author Oleg Zhurakousky
|
||||
* @since 2.1
|
||||
**/
|
||||
class FunctionCatalogWrapper {
|
||||
|
||||
private final FunctionCatalog catalog;
|
||||
|
||||
FunctionCatalogWrapper(FunctionCatalog catalog) {
|
||||
this.catalog = catalog;
|
||||
}
|
||||
|
||||
<T> T lookup(Class<T> functionType, String name) {
|
||||
T function = this.catalog.lookup(functionType, name);
|
||||
Assert.notNull(function,
|
||||
functionType == null
|
||||
? String.format("User provided Function '%s' cannot be located.",
|
||||
name)
|
||||
: String.format("User provided %s '%s' cannot be located.",
|
||||
functionType.getSimpleName(), name));
|
||||
return function;
|
||||
}
|
||||
|
||||
<T> T lookup(String name) {
|
||||
return lookup(null, name);
|
||||
}
|
||||
|
||||
<T> boolean contains(Class<T> functionType, String name) {
|
||||
return this.catalog.lookup(functionType, name) != null;
|
||||
}
|
||||
|
||||
}
|
||||
@@ -49,6 +49,7 @@ import org.springframework.messaging.SubscribableChannel;
|
||||
@EnableConfigurationProperties(StreamFunctionProperties.class)
|
||||
public class FunctionConfiguration {
|
||||
|
||||
|
||||
@Autowired(required = false)
|
||||
private Source source;
|
||||
|
||||
@@ -60,7 +61,7 @@ public class FunctionConfiguration {
|
||||
|
||||
@Bean
|
||||
public IntegrationFlowFunctionSupport functionSupport(
|
||||
FunctionCatalogWrapper functionCatalog, FunctionInspector functionInspector,
|
||||
FunctionCatalog functionCatalog, FunctionInspector functionInspector,
|
||||
CompositeMessageConverterFactory messageConverterFactory,
|
||||
StreamFunctionProperties functionProperties,
|
||||
BindingServiceProperties bindingServiceProperties) {
|
||||
@@ -68,10 +69,10 @@ public class FunctionConfiguration {
|
||||
messageConverterFactory, functionProperties, bindingServiceProperties);
|
||||
}
|
||||
|
||||
@Bean
|
||||
public FunctionCatalogWrapper functionCatalogWrapper(FunctionCatalog catalog) {
|
||||
return new FunctionCatalogWrapper(catalog);
|
||||
}
|
||||
// @Bean
|
||||
// public FunctionCatalogWrapper functionCatalogWrapper(FunctionCatalog catalog) {
|
||||
// return new FunctionCatalogWrapper(catalog);
|
||||
// }
|
||||
|
||||
/**
|
||||
* This configuration creates an instance of the {@link IntegrationFlow} from standard
|
||||
|
||||
@@ -20,7 +20,6 @@ import java.lang.reflect.Field;
|
||||
import java.time.Duration;
|
||||
import java.util.Map;
|
||||
import java.util.concurrent.atomic.AtomicReference;
|
||||
import java.util.function.Consumer;
|
||||
import java.util.function.Function;
|
||||
|
||||
import org.apache.commons.logging.Log;
|
||||
@@ -28,6 +27,7 @@ import org.apache.commons.logging.LogFactory;
|
||||
import reactor.core.publisher.Flux;
|
||||
import reactor.core.publisher.Mono;
|
||||
|
||||
import org.springframework.cloud.function.context.FunctionCatalog;
|
||||
import org.springframework.cloud.function.context.FunctionType;
|
||||
import org.springframework.cloud.function.context.catalog.FunctionInspector;
|
||||
import org.springframework.cloud.stream.binder.ConsumerProperties;
|
||||
@@ -85,15 +85,15 @@ class FunctionInvoker<I, O> implements Function<Flux<Message<I>>, Flux<Message<O
|
||||
private final StreamFunctionProperties functionProperties;
|
||||
|
||||
FunctionInvoker(StreamFunctionProperties functionProperties,
|
||||
FunctionCatalogWrapper functionCatalog, FunctionInspector functionInspector,
|
||||
FunctionCatalog functionCatalog, FunctionInspector functionInspector,
|
||||
CompositeMessageConverterFactory compositeMessageConverterFactory) {
|
||||
this(functionProperties, functionCatalog, functionInspector,
|
||||
compositeMessageConverterFactory, null);
|
||||
}
|
||||
|
||||
@SuppressWarnings({ "unchecked", "rawtypes" })
|
||||
@SuppressWarnings("unchecked")
|
||||
FunctionInvoker(StreamFunctionProperties functionProperties,
|
||||
FunctionCatalogWrapper functionCatalog, FunctionInspector functionInspector,
|
||||
FunctionCatalog functionCatalog, FunctionInspector functionInspector,
|
||||
CompositeMessageConverterFactory compositeMessageConverterFactory,
|
||||
MessageChannel errorChannel) {
|
||||
|
||||
@@ -101,9 +101,7 @@ class FunctionInvoker<I, O> implements Function<Flux<Message<I>>, Flux<Message<O
|
||||
Object originalUserFunction = functionCatalog
|
||||
.lookup(functionProperties.getDefinition());
|
||||
|
||||
this.userFunction = originalUserFunction instanceof Consumer
|
||||
? new FluxedConsumerWrapper<>((Consumer) originalUserFunction)
|
||||
: (Function<Flux<?>, Flux<?>>) originalUserFunction;
|
||||
this.userFunction = (Function<Flux<?>, Flux<?>>) originalUserFunction;
|
||||
|
||||
Assert.isInstanceOf(Function.class, this.userFunction);
|
||||
this.messageConverter = compositeMessageConverterFactory
|
||||
|
||||
@@ -48,7 +48,7 @@ import org.springframework.util.StringUtils;
|
||||
*/
|
||||
public class IntegrationFlowFunctionSupport {
|
||||
|
||||
private final FunctionCatalogWrapper functionCatalog;
|
||||
private final FunctionCatalog functionCatalog;
|
||||
|
||||
private final FunctionInspector functionInspector;
|
||||
|
||||
@@ -59,7 +59,7 @@ public class IntegrationFlowFunctionSupport {
|
||||
@Autowired
|
||||
private MessageChannel errorChannel;
|
||||
|
||||
IntegrationFlowFunctionSupport(FunctionCatalogWrapper functionCatalog,
|
||||
IntegrationFlowFunctionSupport(FunctionCatalog functionCatalog,
|
||||
FunctionInspector functionInspector,
|
||||
CompositeMessageConverterFactory messageConverterFactory,
|
||||
StreamFunctionProperties functionProperties,
|
||||
@@ -85,7 +85,7 @@ public class IntegrationFlowFunctionSupport {
|
||||
*/
|
||||
public <T> boolean containsFunction(Class<T> typeOfFunction) {
|
||||
return StringUtils.hasText(this.functionProperties.getDefinition())
|
||||
&& this.functionCatalog.contains(typeOfFunction,
|
||||
&& this.catalogContains(typeOfFunction,
|
||||
this.functionProperties.getDefinition());
|
||||
}
|
||||
|
||||
@@ -99,7 +99,7 @@ public class IntegrationFlowFunctionSupport {
|
||||
*/
|
||||
public <T> boolean containsFunction(Class<T> typeOfFunction, String functionName) {
|
||||
return StringUtils.hasText(functionName)
|
||||
&& this.functionCatalog.contains(typeOfFunction, functionName);
|
||||
&& this.catalogContains(typeOfFunction, functionName);
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -231,6 +231,10 @@ public class IntegrationFlowFunctionSupport {
|
||||
return true;
|
||||
}
|
||||
|
||||
private <T> boolean catalogContains(Class<T> functionType, String name) {
|
||||
return this.functionCatalog.lookup(functionType, name) != null;
|
||||
}
|
||||
|
||||
private <O> Mono<Void> subscribeToOutput(Consumer<Message<O>> outputProcessor,
|
||||
Publisher<Message<O>> outputPublisher) {
|
||||
|
||||
|
||||
@@ -125,7 +125,7 @@ public class FunctionInvokerTests {
|
||||
functionProperties.setDefinition("messageToMessageSameType");
|
||||
FunctionInvoker<Foo, Foo> messageToMessageSameType = new FunctionInvoker<>(
|
||||
functionProperties,
|
||||
new FunctionCatalogWrapper(context.getBean(FunctionCatalog.class)),
|
||||
context.getBean(FunctionCatalog.class),
|
||||
context.getBean(FunctionInspector.class),
|
||||
context.getBean(CompositeMessageConverterFactory.class));
|
||||
Message<Foo> outputMessage = messageToMessageSameType
|
||||
@@ -135,7 +135,7 @@ public class FunctionInvokerTests {
|
||||
functionProperties.setDefinition("pojoToPojoSameType");
|
||||
FunctionInvoker<Foo, Foo> pojoToPojoSameType = new FunctionInvoker<>(
|
||||
functionProperties,
|
||||
new FunctionCatalogWrapper(context.getBean(FunctionCatalog.class)),
|
||||
context.getBean(FunctionCatalog.class),
|
||||
context.getBean(FunctionInspector.class),
|
||||
context.getBean(CompositeMessageConverterFactory.class));
|
||||
outputMessage = pojoToPojoSameType.apply(Flux.just(inputMessage))
|
||||
@@ -145,7 +145,7 @@ public class FunctionInvokerTests {
|
||||
functionProperties.setDefinition("messageToMessageNoType");
|
||||
FunctionInvoker<Foo, Foo> messageToMessageNoType = new FunctionInvoker<>(
|
||||
functionProperties,
|
||||
new FunctionCatalogWrapper(context.getBean(FunctionCatalog.class)),
|
||||
context.getBean(FunctionCatalog.class),
|
||||
context.getBean(FunctionInspector.class),
|
||||
context.getBean(CompositeMessageConverterFactory.class));
|
||||
outputMessage = messageToMessageNoType.apply(Flux.just(inputMessage))
|
||||
@@ -155,7 +155,7 @@ public class FunctionInvokerTests {
|
||||
functionProperties.setDefinition("withException");
|
||||
FunctionInvoker<Foo, Foo> withException = new FunctionInvoker<>(
|
||||
functionProperties,
|
||||
new FunctionCatalogWrapper(context.getBean(FunctionCatalog.class)),
|
||||
context.getBean(FunctionCatalog.class),
|
||||
context.getBean(FunctionInspector.class),
|
||||
context.getBean(CompositeMessageConverterFactory.class));
|
||||
|
||||
@@ -181,7 +181,7 @@ public class FunctionInvokerTests {
|
||||
functionProperties.setDefinition("pojoToPojoNonEmptyPojo");
|
||||
FunctionInvoker<Baz, Baz> pojoToPojoSameType = new FunctionInvoker<>(
|
||||
functionProperties,
|
||||
new FunctionCatalogWrapper(context.getBean(FunctionCatalog.class)),
|
||||
context.getBean(FunctionCatalog.class),
|
||||
context.getBean(FunctionInspector.class),
|
||||
context.getBean(CompositeMessageConverterFactory.class));
|
||||
Message<Baz> outputMessage = pojoToPojoSameType.apply(Flux.just(inputMessage))
|
||||
@@ -193,7 +193,7 @@ public class FunctionInvokerTests {
|
||||
functionProperties.setDefinition("messageToMessageNoType");
|
||||
FunctionInvoker<Baz, Baz> messageToMessageNoType = new FunctionInvoker<>(
|
||||
functionProperties,
|
||||
new FunctionCatalogWrapper(context.getBean(FunctionCatalog.class)),
|
||||
context.getBean(FunctionCatalog.class),
|
||||
context.getBean(FunctionInspector.class),
|
||||
context.getBean(CompositeMessageConverterFactory.class));
|
||||
outputMessage = messageToMessageNoType.apply(Flux.just(inputMessageWithBaz))
|
||||
@@ -203,7 +203,7 @@ public class FunctionInvokerTests {
|
||||
functionProperties.setDefinition("withExceptionNativeEncodingEnabled");
|
||||
FunctionInvoker<Baz, Baz> withException = new FunctionInvoker<>(
|
||||
functionProperties,
|
||||
new FunctionCatalogWrapper(context.getBean(FunctionCatalog.class)),
|
||||
context.getBean(FunctionCatalog.class),
|
||||
context.getBean(FunctionInspector.class),
|
||||
context.getBean(CompositeMessageConverterFactory.class));
|
||||
|
||||
@@ -229,7 +229,7 @@ public class FunctionInvokerTests {
|
||||
functionProperties.setDefinition("pojoToPojoNonEmptyPojo");
|
||||
FunctionInvoker<Baz, Baz> pojoToPojoSameType = new FunctionInvoker<>(
|
||||
functionProperties,
|
||||
new FunctionCatalogWrapper(context.getBean(FunctionCatalog.class)),
|
||||
context.getBean(FunctionCatalog.class),
|
||||
context.getBean(FunctionInspector.class),
|
||||
context.getBean(CompositeMessageConverterFactory.class));
|
||||
Message<Baz> outputMessage = pojoToPojoSameType.apply(Flux.just(inputMessage))
|
||||
@@ -256,7 +256,7 @@ public class FunctionInvokerTests {
|
||||
functionProperties.setDefinition("fluxConsumer");
|
||||
FunctionInvoker<String, Void> fluxedConsumer = new FunctionInvoker<>(
|
||||
functionProperties,
|
||||
new FunctionCatalogWrapper(context.getBean(FunctionCatalog.class)),
|
||||
context.getBean(FunctionCatalog.class),
|
||||
context.getBean(FunctionInspector.class),
|
||||
context.getBean(CompositeMessageConverterFactory.class));
|
||||
|
||||
|
||||
Reference in New Issue
Block a user