GH-1740 Added initial support for Supplier of Publisher

Resolves #1740
This commit is contained in:
Oleg Zhurakousky
2019-06-20 17:38:49 +02:00
parent 25a6ae7114
commit 1b7345b331
4 changed files with 82 additions and 6 deletions

View File

@@ -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());
}
}

View File

@@ -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<MonoSink<Object>> triggerRef = new AtomicReference<>();
private final Publisher<Object> 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);
}

View File

@@ -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))

View File

@@ -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<Publisher<Object>> 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 {