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 0a78376c5..c5222fd08 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 @@ -35,6 +35,7 @@ import org.springframework.cloud.stream.messaging.Source; import org.springframework.context.annotation.Bean; import org.springframework.context.annotation.Configuration; import org.springframework.context.annotation.Import; +import org.springframework.context.support.GenericApplicationContext; import org.springframework.integration.channel.NullChannel; import org.springframework.integration.dsl.IntegrationFlow; import org.springframework.lang.Nullable; @@ -58,7 +59,8 @@ public class FunctionConfiguration { FunctionCatalog functionCatalog, FunctionInspector functionInspector, CompositeMessageConverterFactory messageConverterFactory, StreamFunctionProperties functionProperties, - BindingServiceProperties bindingServiceProperties) { + BindingServiceProperties bindingServiceProperties, + GenericApplicationContext context) { ((SmartInitializingSingleton) functionCatalog).afterSingletonsInstantiated(); // if (functionCatalog.size() > 0) { @@ -70,7 +72,7 @@ public class FunctionConfiguration { // } return new IntegrationFlowFunctionSupport(functionCatalog, functionInspector, - messageConverterFactory, functionProperties, bindingServiceProperties); + messageConverterFactory, functionProperties, bindingServiceProperties, context); } /** @@ -122,5 +124,4 @@ public class FunctionConfiguration { return processor != null ? processor.output() : (source != null ? source.output() : new NullChannel()); } - } diff --git a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/function/IntegrationFlowFunctionSupport.java b/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/function/IntegrationFlowFunctionSupport.java index 2ff7a8dca..339df5b0d 100644 --- a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/function/IntegrationFlowFunctionSupport.java +++ b/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/function/IntegrationFlowFunctionSupport.java @@ -16,6 +16,7 @@ package org.springframework.cloud.stream.function; +import java.util.concurrent.atomic.AtomicReference; import java.util.function.Consumer; import java.util.function.Function; import java.util.function.Supplier; @@ -23,14 +24,18 @@ import java.util.function.Supplier; import org.reactivestreams.Publisher; import reactor.core.publisher.Flux; import reactor.core.publisher.Mono; +import reactor.core.publisher.MonoSink; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.cloud.function.context.FunctionCatalog; +import org.springframework.cloud.function.context.FunctionRegistration; import org.springframework.cloud.function.context.FunctionType; import org.springframework.cloud.function.context.catalog.FunctionInspector; import org.springframework.cloud.function.core.FluxSupplier; +import org.springframework.cloud.stream.binder.BindingCreatedEvent; import org.springframework.cloud.stream.config.BindingServiceProperties; import org.springframework.cloud.stream.converter.CompositeMessageConverterFactory; +import org.springframework.context.support.GenericApplicationContext; import org.springframework.integration.context.IntegrationObjectSupport; import org.springframework.integration.dsl.IntegrationFlowBuilder; import org.springframework.integration.dsl.IntegrationFlows; @@ -46,7 +51,7 @@ import org.springframework.util.StringUtils; * @author Ilayaperumal Gopinathan * @since 2.1 */ -public class IntegrationFlowFunctionSupport { +public class IntegrationFlowFunctionSupport { private final FunctionCatalog functionCatalog; @@ -56,6 +61,10 @@ public class IntegrationFlowFunctionSupport { private final StreamFunctionProperties functionProperties; + private final AtomicReference> triggerRef = new AtomicReference<>(); + + private final Publisher trigger; + @Autowired private MessageChannel errorChannel; @@ -63,7 +72,8 @@ public class IntegrationFlowFunctionSupport { FunctionInspector functionInspector, CompositeMessageConverterFactory messageConverterFactory, StreamFunctionProperties functionProperties, - BindingServiceProperties bindingServiceProperties) { + BindingServiceProperties bindingServiceProperties, + GenericApplicationContext context) { Assert.notNull(functionCatalog, "'functionCatalog' must not be null"); Assert.notNull(functionInspector, "'functionInspector' must not be null"); Assert.notNull(messageConverterFactory, @@ -74,8 +84,19 @@ public class IntegrationFlowFunctionSupport { this.messageConverterFactory = messageConverterFactory; this.functionProperties = functionProperties; this.functionProperties.setBindingServiceProperties(bindingServiceProperties); + trigger = Mono.create(emmiter -> { + triggerRef.set(emmiter); + }); + context.addApplicationListener(event -> { + if (event instanceof BindingCreatedEvent) { + if (triggerRef.get() != null) { + triggerRef.get().success(); + } + } + }); } + /** * Determines if function specified via 'spring.cloud.stream.function.definition' * property can be located in {@link FunctionCatalog}s. @@ -139,8 +160,19 @@ public class IntegrationFlowFunctionSupport { * @param supplier supplier from which the flow builder will be built * @return instance of {@link IntegrationFlowBuilder} */ + @SuppressWarnings({ "rawtypes", "unchecked" }) public IntegrationFlowBuilder integrationFlowFromProvidedSupplier( Supplier supplier) { + String supplierName = this.functionProperties.getDefinition().split("\\|")[0]; + FunctionRegistration fr = this.functionInspector.getRegistration(this.functionCatalog.lookup(supplierName)); + if (fr != null && fr.getType().isWrapper()) { + Publisher publisher = (Publisher) supplier.get(); + publisher = publisher instanceof Flux + ? ((Flux) publisher).delaySubscription(trigger) + : ((Mono) publisher).delaySubscription(trigger); + + return IntegrationFlows.from(publisher); + } return IntegrationFlows.from(supplier); } 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 38508bd54..118b596f7 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 @@ -93,7 +93,7 @@ public class ImplicitFunctionBindingTests { @Test public void testBindingWithNoEnableBindingAndNoDefinitionPropertyConfiguration() { - + System.clearProperty("spring.cloud.stream.function.definition"); try (ConfigurableApplicationContext context = new SpringApplicationBuilder( TestChannelBinderConfiguration.getCompleteConfiguration( SingleFunctionConfiguration.class)) 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 b968e4a30..415208b71 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 @@ -17,13 +17,17 @@ package org.springframework.cloud.stream.function; import java.nio.charset.StandardCharsets; +import java.util.concurrent.ExecutorService; +import java.util.concurrent.Executors; import java.util.concurrent.atomic.AtomicInteger; import java.util.function.Function; import java.util.function.Supplier; + import org.junit.Rule; import org.junit.Test; import org.junit.rules.ExpectedException; +import org.reactivestreams.Publisher; import reactor.core.publisher.Flux; import org.springframework.beans.factory.BeanCreationException; @@ -176,6 +180,28 @@ public class SourceToFunctionsSupportTests { } } + @Test + public void testMessageSourceIsCreatedFromProvidedStreamSupplier() { + try (ConfigurableApplicationContext context = new SpringApplicationBuilder( + TestChannelBinderConfiguration.getCompleteConfiguration( + StreamSupplierConfiguration.class)).web(WebApplicationType.NONE).run( + "--spring.cloud.stream.function.definition=stream", + "--spring.jmx.enabled=false")) { + + OutputDestination target = context.getBean(OutputDestination.class); + assertThat(target.receive(1000).getPayload()) + .isEqualTo("0".getBytes(StandardCharsets.UTF_8)); + assertThat(target.receive(1000).getPayload()) + .isEqualTo("1".getBytes(StandardCharsets.UTF_8)); + assertThat(target.receive(1000).getPayload()) + .isEqualTo("2".getBytes(StandardCharsets.UTF_8)); + assertThat(target.receive(1000).getPayload()) + .isEqualTo("3".getBytes(StandardCharsets.UTF_8)); + + // etc + } + } + @Test public void testFunctionDoesNotExist() { @@ -227,6 +253,23 @@ public class SourceToFunctionsSupportTests { } + @EnableAutoConfiguration + public static class StreamSupplierConfiguration { + + @Bean + public Supplier> stream() { + ExecutorService executor = Executors.newFixedThreadPool(1); + return () -> Flux.create(emitter -> { + executor.execute(() -> { + for (int i = 0; i < 10; i++) { + emitter.next(MessageBuilder.withPayload(String.valueOf(i)).build()); + } + }); + }); + } + + } + @EnableAutoConfiguration @Import(ExistingMessageSourceConfigurationNoContentTypeSet.class) public static class FunctionsConfigurationNoConversionPossible {