interim
This commit is contained in:
@@ -79,7 +79,7 @@ import org.springframework.util.Assert;
|
||||
*/
|
||||
@Configuration
|
||||
@EnableConfigurationProperties({ BindingServiceProperties.class, SpringIntegrationProperties.class, FunctionProperties.class })
|
||||
@Import({ DestinationPublishingMetricsAutoConfiguration.class, SpelExpressionConverterConfiguration.class, FunctionConfiguration.class })
|
||||
@Import({ DestinationPublishingMetricsAutoConfiguration.class, SpelExpressionConverterConfiguration.class })
|
||||
@Role(BeanDefinition.ROLE_INFRASTRUCTURE)
|
||||
@ConditionalOnBean(value = BinderTypeRegistry.class, search = SearchStrategy.CURRENT)
|
||||
public class BindingServiceConfiguration {
|
||||
|
||||
@@ -16,12 +16,25 @@
|
||||
|
||||
package org.springframework.cloud.stream.function;
|
||||
|
||||
import org.springframework.beans.factory.annotation.Autowired;
|
||||
import org.springframework.beans.factory.config.ConfigurableListableBeanFactory;
|
||||
import org.springframework.beans.factory.support.BeanDefinitionRegistry;
|
||||
import org.springframework.boot.autoconfigure.condition.ConditionalOnMissingBean;
|
||||
import org.springframework.boot.autoconfigure.condition.ConditionalOnProperty;
|
||||
import org.springframework.cloud.function.context.FunctionCatalog;
|
||||
import org.springframework.cloud.function.context.FunctionType;
|
||||
import org.springframework.cloud.function.context.catalog.FunctionInspector;
|
||||
import org.springframework.cloud.stream.binding.BindingBeanDefinitionRegistryUtils;
|
||||
import org.springframework.cloud.stream.converter.CompositeMessageConverterFactory;
|
||||
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.annotation.Bean;
|
||||
import org.springframework.context.annotation.Configuration;
|
||||
import org.springframework.integration.dsl.IntegrationFlow;
|
||||
import org.springframework.messaging.MessageChannel;
|
||||
import org.springframework.messaging.SubscribableChannel;
|
||||
import org.springframework.util.ClassUtils;
|
||||
|
||||
/**
|
||||
*
|
||||
@@ -33,6 +46,18 @@ import org.springframework.context.annotation.Configuration;
|
||||
@ConditionalOnProperty("spring.cloud.stream.function.name")
|
||||
public class FunctionConfiguration {
|
||||
|
||||
@Autowired(required=false)
|
||||
private Source source;
|
||||
|
||||
@Autowired(required=false)
|
||||
private Processor processor;
|
||||
|
||||
@Autowired(required=false)
|
||||
private Sink sink;
|
||||
|
||||
@Autowired
|
||||
private ConfigurableListableBeanFactory registry;
|
||||
|
||||
@Bean
|
||||
public IntegrationFlowFunctionSupport functionSupport(FunctionCatalogWrapper functionCatalog,
|
||||
FunctionInspector functionInspector, CompositeMessageConverterFactory messageConverterFactory,
|
||||
@@ -47,4 +72,28 @@ public class FunctionConfiguration {
|
||||
return new FunctionCatalogWrapper(catalog);
|
||||
}
|
||||
|
||||
|
||||
@ConditionalOnProperty("spring.cloud.stream.function.name")
|
||||
@ConditionalOnMissingBean
|
||||
@Bean
|
||||
public IntegrationFlow foo(IntegrationFlowFunctionSupport functionSupport) {
|
||||
if (processor != null) {
|
||||
return functionSupport.integrationFlowForFunction(processor.input(), processor.output()).get();
|
||||
}
|
||||
else if (sink != null) {
|
||||
return functionSupport.integrationFlowForFunction(sink.input(), null).get();
|
||||
}
|
||||
else if (source != null) {
|
||||
return functionSupport.integrationFlowFromNamedSupplier().channel(this.source.output()).get();
|
||||
}
|
||||
|
||||
FunctionType ft = functionSupport.getCurrentFunctionType();
|
||||
BindingBeanDefinitionRegistryUtils.registerBindingTargetBeanDefinitions(Sink.class,
|
||||
Sink.class.getName(), (BeanDefinitionRegistry) registry);
|
||||
BindingBeanDefinitionRegistryUtils.registerBindingTargetsQualifiedBeanDefinitions(
|
||||
ClassUtils.resolveClassName(this.getClass().getName(), null), Sink.class,
|
||||
(BeanDefinitionRegistry) registry);
|
||||
return functionSupport.integrationFlowForFunction(registry.getBean("input", SubscribableChannel.class), null).get();
|
||||
//throw new UnsupportedOperationException("Not yet supotrted");
|
||||
}
|
||||
}
|
||||
|
||||
@@ -26,6 +26,7 @@ import reactor.core.publisher.Flux;
|
||||
import reactor.core.publisher.Mono;
|
||||
|
||||
import org.springframework.beans.factory.annotation.Autowired;
|
||||
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.converter.CompositeMessageConverterFactory;
|
||||
@@ -76,6 +77,10 @@ public class IntegrationFlowFunctionSupport {
|
||||
this.functionProperties = functionProperties;
|
||||
}
|
||||
|
||||
public FunctionType getCurrentFunctionType() {
|
||||
return functionInspector.getRegistration(functionCatalog.lookup(this.functionProperties.getName())).getType();
|
||||
}
|
||||
|
||||
/**
|
||||
* Create an instance of the {@link IntegrationFlowBuilder} from a {@link Supplier} bean available in the context.
|
||||
* The name of the bean must be provided via `spring.cloud.stream.function.name` property.
|
||||
@@ -113,6 +118,21 @@ public class IntegrationFlowFunctionSupport {
|
||||
return flowBuilder;
|
||||
}
|
||||
|
||||
/**
|
||||
*
|
||||
* @param inputChannel
|
||||
* @param outputChannel
|
||||
* @return
|
||||
*/
|
||||
public <O> IntegrationFlowBuilder integrationFlowForFunction(SubscribableChannel inputChannel, MessageChannel outputChannel) {
|
||||
IntegrationFlowBuilder flowBuilder = IntegrationFlows.from(inputChannel).bridge();
|
||||
|
||||
if (!this.andThenFunction(flowBuilder, outputChannel)) {
|
||||
flowBuilder = flowBuilder.channel(outputChannel);
|
||||
}
|
||||
return flowBuilder;
|
||||
}
|
||||
|
||||
/**
|
||||
* Add a {@link Function} bean to the end of an integration flow.
|
||||
* The name of the bean must be provided via `spring.cloud.stream.function.name` property.
|
||||
@@ -131,7 +151,12 @@ public class IntegrationFlowFunctionSupport {
|
||||
new FunctionInvoker<>(this.functionProperties.getName(), this.functionCatalog,
|
||||
this.functionInspector, this.messageConverterFactory, this.errorChannel);
|
||||
|
||||
subscribeToInput(functionInvoker, flowBuilder.toReactivePublisher(), outputChannel::send);
|
||||
if (outputChannel != null) {
|
||||
subscribeToInput(functionInvoker, flowBuilder.toReactivePublisher(), outputChannel::send);
|
||||
}
|
||||
else {
|
||||
subscribeToInput(functionInvoker, flowBuilder.toReactivePublisher(), null);
|
||||
}
|
||||
return true;
|
||||
}
|
||||
return false;
|
||||
|
||||
@@ -3,6 +3,8 @@ org.springframework.cloud.stream.config.ChannelBindingAutoConfiguration,\
|
||||
org.springframework.cloud.stream.config.BindersHealthIndicatorAutoConfiguration,\
|
||||
org.springframework.cloud.stream.config.ChannelsEndpointAutoConfiguration,\
|
||||
org.springframework.cloud.stream.config.BindingsEndpointAutoConfiguration,\
|
||||
org.springframework.cloud.stream.config.BindingServiceConfiguration
|
||||
org.springframework.cloud.stream.config.BindingServiceConfiguration,\
|
||||
org.springframework.cloud.stream.function.FunctionConfiguration
|
||||
|
||||
|
||||
|
||||
|
||||
@@ -0,0 +1,132 @@
|
||||
package org.springframework.cloud.stream.function;
|
||||
|
||||
import java.util.Date;
|
||||
import java.util.function.Consumer;
|
||||
import java.util.function.Function;
|
||||
import java.util.function.Supplier;
|
||||
|
||||
import org.junit.Test;
|
||||
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.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.Processor;
|
||||
import org.springframework.cloud.stream.messaging.Sink;
|
||||
import org.springframework.cloud.stream.messaging.Source;
|
||||
import org.springframework.context.ConfigurableApplicationContext;
|
||||
import org.springframework.context.annotation.Bean;
|
||||
import org.springframework.messaging.Message;
|
||||
import org.springframework.messaging.support.GenericMessage;
|
||||
|
||||
public class NewSourceAsSupplierTests {
|
||||
|
||||
@Test
|
||||
public void testSourceFromSupplier() {
|
||||
try (ConfigurableApplicationContext context = new SpringApplicationBuilder(
|
||||
TestChannelBinderConfiguration.getCompleteConfiguration(SourceFromSupplier.class)).web(
|
||||
WebApplicationType.NONE).run("--spring.cloud.stream.function.name=date", "--spring.jmx.enabled=false")) {
|
||||
|
||||
OutputDestination target = context.getBean(OutputDestination.class);
|
||||
Message<byte[]> sourceMessage = target.receive(10000);
|
||||
System.out.println(sourceMessage);
|
||||
// assertThat(target.receive(10000).getPayload()).isEqualTo("1".getBytes(StandardCharsets.UTF_8));
|
||||
// assertThat(target.receive(10000).getPayload()).isEqualTo("2".getBytes(StandardCharsets.UTF_8));
|
||||
// assertThat(target.receive(10000).getPayload()).isEqualTo("3".getBytes(StandardCharsets.UTF_8));
|
||||
//etc
|
||||
}
|
||||
}
|
||||
|
||||
@Test
|
||||
public void testProcessorFromFunction() {
|
||||
try (ConfigurableApplicationContext context = new SpringApplicationBuilder(
|
||||
TestChannelBinderConfiguration.getCompleteConfiguration(ProcessorFromFunction.class)).web(
|
||||
WebApplicationType.NONE).run("--spring.cloud.stream.function.name=toUpperCase", "--spring.jmx.enabled=false")) {
|
||||
|
||||
InputDestination source = context.getBean(InputDestination.class);
|
||||
source.send(new GenericMessage<byte[]>("fopo".getBytes()));
|
||||
OutputDestination target = context.getBean(OutputDestination.class);
|
||||
Message<byte[]> targetMessage = target.receive(10000);
|
||||
System.out.println(new String(targetMessage.getPayload()));
|
||||
// assertThat(target.receive(10000).getPayload()).isEqualTo("1".getBytes(StandardCharsets.UTF_8));
|
||||
// assertThat(target.receive(10000).getPayload()).isEqualTo("2".getBytes(StandardCharsets.UTF_8));
|
||||
// assertThat(target.receive(10000).getPayload()).isEqualTo("3".getBytes(StandardCharsets.UTF_8));
|
||||
//etc
|
||||
}
|
||||
}
|
||||
|
||||
@Test
|
||||
public void testSinkFromConsumer() {
|
||||
try (ConfigurableApplicationContext context = new SpringApplicationBuilder(
|
||||
TestChannelBinderConfiguration.getCompleteConfiguration(SinkFromConsumer.class)).web(
|
||||
WebApplicationType.NONE).run("--spring.cloud.stream.function.name=sink", "--spring.jmx.enabled=false")) {
|
||||
|
||||
InputDestination source = context.getBean(InputDestination.class);
|
||||
source.send(new GenericMessage<byte[]>("fopo".getBytes()));
|
||||
// OutputDestination target = context.getBean(OutputDestination.class);
|
||||
// Message<byte[]> targetMessage = target.receive(10000);
|
||||
// System.out.println(new String(targetMessage.getPayload()));
|
||||
// assertThat(target.receive(10000).getPayload()).isEqualTo("1".getBytes(StandardCharsets.UTF_8));
|
||||
// assertThat(target.receive(10000).getPayload()).isEqualTo("2".getBytes(StandardCharsets.UTF_8));
|
||||
// assertThat(target.receive(10000).getPayload()).isEqualTo("3".getBytes(StandardCharsets.UTF_8));
|
||||
//etc
|
||||
}
|
||||
}
|
||||
|
||||
@Test
|
||||
public void testSinkFromConsumerNoEnableBinding() {
|
||||
try (ConfigurableApplicationContext context = new SpringApplicationBuilder(
|
||||
TestChannelBinderConfiguration.getCompleteConfiguration(SinkFromConsumerNoEnableBinding.class)).web(
|
||||
WebApplicationType.NONE).run("--spring.cloud.stream.function.name=sink", "--spring.jmx.enabled=false")) {
|
||||
|
||||
InputDestination source = context.getBean(InputDestination.class);
|
||||
source.send(new GenericMessage<byte[]>("Hello No Binding".getBytes()));
|
||||
// OutputDestination target = context.getBean(OutputDestination.class);
|
||||
// Message<byte[]> targetMessage = target.receive(10000);
|
||||
// System.out.println(new String(targetMessage.getPayload()));
|
||||
// assertThat(target.receive(10000).getPayload()).isEqualTo("1".getBytes(StandardCharsets.UTF_8));
|
||||
// assertThat(target.receive(10000).getPayload()).isEqualTo("2".getBytes(StandardCharsets.UTF_8));
|
||||
// assertThat(target.receive(10000).getPayload()).isEqualTo("3".getBytes(StandardCharsets.UTF_8));
|
||||
//etc
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@EnableAutoConfiguration
|
||||
@EnableBinding(Source.class)
|
||||
public static class SourceFromSupplier {
|
||||
@Bean
|
||||
public Supplier<Date> date() {
|
||||
return () -> new Date();
|
||||
}
|
||||
}
|
||||
|
||||
@EnableAutoConfiguration
|
||||
@EnableBinding(Processor.class)
|
||||
public static class ProcessorFromFunction {
|
||||
@Bean
|
||||
public Function<String, String> toUpperCase() {
|
||||
return s -> s.toUpperCase();
|
||||
}
|
||||
}
|
||||
|
||||
@EnableAutoConfiguration
|
||||
@EnableBinding(Sink.class)
|
||||
public static class SinkFromConsumer {
|
||||
@Bean
|
||||
public Consumer<String> sink() {
|
||||
return s -> System.out.println(s);
|
||||
}
|
||||
}
|
||||
|
||||
@EnableAutoConfiguration
|
||||
// @EnableBinding(Sink.class)
|
||||
public static class SinkFromConsumerNoEnableBinding {
|
||||
@Bean
|
||||
public Consumer<String> sink() {
|
||||
return s -> System.out.println("==> " + s);
|
||||
}
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user