From 6e333693184b2cb992ccdc2f4d72bb96c04e1939 Mon Sep 17 00:00:00 2001 From: Dave Syer Date: Thu, 9 Jul 2015 14:02:20 +0100 Subject: [PATCH] Add aggregate builder and sample applications No need to specify channel names (just spring.cloud.streams.name). Example app: @SpringBootApplication @EnableChannelBinding public class ExtendedApplication implements AggregateConfigurer { @Override public void configure(AggregateBuilder builder) { builder .from(TimeSource.class).as("source") .via(LoggingTransformer.class) .via(LoggingTransformer.class).profiles("other") .to(LogSink.class); } } Fixes gh-2 --- .../cloud/streams/EnableChannelBinding.java | 6 +- .../streams/aggregate/AggregateBuilder.java | 257 ++++++++++++++++++ .../aggregate/AggregateConfigurer.java | 27 ++ .../config/AggregateBuilderConfiguration.java | 50 ++++ .../ChannelBindingAdapterConfiguration.java | 9 +- .../config/RabbitServiceConfiguration.java | 2 + .../config/RedisServiceConfiguration.java | 2 + ...oduleOptionsPropertySourceInitializer.java | 3 + vanilla-samples/double/pom.xml | 69 +++++ .../java/config/SinkModuleDefinition.java} | 4 +- .../java/config/SourceModuleDefinition.java | 55 ++++ .../src/main/java/demo/DoubleApplication.java | 26 ++ .../double/src/main/resources/application.yml | 1 + .../double/src/main/resources/sink.yml | 5 + .../double/src/main/resources/source.yml | 6 + .../java/demo/ModuleApplicationTests.java | 20 ++ vanilla-samples/extended/pom.xml | 81 ++++++ .../java/extended/ExtendedApplication.java | 32 +++ .../src/main/resources/application.yml | 6 + .../extended/src/main/resources/source.yml | 1 + .../java/demo/ModuleApplicationTests.java | 22 ++ vanilla-samples/pom.xml | 22 ++ vanilla-samples/sink/pom.xml | 2 +- .../src/main/java/demo/SinkApplication.java | 4 +- .../sink/src/main/java/sink/LogSink.java | 50 ++++ vanilla-samples/source/pom.xml | 2 +- .../src/main/java/demo/SourceApplication.java | 4 +- .../TimeSource.java} | 4 +- .../TimeSourceOptionsMetadata.java | 2 +- .../source/src/main/resources/application.yml | 14 +- vanilla-samples/transform/pom.xml | 69 +++++ .../main/java/demo/TransformApplication.java | 17 ++ .../java/transform/LoggingTransformer.java | 72 +++++ .../src/main/resources/application.yml | 8 + .../java/demo/ModuleApplicationTests.java | 20 ++ 35 files changed, 954 insertions(+), 20 deletions(-) create mode 100644 spring-cloud-streams/src/main/java/org/springframework/cloud/streams/aggregate/AggregateBuilder.java create mode 100644 spring-cloud-streams/src/main/java/org/springframework/cloud/streams/aggregate/AggregateConfigurer.java create mode 100644 spring-cloud-streams/src/main/java/org/springframework/cloud/streams/config/AggregateBuilderConfiguration.java create mode 100644 vanilla-samples/double/pom.xml rename vanilla-samples/{sink/src/main/java/config/ModuleDefinition.java => double/src/main/java/config/SinkModuleDefinition.java} (92%) create mode 100644 vanilla-samples/double/src/main/java/config/SourceModuleDefinition.java create mode 100644 vanilla-samples/double/src/main/java/demo/DoubleApplication.java create mode 100644 vanilla-samples/double/src/main/resources/application.yml create mode 100644 vanilla-samples/double/src/main/resources/sink.yml create mode 100644 vanilla-samples/double/src/main/resources/source.yml create mode 100644 vanilla-samples/double/src/test/java/demo/ModuleApplicationTests.java create mode 100644 vanilla-samples/extended/pom.xml create mode 100644 vanilla-samples/extended/src/main/java/extended/ExtendedApplication.java create mode 100644 vanilla-samples/extended/src/main/resources/application.yml create mode 100644 vanilla-samples/extended/src/main/resources/source.yml create mode 100644 vanilla-samples/extended/src/test/java/demo/ModuleApplicationTests.java create mode 100644 vanilla-samples/sink/src/main/java/sink/LogSink.java rename vanilla-samples/source/src/main/java/{config/ModuleDefinition.java => source/TimeSource.java} (97%) rename vanilla-samples/source/src/main/java/{config => source}/TimeSourceOptionsMetadata.java (99%) create mode 100644 vanilla-samples/transform/pom.xml create mode 100644 vanilla-samples/transform/src/main/java/demo/TransformApplication.java create mode 100644 vanilla-samples/transform/src/main/java/transform/LoggingTransformer.java create mode 100644 vanilla-samples/transform/src/main/resources/application.yml create mode 100644 vanilla-samples/transform/src/test/java/demo/ModuleApplicationTests.java diff --git a/spring-cloud-streams/src/main/java/org/springframework/cloud/streams/EnableChannelBinding.java b/spring-cloud-streams/src/main/java/org/springframework/cloud/streams/EnableChannelBinding.java index 7c88e27bf..8fb515ac4 100644 --- a/spring-cloud-streams/src/main/java/org/springframework/cloud/streams/EnableChannelBinding.java +++ b/spring-cloud-streams/src/main/java/org/springframework/cloud/streams/EnableChannelBinding.java @@ -23,8 +23,9 @@ import java.lang.annotation.Retention; import java.lang.annotation.RetentionPolicy; import java.lang.annotation.Target; -import org.springframework.cloud.streams.config.LifecycleConfiguration; +import org.springframework.cloud.streams.config.AggregateBuilderConfiguration; import org.springframework.cloud.streams.config.ChannelBindingAdapterConfiguration; +import org.springframework.cloud.streams.config.LifecycleConfiguration; import org.springframework.cloud.streams.config.RabbitServiceConfiguration; import org.springframework.cloud.streams.config.RedisServiceConfiguration; import org.springframework.context.annotation.Configuration; @@ -40,7 +41,8 @@ import org.springframework.context.annotation.Import; @Inherited @Configuration @Import({ RedisServiceConfiguration.class, RabbitServiceConfiguration.class, - ChannelBindingAdapterConfiguration.class, LifecycleConfiguration.class }) + ChannelBindingAdapterConfiguration.class, LifecycleConfiguration.class, + AggregateBuilderConfiguration.class }) public @interface EnableChannelBinding { } diff --git a/spring-cloud-streams/src/main/java/org/springframework/cloud/streams/aggregate/AggregateBuilder.java b/spring-cloud-streams/src/main/java/org/springframework/cloud/streams/aggregate/AggregateBuilder.java new file mode 100644 index 000000000..5a4c6995f --- /dev/null +++ b/spring-cloud-streams/src/main/java/org/springframework/cloud/streams/aggregate/AggregateBuilder.java @@ -0,0 +1,257 @@ +/* + * Copyright 2015 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 + * + * http://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.streams.aggregate; + +import java.util.ArrayList; +import java.util.HashSet; +import java.util.List; + +import org.springframework.beans.BeansException; +import org.springframework.boot.autoconfigure.EnableAutoConfiguration; +import org.springframework.boot.builder.SpringApplicationBuilder; +import org.springframework.boot.context.properties.ConfigurationProperties; +import org.springframework.context.ApplicationContext; +import org.springframework.context.ApplicationContextAware; +import org.springframework.context.ApplicationContextInitializer; +import org.springframework.context.ConfigurableApplicationContext; +import org.springframework.context.annotation.Configuration; +import org.springframework.util.StringUtils; + +/** + * @author Dave Syer + * + */ +@ConfigurationProperties("spring.cloud.streams") +public class AggregateBuilder implements ApplicationContextAware { + + public static final String DEFAULT_NAME = "application"; + + private ConfigurableApplicationContext parent; + + private int index = 0; + + private String streamName = "stream"; + + private HashSet sinks = new HashSet(); + + public String getName() { + return this.streamName; + } + + public void setName(String name) { + this.streamName = name; + } + + @Override + public void setApplicationContext(ApplicationContext applicationContext) + throws BeansException { + this.parent = (ConfigurableApplicationContext) applicationContext; + } + + public void build() { + for (SinkConfigurer sink : this.sinks) { + sink.build(); + } + } + + public SourceConfigurer from(Class module) { + return new SourceConfigurer(module); + } + + private String channelName() { + return this.streamName + "." + this.index; + } + + private String incrementChannelName() { + return this.streamName + "." + (this.index++); + } + + public class SourceConfigurer { + + private Class module; + private String[] names = null; + private String[] profiles = null; + + public SourceConfigurer(Class module) { + this.module = module; + } + + public SourceConfigurer as(String... names) { + this.names = names; + return this; + } + + public SourceConfigurer profiles(String... profiles) { + this.profiles = profiles; + return this; + } + + public SinkConfigurer to(Class sink) { + build(); + return new SinkConfigurer(sink); + } + + public ProcessorConfigurer via(Class processor) { + build(); + return new ProcessorConfigurer(processor); + } + + private void build() { + childContext(this.module).config(this.names).profiles(this.profiles) + .output(channelName()).build(); + } + + } + + public class SinkConfigurer { + + private Class module; + private String[] names = null; + private String[] profiles = null; + + public SinkConfigurer profiles(String... profiles) { + this.profiles = profiles; + return this; + } + + public SinkConfigurer(Class module) { + AggregateBuilder.this.sinks.add(this); + this.module = module; + } + + public SinkConfigurer as(String... names) { + this.names = names; + return this; + } + + void build() { + childContext(this.module).config(this.names).profiles(this.profiles) + .input(channelName()).build(); + } + + } + + public class ProcessorConfigurer { + + private Class module; + private String[] names = null; + private String[] profiles = null; + + public ProcessorConfigurer(Class module) { + this.module = module; + } + + public ProcessorConfigurer as(String... names) { + this.names = names; + return this; + } + + public ProcessorConfigurer profiles(String... profiles) { + this.profiles = profiles; + return this; + } + + public SinkConfigurer to(Class sink) { + build(); + return new SinkConfigurer(sink); + } + + public ProcessorConfigurer via(Class processor) { + build(); + return new ProcessorConfigurer(processor); + } + + private void build() { + childContext(this.module).config(this.names).profiles(this.profiles) + .input(incrementChannelName()).output(channelName()).build(); + } + + } + + private ChildContextBuilder childContext(Class type) { + return new ChildContextBuilder(new SpringApplicationBuilder(type, + SeedConfiguration.class).parent(AggregateBuilder.this.parent) + .showBanner(false).web(false) + .initializers(new BeanPostProcessorInitializer())); + } + + private class ChildContextBuilder { + + private SpringApplicationBuilder builder; + private String configName; + private String input; + private String output; + + public ChildContextBuilder(SpringApplicationBuilder builder) { + this.builder = builder; + } + + public ChildContextBuilder profiles(String... profiles) { + if (profiles != null) { + this.builder.profiles(profiles); + } + return this; + } + + public ChildContextBuilder config(String... configs) { + if (configs != null) { + this.configName = StringUtils.arrayToCommaDelimitedString(configs); + } + return this; + } + + public ChildContextBuilder input(String input) { + this.input = input; + return this; + } + + public ChildContextBuilder output(String output) { + this.output = output; + return this; + } + + public void build() { + List args = new ArrayList(); + if (this.configName != null) { + args.add("--spring.config.name=" + this.configName); + } + if (this.input != null) { + args.add("--spring.cloud.channels.inputChannelName=" + this.input); + } + if (this.output != null) { + args.add("--spring.cloud.channels.outputChannelName=" + this.output); + } + this.builder.run(args.toArray(new String[0])); + } + + } + + private class BeanPostProcessorInitializer implements + ApplicationContextInitializer { + + @Override + public void initialize(ConfigurableApplicationContext applicationContext) { + } + + } + + @Configuration + @EnableAutoConfiguration + protected static class SeedConfiguration { + } + +} diff --git a/spring-cloud-streams/src/main/java/org/springframework/cloud/streams/aggregate/AggregateConfigurer.java b/spring-cloud-streams/src/main/java/org/springframework/cloud/streams/aggregate/AggregateConfigurer.java new file mode 100644 index 000000000..735f57e4c --- /dev/null +++ b/spring-cloud-streams/src/main/java/org/springframework/cloud/streams/aggregate/AggregateConfigurer.java @@ -0,0 +1,27 @@ +/* + * Copyright 2015 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 + * + * http://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.streams.aggregate; + +/** + * @author Dave Syer + * + */ +public interface AggregateConfigurer { + + void configure(AggregateBuilder builder); + +} diff --git a/spring-cloud-streams/src/main/java/org/springframework/cloud/streams/config/AggregateBuilderConfiguration.java b/spring-cloud-streams/src/main/java/org/springframework/cloud/streams/config/AggregateBuilderConfiguration.java new file mode 100644 index 000000000..1d3717bf7 --- /dev/null +++ b/spring-cloud-streams/src/main/java/org/springframework/cloud/streams/config/AggregateBuilderConfiguration.java @@ -0,0 +1,50 @@ +/* + * Copyright 2015 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 + * + * http://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.streams.config; + +import org.springframework.beans.factory.ListableBeanFactory; +import org.springframework.beans.factory.annotation.Autowired; +import org.springframework.boot.CommandLineRunner; +import org.springframework.cloud.streams.aggregate.AggregateBuilder; +import org.springframework.cloud.streams.aggregate.AggregateConfigurer; +import org.springframework.context.annotation.Bean; +import org.springframework.context.annotation.Configuration; + +/** + * @author Dave Syer + * + */ +@Configuration +public class AggregateBuilderConfiguration implements CommandLineRunner { + + @Autowired + private ListableBeanFactory beanFactory; + + @Bean + public AggregateBuilder aggregateBuilder() { + return new AggregateBuilder(); + } + + @Override + public void run(String... args) throws Exception { + for (AggregateConfigurer configurer : this.beanFactory.getBeansOfType(AggregateConfigurer.class).values()) { + configurer.configure(aggregateBuilder()); + } + aggregateBuilder().build(); + } + +} diff --git a/spring-cloud-streams/src/main/java/org/springframework/cloud/streams/config/ChannelBindingAdapterConfiguration.java b/spring-cloud-streams/src/main/java/org/springframework/cloud/streams/config/ChannelBindingAdapterConfiguration.java index 3f4713606..e8d817cf9 100644 --- a/spring-cloud-streams/src/main/java/org/springframework/cloud/streams/config/ChannelBindingAdapterConfiguration.java +++ b/spring-cloud-streams/src/main/java/org/springframework/cloud/streams/config/ChannelBindingAdapterConfiguration.java @@ -32,6 +32,7 @@ import org.springframework.beans.factory.BeanFactoryUtils; import org.springframework.beans.factory.ListableBeanFactory; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.boot.autoconfigure.condition.ConditionalOnMissingBean; +import org.springframework.boot.autoconfigure.condition.SearchStrategy; import org.springframework.cloud.streams.adapter.ChannelBindingAdapter; import org.springframework.cloud.streams.adapter.ChannelLocator; import org.springframework.cloud.streams.adapter.Input; @@ -104,8 +105,7 @@ public class ChannelBindingAdapterConfiguration { protected Collection getOutputChannels() { Set channels = new LinkedHashSet(); - String[] names = BeanFactoryUtils.beanNamesForTypeIncludingAncestors( - this.beanFactory, MessageChannel.class); + String[] names = this.beanFactory.getBeanNamesForType(MessageChannel.class); for (String name : names) { if (name.startsWith("output")) { channels.add(new OutputChannelBinding(name)); @@ -116,8 +116,7 @@ public class ChannelBindingAdapterConfiguration { protected Collection getInputChannels() { Set channels = new LinkedHashSet(); - String[] names = BeanFactoryUtils.beanNamesForTypeIncludingAncestors( - this.beanFactory, MessageChannel.class); + String[] names = this.beanFactory.getBeanNamesForType(MessageChannel.class); for (String name : names) { if (name.startsWith("input")) { channels.add(new InputChannelBinding(name)); @@ -174,7 +173,7 @@ public class ChannelBindingAdapterConfiguration { } @Configuration - @ConditionalOnMissingBean(ChannelBindingProperties.class) + @ConditionalOnMissingBean(value=ChannelBindingProperties.class, search=SearchStrategy.CURRENT) protected static class ModulePropertiesConfiguration { @Bean(name = "spring.cloud.channels.CONFIGURATION_PROPERTIES") public ChannelBindingProperties moduleProperties() { diff --git a/spring-cloud-streams/src/main/java/org/springframework/cloud/streams/config/RabbitServiceConfiguration.java b/spring-cloud-streams/src/main/java/org/springframework/cloud/streams/config/RabbitServiceConfiguration.java index abf577b29..926e89cff 100644 --- a/spring-cloud-streams/src/main/java/org/springframework/cloud/streams/config/RabbitServiceConfiguration.java +++ b/spring-cloud-streams/src/main/java/org/springframework/cloud/streams/config/RabbitServiceConfiguration.java @@ -18,6 +18,7 @@ package org.springframework.cloud.streams.config; import org.springframework.amqp.rabbit.connection.ConnectionFactory; import org.springframework.boot.autoconfigure.condition.ConditionalOnClass; +import org.springframework.boot.autoconfigure.condition.ConditionalOnMissingBean; import org.springframework.cloud.Cloud; import org.springframework.cloud.CloudFactory; import org.springframework.context.annotation.Bean; @@ -35,6 +36,7 @@ import org.springframework.xd.dirt.integration.rabbit.RabbitMessageBus; */ @Configuration @ConditionalOnClass(RabbitMessageBus.class) +@ConditionalOnMissingBean(RabbitMessageBus.class) @ImportResource({ "classpath*:/META-INF/spring-xd/bus/rabbit-bus.xml", "classpath*:/META-INF/spring-xd/analytics/rabbit-analytics.xml" }) @PropertySource("classpath:/META-INF/spring-cloud-streams/rabbit-bus.properties") diff --git a/spring-cloud-streams/src/main/java/org/springframework/cloud/streams/config/RedisServiceConfiguration.java b/spring-cloud-streams/src/main/java/org/springframework/cloud/streams/config/RedisServiceConfiguration.java index 3d94ab9db..26ce4d2b6 100644 --- a/spring-cloud-streams/src/main/java/org/springframework/cloud/streams/config/RedisServiceConfiguration.java +++ b/spring-cloud-streams/src/main/java/org/springframework/cloud/streams/config/RedisServiceConfiguration.java @@ -17,6 +17,7 @@ package org.springframework.cloud.streams.config; import org.springframework.boot.autoconfigure.condition.ConditionalOnClass; +import org.springframework.boot.autoconfigure.condition.ConditionalOnMissingBean; import org.springframework.cloud.Cloud; import org.springframework.cloud.CloudFactory; import org.springframework.context.annotation.Bean; @@ -35,6 +36,7 @@ import org.springframework.xd.dirt.integration.redis.RedisMessageBus; */ @Configuration @ConditionalOnClass(RedisMessageBus.class) +@ConditionalOnMissingBean(RedisMessageBus.class) @ImportResource({ "classpath*:/META-INF/spring-xd/bus/redis-bus.xml", "classpath*:/META-INF/spring-xd/analytics/redis-analytics.xml" }) @PropertySource("classpath:/META-INF/spring-cloud-streams/redis-bus.properties") diff --git a/spring-xd-runner/src/main/java/org/springframework/cloud/streams/xd/ModuleOptionsPropertySourceInitializer.java b/spring-xd-runner/src/main/java/org/springframework/cloud/streams/xd/ModuleOptionsPropertySourceInitializer.java index 93414a8bf..6b67eb23d 100644 --- a/spring-xd-runner/src/main/java/org/springframework/cloud/streams/xd/ModuleOptionsPropertySourceInitializer.java +++ b/spring-xd-runner/src/main/java/org/springframework/cloud/streams/xd/ModuleOptionsPropertySourceInitializer.java @@ -25,6 +25,8 @@ import java.util.Map; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.boot.autoconfigure.condition.ConditionalOnExpression; +import org.springframework.boot.autoconfigure.condition.ConditionalOnMissingBean; +import org.springframework.boot.autoconfigure.condition.SearchStrategy; import org.springframework.boot.context.properties.EnableConfigurationProperties; import org.springframework.context.ApplicationContextInitializer; import org.springframework.context.ConfigurableApplicationContext; @@ -129,6 +131,7 @@ ApplicationContextInitializer { } @Configuration + @ConditionalOnMissingBean(value=ModuleProperties.class, search=SearchStrategy.CURRENT) protected static class ModulePropertiesConfiguration { @Bean(name = "spring.cloud.channels.CONFIGURATION_PROPERTIES") public ModuleProperties moduleProperties() { diff --git a/vanilla-samples/double/pom.xml b/vanilla-samples/double/pom.xml new file mode 100644 index 000000000..d2e4dd95b --- /dev/null +++ b/vanilla-samples/double/pom.xml @@ -0,0 +1,69 @@ + + + 4.0.0 + + org.springframework.cloud + spring-cloud-streams-sample-double + 1.0.0.BUILD-SNAPSHOT + jar + + spring-cloud-streams-sample-double + Demo project for Spring XD module + + + org.springframework.cloud + spring-cloud-streams-samples + 1.0.0.BUILD-SNAPSHOT + + + + UTF-8 + demo.SinkApplication + 1.8 + + + + + org.springframework.cloud + spring-cloud-streams + + + org.springframework.xd + spring-xd-messagebus-redis + + + org.springframework.cloud + spring-cloud-lattice-connector + 1.0.2.BUILD-SNAPSHOT + + + org.springframework.boot + spring-boot-starter-redis + + + org.springframework.boot + spring-boot-configuration-processor + true + + + + org.springframework.boot + spring-boot-starter-test + test + + + + + + + org.springframework.boot + spring-boot-maven-plugin + + exec + + + + + + diff --git a/vanilla-samples/sink/src/main/java/config/ModuleDefinition.java b/vanilla-samples/double/src/main/java/config/SinkModuleDefinition.java similarity index 92% rename from vanilla-samples/sink/src/main/java/config/ModuleDefinition.java rename to vanilla-samples/double/src/main/java/config/SinkModuleDefinition.java index f7a9abd67..81d8fc0e1 100644 --- a/vanilla-samples/sink/src/main/java/config/ModuleDefinition.java +++ b/vanilla-samples/double/src/main/java/config/SinkModuleDefinition.java @@ -33,9 +33,9 @@ import org.springframework.messaging.MessageChannel; @Configuration @EnableChannelBinding @MessageEndpoint -public class ModuleDefinition { +public class SinkModuleDefinition { - private static Logger logger = LoggerFactory.getLogger(ModuleDefinition.class); + private static Logger logger = LoggerFactory.getLogger(SinkModuleDefinition.class); @Bean public MessageChannel input() { diff --git a/vanilla-samples/double/src/main/java/config/SourceModuleDefinition.java b/vanilla-samples/double/src/main/java/config/SourceModuleDefinition.java new file mode 100644 index 000000000..c111f748f --- /dev/null +++ b/vanilla-samples/double/src/main/java/config/SourceModuleDefinition.java @@ -0,0 +1,55 @@ +/* + * Copyright 2015 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 + * + * http://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 config; + +import java.text.SimpleDateFormat; +import java.util.Date; + +import org.springframework.beans.factory.annotation.Value; +import org.springframework.cloud.streams.EnableChannelBinding; +import org.springframework.context.annotation.Bean; +import org.springframework.context.annotation.Configuration; +import org.springframework.integration.annotation.InboundChannelAdapter; +import org.springframework.integration.annotation.Poller; +import org.springframework.integration.channel.DirectChannel; +import org.springframework.integration.core.MessageSource; +import org.springframework.messaging.MessageChannel; +import org.springframework.messaging.support.GenericMessage; + +/** + * @author Dave Syer + * + */ +@Configuration +@EnableChannelBinding +public class SourceModuleDefinition { + + @Value("${format:YYYY/MM/dd hh:mm:ss}") + private String format; + + @Bean + public MessageChannel output() { + return new DirectChannel(); + } + + @Bean + @InboundChannelAdapter(value = "output", autoStartup = "false", poller = @Poller(fixedDelay = "${fixedDelay}", maxMessagesPerPoll = "1")) + public MessageSource timerMessageSource() { + return () -> new GenericMessage<>(new SimpleDateFormat(this.format).format(new Date())); + } + +} diff --git a/vanilla-samples/double/src/main/java/demo/DoubleApplication.java b/vanilla-samples/double/src/main/java/demo/DoubleApplication.java new file mode 100644 index 000000000..6cd23443c --- /dev/null +++ b/vanilla-samples/double/src/main/java/demo/DoubleApplication.java @@ -0,0 +1,26 @@ +package demo; + +import org.springframework.boot.SpringApplication; +import org.springframework.boot.autoconfigure.SpringBootApplication; +import org.springframework.cloud.streams.EnableChannelBinding; +import org.springframework.cloud.streams.aggregate.AggregateBuilder; +import org.springframework.cloud.streams.aggregate.AggregateConfigurer; + +import config.SinkModuleDefinition; +import config.SourceModuleDefinition; + +@SpringBootApplication +@EnableChannelBinding +public class DoubleApplication implements AggregateConfigurer { + + @Override + public void configure(AggregateBuilder builder) { + builder.from(SourceModuleDefinition.class).as("source") + .to(SinkModuleDefinition.class).as("sink"); + } + + public static void main(String[] args) throws InterruptedException { + SpringApplication.run(DoubleApplication.class, args); + } + +} diff --git a/vanilla-samples/double/src/main/resources/application.yml b/vanilla-samples/double/src/main/resources/application.yml new file mode 100644 index 000000000..136d06384 --- /dev/null +++ b/vanilla-samples/double/src/main/resources/application.yml @@ -0,0 +1 @@ + \ No newline at end of file diff --git a/vanilla-samples/double/src/main/resources/sink.yml b/vanilla-samples/double/src/main/resources/sink.yml new file mode 100644 index 000000000..508af83e3 --- /dev/null +++ b/vanilla-samples/double/src/main/resources/sink.yml @@ -0,0 +1,5 @@ +spring: + cloud: + channels: + inputChannelName: testtock + \ No newline at end of file diff --git a/vanilla-samples/double/src/main/resources/source.yml b/vanilla-samples/double/src/main/resources/source.yml new file mode 100644 index 000000000..8d7308b20 --- /dev/null +++ b/vanilla-samples/double/src/main/resources/source.yml @@ -0,0 +1,6 @@ +fixedDelay: 5000 +spring: + cloud: + channels: + outputChannelName: testtock + \ No newline at end of file diff --git a/vanilla-samples/double/src/test/java/demo/ModuleApplicationTests.java b/vanilla-samples/double/src/test/java/demo/ModuleApplicationTests.java new file mode 100644 index 000000000..f62135c8f --- /dev/null +++ b/vanilla-samples/double/src/test/java/demo/ModuleApplicationTests.java @@ -0,0 +1,20 @@ +package demo; + +import org.junit.Test; +import org.junit.runner.RunWith; +import org.springframework.test.annotation.DirtiesContext; +import org.springframework.test.context.web.WebAppConfiguration; +import org.springframework.boot.test.SpringApplicationConfiguration; +import org.springframework.test.context.junit4.SpringJUnit4ClassRunner; + +@RunWith(SpringJUnit4ClassRunner.class) +@SpringApplicationConfiguration(classes = DoubleApplication.class) +@WebAppConfiguration +@DirtiesContext +public class ModuleApplicationTests { + + @Test + public void contextLoads() { + } + +} diff --git a/vanilla-samples/extended/pom.xml b/vanilla-samples/extended/pom.xml new file mode 100644 index 000000000..e5777d82e --- /dev/null +++ b/vanilla-samples/extended/pom.xml @@ -0,0 +1,81 @@ + + + 4.0.0 + + org.springframework.cloud + spring-cloud-streams-sample-extended + 1.0.0.BUILD-SNAPSHOT + jar + + spring-cloud-streams-sample-extended + Demo project for Spring XD module + + + org.springframework.cloud + spring-cloud-streams-samples + 1.0.0.BUILD-SNAPSHOT + + + + UTF-8 + demo.SinkApplication + 1.8 + + + + + org.springframework.cloud + spring-cloud-streams + + + org.springframework.cloud + spring-cloud-streams-sample-source + + + org.springframework.cloud + spring-cloud-streams-sample-transform + + + org.springframework.cloud + spring-cloud-streams-sample-sink + + + org.springframework.xd + spring-xd-messagebus-redis + + + org.springframework.cloud + spring-cloud-lattice-connector + 1.0.2.BUILD-SNAPSHOT + + + org.springframework.boot + spring-boot-starter-redis + + + org.springframework.boot + spring-boot-configuration-processor + true + + + + org.springframework.boot + spring-boot-starter-test + test + + + + + + + org.springframework.boot + spring-boot-maven-plugin + + exec + + + + + + diff --git a/vanilla-samples/extended/src/main/java/extended/ExtendedApplication.java b/vanilla-samples/extended/src/main/java/extended/ExtendedApplication.java new file mode 100644 index 000000000..5065d2b64 --- /dev/null +++ b/vanilla-samples/extended/src/main/java/extended/ExtendedApplication.java @@ -0,0 +1,32 @@ +package extended; + +import org.springframework.boot.SpringApplication; +import org.springframework.boot.autoconfigure.SpringBootApplication; +import org.springframework.cloud.streams.EnableChannelBinding; +import org.springframework.cloud.streams.aggregate.AggregateBuilder; +import org.springframework.cloud.streams.aggregate.AggregateConfigurer; + +import sink.LogSink; +import source.TimeSource; +import transform.LoggingTransformer; + +@SpringBootApplication +@EnableChannelBinding +public class ExtendedApplication implements AggregateConfigurer { + + @Override + public void configure(AggregateBuilder builder) { + // @formatter:off + builder + .from(TimeSource.class).as("source") + .via(LoggingTransformer.class) + .via(LoggingTransformer.class).profiles("other") + .to(LogSink.class); + // @formatter:on + } + + public static void main(String[] args) throws InterruptedException { + SpringApplication.run(ExtendedApplication.class, args); + } + +} diff --git a/vanilla-samples/extended/src/main/resources/application.yml b/vanilla-samples/extended/src/main/resources/application.yml new file mode 100644 index 000000000..a70e9bc69 --- /dev/null +++ b/vanilla-samples/extended/src/main/resources/application.yml @@ -0,0 +1,6 @@ +--- +spring: + profiles: other +module: + logging: + name: other \ No newline at end of file diff --git a/vanilla-samples/extended/src/main/resources/source.yml b/vanilla-samples/extended/src/main/resources/source.yml new file mode 100644 index 000000000..daf5166e5 --- /dev/null +++ b/vanilla-samples/extended/src/main/resources/source.yml @@ -0,0 +1 @@ +fixedDelay: 5000 diff --git a/vanilla-samples/extended/src/test/java/demo/ModuleApplicationTests.java b/vanilla-samples/extended/src/test/java/demo/ModuleApplicationTests.java new file mode 100644 index 000000000..65e2b9015 --- /dev/null +++ b/vanilla-samples/extended/src/test/java/demo/ModuleApplicationTests.java @@ -0,0 +1,22 @@ +package demo; + +import org.junit.Test; +import org.junit.runner.RunWith; +import org.springframework.test.annotation.DirtiesContext; +import org.springframework.test.context.web.WebAppConfiguration; +import org.springframework.boot.test.SpringApplicationConfiguration; +import org.springframework.test.context.junit4.SpringJUnit4ClassRunner; + +import extended.ExtendedApplication; + +@RunWith(SpringJUnit4ClassRunner.class) +@SpringApplicationConfiguration(classes = ExtendedApplication.class) +@WebAppConfiguration +@DirtiesContext +public class ModuleApplicationTests { + + @Test + public void contextLoads() { + } + +} diff --git a/vanilla-samples/pom.xml b/vanilla-samples/pom.xml index d9b1d3fd7..557c889f8 100644 --- a/vanilla-samples/pom.xml +++ b/vanilla-samples/pom.xml @@ -18,7 +18,29 @@ source sink + transform + double + extended + + + + org.springframework.cloud + spring-cloud-streams-sample-source + 1.0.0.BUILD-SNAPSHOT + + + org.springframework.cloud + spring-cloud-streams-sample-sink + 1.0.0.BUILD-SNAPSHOT + + + org.springframework.cloud + spring-cloud-streams-sample-transform + 1.0.0.BUILD-SNAPSHOT + + + diff --git a/vanilla-samples/sink/pom.xml b/vanilla-samples/sink/pom.xml index aa5b7e4b8..3fa7a16bd 100644 --- a/vanilla-samples/sink/pom.xml +++ b/vanilla-samples/sink/pom.xml @@ -13,7 +13,7 @@ org.springframework.cloud - spring-xd-samples + spring-cloud-streams-samples 1.0.0.BUILD-SNAPSHOT diff --git a/vanilla-samples/sink/src/main/java/demo/SinkApplication.java b/vanilla-samples/sink/src/main/java/demo/SinkApplication.java index b9f605791..bd26f4e77 100644 --- a/vanilla-samples/sink/src/main/java/demo/SinkApplication.java +++ b/vanilla-samples/sink/src/main/java/demo/SinkApplication.java @@ -4,10 +4,10 @@ import org.springframework.boot.SpringApplication; import org.springframework.boot.autoconfigure.SpringBootApplication; import org.springframework.context.annotation.ComponentScan; -import config.ModuleDefinition; +import sink.LogSink; @SpringBootApplication -@ComponentScan(basePackageClasses=ModuleDefinition.class) +@ComponentScan(basePackageClasses=LogSink.class) public class SinkApplication { public static void main(String[] args) throws InterruptedException { diff --git a/vanilla-samples/sink/src/main/java/sink/LogSink.java b/vanilla-samples/sink/src/main/java/sink/LogSink.java new file mode 100644 index 000000000..b28e4a03c --- /dev/null +++ b/vanilla-samples/sink/src/main/java/sink/LogSink.java @@ -0,0 +1,50 @@ +/* + * Copyright 2015 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 + * + * http://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 sink; + +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; +import org.springframework.cloud.streams.EnableChannelBinding; +import org.springframework.context.annotation.Bean; +import org.springframework.context.annotation.Configuration; +import org.springframework.integration.annotation.MessageEndpoint; +import org.springframework.integration.annotation.ServiceActivator; +import org.springframework.integration.channel.DirectChannel; +import org.springframework.messaging.MessageChannel; + +/** + * @author Dave Syer + * + */ +@Configuration +@EnableChannelBinding +@MessageEndpoint +public class LogSink { + + private static Logger logger = LoggerFactory.getLogger(LogSink.class); + + @Bean + public MessageChannel input() { + return new DirectChannel(); + } + + @ServiceActivator(inputChannel="input") + public void loggerSink(Object payload) { + logger.info("Received: " + payload); + } + +} diff --git a/vanilla-samples/source/pom.xml b/vanilla-samples/source/pom.xml index 28f0af948..2d5e39bbe 100644 --- a/vanilla-samples/source/pom.xml +++ b/vanilla-samples/source/pom.xml @@ -13,7 +13,7 @@ org.springframework.cloud - spring-xd-samples + spring-cloud-streams-samples 1.0.0.BUILD-SNAPSHOT diff --git a/vanilla-samples/source/src/main/java/demo/SourceApplication.java b/vanilla-samples/source/src/main/java/demo/SourceApplication.java index 38d7f2162..57f9eb186 100644 --- a/vanilla-samples/source/src/main/java/demo/SourceApplication.java +++ b/vanilla-samples/source/src/main/java/demo/SourceApplication.java @@ -4,10 +4,10 @@ import org.springframework.boot.SpringApplication; import org.springframework.boot.autoconfigure.SpringBootApplication; import org.springframework.context.annotation.ComponentScan; -import config.ModuleDefinition; +import source.TimeSource; @SpringBootApplication -@ComponentScan(basePackageClasses=ModuleDefinition.class) +@ComponentScan(basePackageClasses=TimeSource.class) public class SourceApplication { public static void main(String[] args) throws InterruptedException { diff --git a/vanilla-samples/source/src/main/java/config/ModuleDefinition.java b/vanilla-samples/source/src/main/java/source/TimeSource.java similarity index 97% rename from vanilla-samples/source/src/main/java/config/ModuleDefinition.java rename to vanilla-samples/source/src/main/java/source/TimeSource.java index cae7d4a94..9826971e5 100644 --- a/vanilla-samples/source/src/main/java/config/ModuleDefinition.java +++ b/vanilla-samples/source/src/main/java/source/TimeSource.java @@ -14,7 +14,7 @@ * limitations under the License. */ -package config; +package source; import java.text.SimpleDateFormat; import java.util.Date; @@ -38,7 +38,7 @@ import org.springframework.messaging.support.GenericMessage; @Configuration @EnableChannelBinding @EnableConfigurationProperties(TimeSourceOptionsMetadata.class) -public class ModuleDefinition { +public class TimeSource { @Autowired private TimeSourceOptionsMetadata options; diff --git a/vanilla-samples/source/src/main/java/config/TimeSourceOptionsMetadata.java b/vanilla-samples/source/src/main/java/source/TimeSourceOptionsMetadata.java similarity index 99% rename from vanilla-samples/source/src/main/java/config/TimeSourceOptionsMetadata.java rename to vanilla-samples/source/src/main/java/source/TimeSourceOptionsMetadata.java index 2e1dfb232..53481844a 100644 --- a/vanilla-samples/source/src/main/java/config/TimeSourceOptionsMetadata.java +++ b/vanilla-samples/source/src/main/java/source/TimeSourceOptionsMetadata.java @@ -14,7 +14,7 @@ * limitations under the License. */ -package config; +package source; import javax.validation.constraints.Min; import javax.validation.constraints.Pattern; diff --git a/vanilla-samples/source/src/main/resources/application.yml b/vanilla-samples/source/src/main/resources/application.yml index 6a079ed13..9d48be8a6 100644 --- a/vanilla-samples/source/src/main/resources/application.yml +++ b/vanilla-samples/source/src/main/resources/application.yml @@ -1,11 +1,21 @@ -fixedDelay: 5000 +server: + port: 8081 spring: cloud: channels: - inpoutChannelName: testtock + inputChannelName: testtock # uncomment below to use the last digit of the seconds as a partition key # hashcode(key) % N is then applied with N being the partitionCount value # thus, even seconds should go to the 0 queue, odd seconds to the 1 queue #producerProperties: # partitionKeyExpression: payload.charAt(payload.length()-1) # partitionCount: 2 + +--- +spring: + profiles: extended +spring: + cloud: + channels: + inputChannelName: xformed + \ No newline at end of file diff --git a/vanilla-samples/transform/pom.xml b/vanilla-samples/transform/pom.xml new file mode 100644 index 000000000..e2465e30e --- /dev/null +++ b/vanilla-samples/transform/pom.xml @@ -0,0 +1,69 @@ + + + 4.0.0 + + org.springframework.cloud + spring-cloud-streams-sample-transform + 1.0.0.BUILD-SNAPSHOT + jar + + spring-cloud-streams-sample-transform + Demo project for Spring XD module + + + org.springframework.cloud + spring-cloud-streams-samples + 1.0.0.BUILD-SNAPSHOT + + + + UTF-8 + demo.SinkApplication + 1.8 + + + + + org.springframework.cloud + spring-cloud-streams + + + org.springframework.xd + spring-xd-messagebus-redis + + + org.springframework.cloud + spring-cloud-lattice-connector + 1.0.2.BUILD-SNAPSHOT + + + org.springframework.boot + spring-boot-starter-redis + + + org.springframework.boot + spring-boot-configuration-processor + true + + + + org.springframework.boot + spring-boot-starter-test + test + + + + + + + org.springframework.boot + spring-boot-maven-plugin + + exec + + + + + + diff --git a/vanilla-samples/transform/src/main/java/demo/TransformApplication.java b/vanilla-samples/transform/src/main/java/demo/TransformApplication.java new file mode 100644 index 000000000..8a2dd6f2b --- /dev/null +++ b/vanilla-samples/transform/src/main/java/demo/TransformApplication.java @@ -0,0 +1,17 @@ +package demo; + +import org.springframework.boot.SpringApplication; +import org.springframework.boot.autoconfigure.SpringBootApplication; +import org.springframework.context.annotation.ComponentScan; + +import transform.LoggingTransformer; + +@SpringBootApplication +@ComponentScan(basePackageClasses=LoggingTransformer.class) +public class TransformApplication { + + public static void main(String[] args) throws InterruptedException { + SpringApplication.run(TransformApplication.class, args); + } + +} diff --git a/vanilla-samples/transform/src/main/java/transform/LoggingTransformer.java b/vanilla-samples/transform/src/main/java/transform/LoggingTransformer.java new file mode 100644 index 000000000..539bec045 --- /dev/null +++ b/vanilla-samples/transform/src/main/java/transform/LoggingTransformer.java @@ -0,0 +1,72 @@ +/* + * Copyright 2015 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 + * + * http://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 transform; + +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; +import org.springframework.boot.context.properties.ConfigurationProperties; +import org.springframework.cloud.streams.EnableChannelBinding; +import org.springframework.context.annotation.Bean; +import org.springframework.context.annotation.Configuration; +import org.springframework.integration.annotation.MessageEndpoint; +import org.springframework.integration.annotation.ServiceActivator; +import org.springframework.integration.channel.DirectChannel; +import org.springframework.messaging.MessageChannel; +import org.springframework.messaging.SubscribableChannel; + +/** + * @author Dave Syer + * + */ +@Configuration +@EnableChannelBinding +@MessageEndpoint +@ConfigurationProperties("module.logging") +public class LoggingTransformer { + + private static Logger logger = LoggerFactory.getLogger(LoggingTransformer.class); + + /** + * The name to include in the log message + */ + private String name = "logging"; + + public String getName() { + return this.name; + } + + public void setName(String name) { + this.name = name; + } + + @Bean + public MessageChannel input() { + return new DirectChannel(); + } + + @Bean + public SubscribableChannel output() { + return new DirectChannel(); + } + + @ServiceActivator(inputChannel = "input", outputChannel = "output") + public Object transform(Object payload) { + logger.info("Transformed by " + this.name + ": " + payload); + return payload; + } + +} diff --git a/vanilla-samples/transform/src/main/resources/application.yml b/vanilla-samples/transform/src/main/resources/application.yml new file mode 100644 index 000000000..edbdf3181 --- /dev/null +++ b/vanilla-samples/transform/src/main/resources/application.yml @@ -0,0 +1,8 @@ +server: + port: 8082 +spring: + cloud: + channels: + outputChannelName: xformed + inputChannelName: testtock + \ No newline at end of file diff --git a/vanilla-samples/transform/src/test/java/demo/ModuleApplicationTests.java b/vanilla-samples/transform/src/test/java/demo/ModuleApplicationTests.java new file mode 100644 index 000000000..4b781c269 --- /dev/null +++ b/vanilla-samples/transform/src/test/java/demo/ModuleApplicationTests.java @@ -0,0 +1,20 @@ +package demo; + +import org.junit.Test; +import org.junit.runner.RunWith; +import org.springframework.test.annotation.DirtiesContext; +import org.springframework.test.context.web.WebAppConfiguration; +import org.springframework.boot.test.SpringApplicationConfiguration; +import org.springframework.test.context.junit4.SpringJUnit4ClassRunner; + +@RunWith(SpringJUnit4ClassRunner.class) +@SpringApplicationConfiguration(classes = TransformApplication.class) +@WebAppConfiguration +@DirtiesContext +public class ModuleApplicationTests { + + @Test + public void contextLoads() { + } + +}