diff --git a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/annotation/BindingProvider.java b/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/annotation/BindingProvider.java new file mode 100644 index 000000000..15dca35aa --- /dev/null +++ b/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/annotation/BindingProvider.java @@ -0,0 +1,48 @@ +/* + * Copyright 2019-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 + * + * https://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.annotation; + +import java.lang.annotation.Documented; +import java.lang.annotation.ElementType; +import java.lang.annotation.Inherited; +import java.lang.annotation.Retention; +import java.lang.annotation.RetentionPolicy; +import java.lang.annotation.Target; + +import org.springframework.integration.config.EnableIntegration; + +/** + * Marker annotation which signals to the framework that the + * actual binder will handle the actual bindings and no action from the core + * framework is necessary to establish bindings. + * + * This annotation must be present on the configuration class of the actual binder. + * + * @author Oleg Zhurakousky + * + * @since 3.0.0 + * + */ +@Target({ ElementType.TYPE}) +@Retention(RetentionPolicy.RUNTIME) +@Documented +@Inherited +@EnableIntegration +public @interface BindingProvider { + + +} diff --git a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/config/BindableProvider.java b/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/config/BindableProvider.java deleted file mode 100644 index 0b7e3d52f..000000000 --- a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/config/BindableProvider.java +++ /dev/null @@ -1,37 +0,0 @@ -/* - * Copyright 2019-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 - * - * https://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.config; - -/** - * If downstream binders want to take over the responsibility of binding the target types (such as the Kafka Streams binder), - * then they can implement this functional interface to signal the core framework to bypass any binding. - * - * @author Soby Chacko - * @since 3.0.0 - */ -@FunctionalInterface -public interface BindableProvider { - - /** - * Based on the type provided, the implementation can determine whether it is capable of - * binding this target type. - * - * @param clazz target type to bind - * @return true if capable of binding - */ - boolean canBind(Class clazz); -} 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 ccd2e7880..e17749d2f 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 @@ -26,6 +26,7 @@ import java.util.LinkedList; import java.util.List; import java.util.Map; import java.util.Properties; +import java.util.stream.Stream; import org.apache.commons.logging.Log; import org.apache.commons.logging.LogFactory; @@ -44,6 +45,7 @@ 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; import org.springframework.cloud.stream.annotation.EnableBinding; import org.springframework.cloud.stream.binder.BinderType; import org.springframework.cloud.stream.binder.BinderTypeRegistry; @@ -63,6 +65,7 @@ import org.springframework.context.annotation.Bean; import org.springframework.context.annotation.Configuration; import org.springframework.context.annotation.Import; import org.springframework.context.annotation.Role; +import org.springframework.core.annotation.AnnotationUtils; import org.springframework.core.env.Environment; import org.springframework.core.env.StandardEnvironment; import org.springframework.core.io.Resource; @@ -248,8 +251,12 @@ public class BinderFactoryAutoConfiguration { } @Bean - public BeanFactoryPostProcessor implicitFunctionBinder(Environment environment, + 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 { @@ -261,25 +268,7 @@ public class BinderFactoryAutoConfiguration { Class inputType = inspector.getInputType(definedFunction); Class outputType = inspector.getOutputType(definedFunction); - boolean bindDownstream = false; - try { - Map bindableProviders = beanFactory.getBeansOfType(BindableProvider.class); - Class inputTypeWrapper = inspector.getInputWrapper(definedFunction); - if (bindableProviders != null) { - Collection values = bindableProviders.values(); - for (BindableProvider bindableProvider : values) { - if (bindableProvider.canBind(inputTypeWrapper)) { - bindDownstream = true; - break; - } - } - } - } - catch (BeansException be) { - // pass through - } - - if (!bindDownstream) { + if (!bindingProvider) { if (Void.class.isAssignableFrom(outputType)) { bind(Sink.class, registry); }