From 35c0e1e9f06d626cf9a66ed371ca1ee7f253474b Mon Sep 17 00:00:00 2001 From: Oleg Zhurakousky Date: Thu, 1 Aug 2019 12:17:17 +0200 Subject: [PATCH] GH-1775 Fix early initialisation of FunctionCatalog Resolves #1775 --- .../BinderFactoryAutoConfiguration.java | 147 +++++++++++------- 1 file changed, 90 insertions(+), 57 deletions(-) diff --git a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/config/BinderFactoryAutoConfiguration.java b/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/config/BinderFactoryAutoConfiguration.java index 8f3636396..4b3a98973 100644 --- a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/config/BinderFactoryAutoConfiguration.java +++ b/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/config/BinderFactoryAutoConfiguration.java @@ -33,15 +33,16 @@ import org.apache.commons.logging.LogFactory; import org.springframework.beans.BeansException; import org.springframework.beans.factory.BeanCreationException; +import org.springframework.beans.factory.BeanFactory; +import org.springframework.beans.factory.BeanFactoryAware; +import org.springframework.beans.factory.InitializingBean; import org.springframework.beans.factory.annotation.Value; import org.springframework.beans.factory.config.BeanDefinition; -import org.springframework.beans.factory.config.BeanFactoryPostProcessor; import org.springframework.beans.factory.config.ConfigurableListableBeanFactory; import org.springframework.beans.factory.support.BeanDefinitionRegistry; import org.springframework.beans.factory.support.RootBeanDefinition; import org.springframework.boot.context.properties.EnableConfigurationProperties; import org.springframework.cloud.function.context.FunctionCatalog; -import org.springframework.cloud.function.context.FunctionRegistry; import org.springframework.cloud.function.context.catalog.FunctionInspector; import org.springframework.cloud.function.context.config.RoutingFunction; import org.springframework.cloud.stream.annotation.BindingProvider; @@ -60,6 +61,7 @@ import org.springframework.cloud.stream.messaging.Processor; import org.springframework.cloud.stream.messaging.Sink; import org.springframework.cloud.stream.messaging.Source; import org.springframework.context.ConfigurableApplicationContext; +import org.springframework.context.EnvironmentAware; import org.springframework.context.annotation.Bean; import org.springframework.context.annotation.Configuration; import org.springframework.context.annotation.Import; @@ -252,67 +254,98 @@ public class BinderFactoryAutoConfiguration { } @Bean - public BeanFactoryPostProcessor implicitFunctionBinder(BinderTypeRegistry bfac, Environment environment, - @Nullable FunctionRegistry functionCatalog, @Nullable FunctionInspector inspector) { - Class[] configurationClasses = bfac.getAll().values().iterator().next().getConfigurationClasses(); - boolean bindingProvider = Stream.of(configurationClasses) - .filter(clazz -> AnnotationUtils.findAnnotation(clazz, BindingProvider.class) != null) - .findFirst().isPresent(); - return new BeanFactoryPostProcessor() { - @Override - public void postProcessBeanFactory(ConfigurableListableBeanFactory beanFactory) throws BeansException { - if (functionCatalog != null && ObjectUtils.isEmpty(beanFactory.getBeanNamesForAnnotation(EnableBinding.class))) { - BeanDefinitionRegistry registry = (BeanDefinitionRegistry) beanFactory; - String name = determineFunctionName(functionCatalog, environment); - if (StringUtils.hasText(name)) { - Object definedFunction = functionCatalog.lookup(name); - Class inputType = inspector.getInputType(definedFunction); - Class outputType = inspector.getOutputType(definedFunction); + public InitializingBean functionToChannelBindingInitializer(@Nullable FunctionCatalog functionCatalog, + @Nullable FunctionInspector functionInspector, BinderTypeRegistry bfac) { + return new ImplicitFunctionToChannelBindingInitializer(functionCatalog, functionInspector, bfac); + } - if (!bindingProvider) { - if (Void.class.isAssignableFrom(outputType)) { - bind(Sink.class, registry); - } - else if (Void.class.isAssignableFrom(inputType)) { - bind(Source.class, registry); - } - else { - bind(Processor.class, registry); - } + + private static class ImplicitFunctionToChannelBindingInitializer implements InitializingBean, BeanFactoryAware, EnvironmentAware { + + private ConfigurableListableBeanFactory beanFactory; + + private Environment environment; + + private final FunctionCatalog functionCatalog; + + private final FunctionInspector functionInspector; + + private final BinderTypeRegistry bfac; + + ImplicitFunctionToChannelBindingInitializer(FunctionCatalog functionCatalog, + FunctionInspector functionInspector, BinderTypeRegistry bfac) { + this.functionCatalog = functionCatalog; + this.functionInspector = functionInspector; + this.bfac = bfac; + } + @Override + public void afterPropertiesSet() { + Class[] configurationClasses = bfac.getAll().values().iterator().next().getConfigurationClasses(); + boolean bindingProvider = Stream.of(configurationClasses) + .filter(clazz -> AnnotationUtils.findAnnotation(clazz, BindingProvider.class) != null) + .findFirst().isPresent(); + if (functionCatalog != null && ObjectUtils.isEmpty(beanFactory.getBeanNamesForAnnotation(EnableBinding.class))) { + BeanDefinitionRegistry registry = (BeanDefinitionRegistry) beanFactory; + String name = determineFunctionName(functionCatalog, environment); + if (StringUtils.hasText(name)) { + Object definedFunction = functionCatalog.lookup(name); + Class inputType = functionInspector.getInputType(definedFunction); + Class outputType = functionInspector.getOutputType(definedFunction); + + if (!bindingProvider) { + if (Void.class.isAssignableFrom(outputType)) { + bind(Sink.class, registry); + } + else if (Void.class.isAssignableFrom(inputType)) { + bind(Source.class, registry); + } + else { + bind(Processor.class, registry); } } } } - }; + } + + @Override + public void setBeanFactory(BeanFactory beanFactory) throws BeansException { + this.beanFactory = (ConfigurableListableBeanFactory) beanFactory; + } + + @Override + public void setEnvironment(Environment environment) { + this.environment = environment; + } + + private String determineFunctionName(FunctionCatalog catalog, Environment environment) { + String name = environment.getProperty("spring.cloud.stream.function.definition"); + if (!StringUtils.hasText(name)) { + name = environment.getProperty("spring.cloud.function.definition"); + } + if (!StringUtils.hasText(name) && Boolean.parseBoolean( + environment.getProperty("spring.cloud.function.routing.enabled", "false"))) { + name = RoutingFunction.FUNCTION_NAME; + } + if (!StringUtils.hasText(name) && catalog.size() >= 1 && catalog.size() <= 2) { + name = ((FunctionInspector) catalog).getName(catalog.lookup("")); + } + if (StringUtils.hasText(name)) { + ((StandardEnvironment) environment).getSystemProperties() + .putIfAbsent("spring.cloud.stream.function.definition", name); + } + return name; + } + + private void bind(Class type, BeanDefinitionRegistry registry) { + if (!registry.containsBeanDefinition(type.getName())) { + RootBeanDefinition rootBeanDefinition = new RootBeanDefinition( + BindableProxyFactory.class); + rootBeanDefinition.getConstructorArgumentValues() + .addGenericArgumentValue(type); + registry.registerBeanDefinition(type.getName(), rootBeanDefinition); + } + } } - private String determineFunctionName(FunctionCatalog catalog, Environment environment) { - String name = environment.getProperty("spring.cloud.stream.function.definition"); - if (!StringUtils.hasText(name)) { - name = environment.getProperty("spring.cloud.function.definition"); - } - if (!StringUtils.hasText(name) && Boolean.parseBoolean( - environment.getProperty("spring.cloud.function.routing.enabled", "false"))) { - name = RoutingFunction.FUNCTION_NAME; - } - if (!StringUtils.hasText(name) && catalog.size() >= 1 && catalog.size() <= 2) { - name = ((FunctionInspector) catalog).getName(catalog.lookup("")); - } - if (StringUtils.hasText(name)) { - ((StandardEnvironment) environment).getSystemProperties() - .putIfAbsent("spring.cloud.stream.function.definition", name); - } - return name; - } - - private void bind(Class type, BeanDefinitionRegistry registry) { - if (!registry.containsBeanDefinition(type.getName())) { - RootBeanDefinition rootBeanDefinition = new RootBeanDefinition( - BindableProxyFactory.class); - rootBeanDefinition.getConstructorArgumentValues() - .addGenericArgumentValue(type); - registry.registerBeanDefinition(type.getName(), rootBeanDefinition); - } - } }