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 62227815b..5a8b08b81 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 @@ -41,6 +41,7 @@ import org.springframework.cloud.function.context.PollableSupplier; import org.springframework.cloud.function.context.catalog.BeanFactoryAwareFunctionRegistry.FunctionInvocationWrapper; import org.springframework.cloud.function.context.catalog.FunctionInspector; import org.springframework.cloud.function.context.catalog.FunctionTypeUtils; +import org.springframework.cloud.stream.annotation.EnableBinding; import org.springframework.cloud.stream.binder.BindingCreatedEvent; import org.springframework.cloud.stream.binding.BindableProxyFactory; import org.springframework.cloud.stream.config.BinderFactoryAutoConfiguration; @@ -82,39 +83,42 @@ public class FunctionConfiguration { @Bean public InitializingBean functionChannelBindingInitializer(FunctionCatalog functionCatalog, FunctionInspector functionInspector, - StreamFunctionProperties functionProperties, @Nullable BindableProxyFactory[] bindableProxyFactory) { + StreamFunctionProperties functionProperties, @Nullable BindableProxyFactory[] bindableProxyFactory, GenericApplicationContext context) { return new FunctionChannelBindingInitializer(functionCatalog, functionInspector, functionProperties, - ObjectUtils.isEmpty(bindableProxyFactory) ? null : bindableProxyFactory[0]); + ObjectUtils.isEmpty(bindableProxyFactory) ? null : bindableProxyFactory[0]); } @Bean public IntegrationFlow standAloneSupplierFlow(FunctionCatalog functionCatalog, FunctionInspector functionInspector, StreamFunctionProperties functionProperties, GenericApplicationContext context) { + IntegrationFlow integrationFlow = null; - FunctionInvocationWrapper functionWrapper = functionCatalog.lookup(functionProperties.getDefinition()); - if (functionWrapper != null) { - AtomicReference> triggerRef = new AtomicReference<>(); - Publisher beginPublishingTrigger = Mono.create(emmiter -> { - triggerRef.set(emmiter); - }); - context.addApplicationListener(event -> { - if (event instanceof BindingCreatedEvent) { - if (triggerRef.get() != null) { - triggerRef.get().success(); + if (functionCatalog != null && ObjectUtils.isEmpty(context.getBeanNamesForAnnotation(EnableBinding.class))) { + FunctionInvocationWrapper functionWrapper = functionCatalog.lookup(functionProperties.getDefinition()); + if (functionWrapper != null) { + AtomicReference> triggerRef = new AtomicReference<>(); + Publisher beginPublishingTrigger = Mono.create(emmiter -> { + triggerRef.set(emmiter); + }); + context.addApplicationListener(event -> { + if (event instanceof BindingCreatedEvent) { + if (triggerRef.get() != null) { + triggerRef.get().success(); + } } + }); + + + RootBeanDefinition bd = (RootBeanDefinition) context.getBeanDefinition(functionProperties.getParsedDefinition()[0]); + Method factoryMethod = bd.getResolvedFactoryMethod(); + PollableSupplier pollable = factoryMethod.getReturnType().isAssignableFrom(Supplier.class) + ? AnnotationUtils.findAnnotation(factoryMethod, PollableSupplier.class) + : null; + + if (!functionProperties.isComposeFrom() && !functionProperties.isComposeTo() && functionWrapper.isSupplier()) { + integrationFlow = this.integrationFlowFromProvidedSupplier(functionWrapper, functionInspector, beginPublishingTrigger, pollable) + .channel("output").get(); } - }); - - - RootBeanDefinition bd = (RootBeanDefinition) context.getBeanDefinition(functionProperties.getParsedDefinition()[0]); - Method factoryMethod = bd.getResolvedFactoryMethod(); - PollableSupplier pollable = factoryMethod.getReturnType().isAssignableFrom(Supplier.class) - ? AnnotationUtils.findAnnotation(factoryMethod, PollableSupplier.class) - : null; - - if (!functionProperties.isComposeFrom() && !functionProperties.isComposeTo() && functionWrapper.isSupplier()) { - integrationFlow = this.integrationFlowFromProvidedSupplier(functionWrapper, functionInspector, beginPublishingTrigger, pollable) - .channel("output").get(); } } diff --git a/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/function/GreenfieldFunctionEnableBindingTests.java b/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/function/GreenfieldFunctionEnableBindingTests.java index 12ab236fd..65fef8b90 100644 --- a/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/function/GreenfieldFunctionEnableBindingTests.java +++ b/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/function/GreenfieldFunctionEnableBindingTests.java @@ -166,7 +166,6 @@ public class GreenfieldFunctionEnableBindingTests { } @EnableAutoConfiguration - @EnableBinding(Source.class) public static class SourceFromSupplier { @Bean 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 917cb7e92..5bd9b4328 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 @@ -26,9 +26,12 @@ import reactor.core.publisher.Flux; import org.springframework.boot.WebApplicationType; import org.springframework.boot.autoconfigure.EnableAutoConfiguration; import org.springframework.boot.builder.SpringApplicationBuilder; +import org.springframework.cloud.stream.annotation.EnableBinding; +import org.springframework.cloud.stream.annotation.StreamListener; import org.springframework.cloud.stream.binder.test.InputDestination; import org.springframework.cloud.stream.binder.test.OutputDestination; import org.springframework.cloud.stream.binder.test.TestChannelBinderConfiguration; +import org.springframework.cloud.stream.messaging.Sink; import org.springframework.context.ConfigurableApplicationContext; import org.springframework.context.annotation.Bean; import org.springframework.integration.support.MessageBuilder; @@ -178,7 +181,19 @@ public class ImplicitFunctionBindingTests { assertThat(outputMessage.getPayload()).isEqualTo("Hello".getBytes()); outputMessage = outputDestination.receive(); assertThat(outputMessage.getPayload()).isEqualTo("Hello Again".getBytes()); + } + } + @Test + public void testFunctionConfigDisabledIfStreamListenerIsUsed() { + System.clearProperty("spring.cloud.stream.function.definition"); + try (ConfigurableApplicationContext context = new SpringApplicationBuilder( + TestChannelBinderConfiguration.getCompleteConfiguration( + LegacyConfiguration.class)) + .web(WebApplicationType.NONE) + .run("--spring.jmx.enabled=false")) { + + assertThat(context.getBean("standAloneSupplierFlow")).isEqualTo(null); } } @@ -237,4 +252,14 @@ public class ImplicitFunctionBindingTests { } } + @EnableAutoConfiguration + @EnableBinding(Sink.class) + public static class LegacyConfiguration { + + @StreamListener(Sink.INPUT) + public void handle(String value) { + + } + } + } diff --git a/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/function/SourceToFunctionsSupportTests.java b/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/function/SourceToFunctionsSupportTests.java index 44d88a70a..1fffb2d71 100644 --- a/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/function/SourceToFunctionsSupportTests.java +++ b/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/function/SourceToFunctionsSupportTests.java @@ -232,6 +232,9 @@ public class SourceToFunctionsSupportTests { assertThat(new String(target.receive(2000).getPayload())).isEqualTo("4"); assertThat(new String(target.receive(2000).getPayload())).isEqualTo("5"); assertThat(new String(target.receive(2000).getPayload())).isEqualTo("6"); + + assertThat(context.getBean("standAloneSupplierFlow")).isNotEqualTo(null); + assertThat(context.getBean("functionChannelBindingInitializer")).isNotEqualTo(null); } }