From b692bffd6245cf4137542dfdec0ebb5437d2853f Mon Sep 17 00:00:00 2001 From: Oleg Zhurakousky Date: Wed, 2 Feb 2022 17:30:45 +0100 Subject: [PATCH] Additional cleanup and simplification --- .../cloud/stream/binder/AbstractBinder.java | 2 +- .../stream/binder/DefaultBinderFactory.java | 2 +- .../function/FunctionConfiguration.java | 19 ++++++++++--------- .../PartitionAwareFunctionWrapper.java | 2 +- .../ImplicitFunctionBindingTests.java | 1 - 5 files changed, 13 insertions(+), 13 deletions(-) diff --git a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/binder/AbstractBinder.java b/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/binder/AbstractBinder.java index 369ffc9db..55a7f0cdc 100644 --- a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/binder/AbstractBinder.java +++ b/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/binder/AbstractBinder.java @@ -136,7 +136,7 @@ public abstract class AbstractBinder bindConsumer(String name, String group, T target, C properties) { - if (StringUtils.isEmpty(group)) { + if (!StringUtils.hasText(group)) { Assert.isTrue(!properties.isPartitioned(), "A consumer group is required for a partitioned subscription"); } diff --git a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/binder/DefaultBinderFactory.java b/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/binder/DefaultBinderFactory.java index 652d4185f..a9c9c8645 100644 --- a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/binder/DefaultBinderFactory.java +++ b/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/binder/DefaultBinderFactory.java @@ -167,7 +167,7 @@ public class DefaultBinderFactory implements BinderFactory, DisposableBean, Appl } String configurationName; // Fall back to a default if no argument is provided - if (StringUtils.isEmpty(name)) { + if (!StringUtils.hasText(name)) { Assert.notEmpty(this.binderConfigurations, "A default binder has been requested, but there is no binder available"); if (!StringUtils.hasText(this.defaultBinder)) { 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 9c9c2b91a..57171e140 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 @@ -118,15 +118,13 @@ import org.springframework.util.StringUtils; */ @Configuration(proxyBeanMethods = false) @EnableConfigurationProperties(StreamFunctionProperties.class) -//@Import({ BindingBeansRegistrar.class, BinderFactoryAutoConfiguration.class }) @Import({ BinderFactoryAutoConfiguration.class }) @AutoConfigureBefore(BindingServiceConfiguration.class) @AutoConfigureAfter(ContextFunctionCatalogAutoConfiguration.class) @ConditionalOnBean(FunctionRegistry.class) public class FunctionConfiguration { - private final static String SOURCE_PROPERY = "spring.cloud.stream.source"; - + @SuppressWarnings("rawtypes") @Bean public StreamBridge streamBridgeUtils(FunctionCatalog functionCatalog, FunctionRegistry functionRegistry, BindingServiceProperties bindingServiceProperties, ConfigurableApplicationContext applicationContext, @@ -229,9 +227,9 @@ public class FunctionConfiguration { PollableBean pollable = extractPollableAnnotation(functionProperties, context, proxyFactory); if (functionWrapper != null) { - Type functionType = functionWrapper.getFunctionType(); +// Type functionType = functionWrapper.getFunctionType(); IntegrationFlow integrationFlow = integrationFlowFromProvidedSupplier(new PartitionAwareFunctionWrapper(functionWrapper, context, producerProperties), - beginPublishingTrigger, pollable, context, taskScheduler, functionType, producerProperties, outputName) + beginPublishingTrigger, pollable, context, taskScheduler, producerProperties, outputName) .route(Message.class, message -> { if (message.getHeaders().get("spring.cloud.stream.sendto.destination") != null) { String destinationName = (String) message.getHeaders().get("spring.cloud.stream.sendto.destination"); @@ -244,9 +242,9 @@ public class FunctionConfiguration { context.registerBean(integrationFlowName, IntegrationFlow.class, () -> postProcessedFlow); } else { - Type functionType = ((FunctionInvocationWrapper) supplier).getFunctionType(); + //Type functionType = ((FunctionInvocationWrapper) supplier).getFunctionType(); IntegrationFlow integrationFlow = integrationFlowFromProvidedSupplier(new PartitionAwareFunctionWrapper(supplier, context, producerProperties), - beginPublishingTrigger, pollable, context, taskScheduler, functionType, producerProperties, outputName) + beginPublishingTrigger, pollable, context, taskScheduler, producerProperties, outputName) .channel(c -> c.direct()) .fluxTransform((Function>, ? extends Publisher>) function) .route(Message.class, message -> { @@ -290,13 +288,16 @@ public class FunctionConfiguration { @SuppressWarnings({ "rawtypes", "unchecked" }) private IntegrationFlowBuilder integrationFlowFromProvidedSupplier(Supplier supplier, Publisher beginPublishingTrigger, PollableBean pollable, GenericApplicationContext context, - TaskScheduler taskScheduler, Type functionType, ProducerProperties producerProperties, String bindingName) { + TaskScheduler taskScheduler, ProducerProperties producerProperties, String bindingName) { IntegrationFlowBuilder integrationFlowBuilder; boolean splittable = pollable != null && (boolean) AnnotationUtils.getAnnotationAttributes(pollable).get("splittable"); - boolean reactive = FunctionTypeUtils.isPublisher(FunctionTypeUtils.getOutputType(functionType)); + + FunctionInvocationWrapper function = (supplier instanceof PartitionAwareFunctionWrapper) + ? (FunctionInvocationWrapper) ((PartitionAwareFunctionWrapper) supplier).function : (FunctionInvocationWrapper) supplier; + boolean reactive = FunctionTypeUtils.isPublisher(function.getOutputType()); if (pollable == null && reactive) { Publisher publisher = (Publisher) supplier.get(); diff --git a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/function/PartitionAwareFunctionWrapper.java b/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/function/PartitionAwareFunctionWrapper.java index ef3cc2fde..e6cfbe5f6 100644 --- a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/function/PartitionAwareFunctionWrapper.java +++ b/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/function/PartitionAwareFunctionWrapper.java @@ -43,7 +43,7 @@ class PartitionAwareFunctionWrapper implements Function, Supplie protected final Log logger = LogFactory.getLog(PartitionAwareFunctionWrapper.class); @SuppressWarnings("rawtypes") - private final Function function; + protected final Function function; private final Function outputMessageEnricher; diff --git a/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/function/ImplicitFunctionBindingTests.java b/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/function/ImplicitFunctionBindingTests.java index 4e1e5f015..fd141486a 100644 --- a/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/function/ImplicitFunctionBindingTests.java +++ b/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/function/ImplicitFunctionBindingTests.java @@ -87,7 +87,6 @@ public class ImplicitFunctionBindingTests { } - @Test public void testFailedApplicationListenerConfiguration() { try (ConfigurableApplicationContext context = new SpringApplicationBuilder(