diff --git a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/config/BindingServiceConfiguration.java b/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/config/BindingServiceConfiguration.java
index 12124b684..cbe058ae9 100644
--- a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/config/BindingServiceConfiguration.java
+++ b/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/config/BindingServiceConfiguration.java
@@ -44,8 +44,7 @@ import org.springframework.cloud.stream.binding.InputBindingLifecycle;
import org.springframework.cloud.stream.binding.MessageChannelStreamListenerResultAdapter;
import org.springframework.cloud.stream.binding.OutputBindingLifecycle;
import org.springframework.cloud.stream.binding.StreamListenerAnnotationBeanPostProcessor;
-import org.springframework.cloud.stream.function.FunctionConfiguration;
-import org.springframework.cloud.stream.function.FunctionProperties;
+import org.springframework.cloud.stream.function.StreamFunctionProperties;
import org.springframework.cloud.stream.micrometer.DestinationPublishingMetricsAutoConfiguration;
import org.springframework.context.ApplicationListener;
import org.springframework.context.annotation.Bean;
@@ -78,7 +77,7 @@ import org.springframework.util.Assert;
* @author Soby Chacko
*/
@Configuration
-@EnableConfigurationProperties({ BindingServiceProperties.class, SpringIntegrationProperties.class, FunctionProperties.class })
+@EnableConfigurationProperties({ BindingServiceProperties.class, SpringIntegrationProperties.class, StreamFunctionProperties.class })
@Import({ DestinationPublishingMetricsAutoConfiguration.class, SpelExpressionConverterConfiguration.class })
@Role(BeanDefinition.ROLE_INFRASTRUCTURE)
@ConditionalOnBean(value = BinderTypeRegistry.class, search = SearchStrategy.CURRENT)
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 37b4a6551..c42654ef4 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
@@ -17,14 +17,10 @@
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;
@@ -32,9 +28,6 @@ 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;
/**
*
@@ -43,7 +36,7 @@ import org.springframework.util.ClassUtils;
* @since 2.1
*/
@Configuration
-@ConditionalOnProperty("spring.cloud.stream.function.name")
+@ConditionalOnProperty("spring.cloud.stream.function.definition")
public class FunctionConfiguration {
@Autowired(required=false)
@@ -55,13 +48,10 @@ public class FunctionConfiguration {
@Autowired(required=false)
private Sink sink;
- @Autowired
- private ConfigurableListableBeanFactory registry;
-
@Bean
public IntegrationFlowFunctionSupport functionSupport(FunctionCatalogWrapper functionCatalog,
FunctionInspector functionInspector, CompositeMessageConverterFactory messageConverterFactory,
- FunctionProperties functionProperties) {
+ StreamFunctionProperties functionProperties) {
return new IntegrationFlowFunctionSupport(functionCatalog, functionInspector, messageConverterFactory,
functionProperties);
@@ -72,11 +62,13 @@ public class FunctionConfiguration {
return new FunctionCatalogWrapper(catalog);
}
-
- @ConditionalOnProperty("spring.cloud.stream.function.name")
- @ConditionalOnMissingBean
+ /**
+ * This configuration creates an instance of {@link IntegrationFlow} appropriate for binding declared using EnableBinding.
+ * At the moment only Source, Processor and Sink are supported.
+ */
+ @ConditionalOnMissingBean // starter apps typically already provide and instance of IntegrationFlow, so we don't need this one.
@Bean
- public IntegrationFlow foo(IntegrationFlowFunctionSupport functionSupport) {
+ public IntegrationFlow integrationFlowCreator(IntegrationFlowFunctionSupport functionSupport) {
if (processor != null) {
return functionSupport.integrationFlowForFunction(processor.input(), processor.output()).get();
}
@@ -86,14 +78,6 @@ public class FunctionConfiguration {
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");
+ throw new UnsupportedOperationException("Bindings other then Source, Processor and Sink are not currently supported");
}
}
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 2bd645dc3..27e977f17 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
@@ -53,7 +53,7 @@ public class IntegrationFlowFunctionSupport {
private final CompositeMessageConverterFactory messageConverterFactory;
- private final FunctionProperties functionProperties;
+ private final StreamFunctionProperties functionProperties;
@Autowired
private MessageChannel errorChannel;
@@ -65,7 +65,7 @@ public class IntegrationFlowFunctionSupport {
* @param functionProperties
*/
public IntegrationFlowFunctionSupport(FunctionCatalogWrapper functionCatalog, FunctionInspector functionInspector,
- CompositeMessageConverterFactory messageConverterFactory, FunctionProperties functionProperties) {
+ CompositeMessageConverterFactory messageConverterFactory, StreamFunctionProperties functionProperties) {
Assert.notNull(functionCatalog, "'functionCatalog' must not be null");
Assert.notNull(functionInspector, "'functionInspector' must not be null");
@@ -78,18 +78,18 @@ public class IntegrationFlowFunctionSupport {
}
public FunctionType getCurrentFunctionType() {
- return functionInspector.getRegistration(functionCatalog.lookup(this.functionProperties.getName())).getType();
+ return functionInspector.getRegistration(functionCatalog.lookup(this.functionProperties.getDefinition())).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.
+ * The name of the bean must be provided via `spring.cloud.stream.function.definition` property.
* @return instance of {@link IntegrationFlowBuilder}
* @throws IllegalStateException if the named Supplier can not be located.
*/
public IntegrationFlowBuilder integrationFlowFromNamedSupplier() {
- if (StringUtils.hasText(this.functionProperties.getName())) {
- Supplier> supplier = functionCatalog.lookup(Supplier.class, this.functionProperties.getName());
+ if (StringUtils.hasText(this.functionProperties.getDefinition())) {
+ Supplier> supplier = functionCatalog.lookup(Supplier.class, this.functionProperties.getDefinition());
if (supplier instanceof FluxSupplier) {
supplier = ((FluxSupplier>)supplier).getTarget();
}
@@ -98,7 +98,7 @@ public class IntegrationFlowFunctionSupport {
}
throw new IllegalStateException(
- "A Supplier is not specified in the 'spring.cloud.stream.function.name' property.");
+ "A Supplier is not specified in the 'spring.cloud.stream.function.definition' property.");
}
/**
@@ -135,7 +135,7 @@ public class IntegrationFlowFunctionSupport {
/**
* 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.
+ * The name of the bean must be provided via `spring.cloud.stream.function.definition` property.
*
* NOTE: If this method returns true, the integration flow is now represented
* as a Reactive Streams {@link Publisher} bean.
@@ -146,9 +146,9 @@ public class IntegrationFlowFunctionSupport {
* @return true if {@link Function} was located and added and false if it wasn't.
*/
public boolean andThenFunction(IntegrationFlowBuilder flowBuilder, MessageChannel outputChannel) {
- if (StringUtils.hasText(this.functionProperties.getName())) {
+ if (StringUtils.hasText(this.functionProperties.getDefinition())) {
FunctionInvoker functionInvoker =
- new FunctionInvoker<>(this.functionProperties.getName(), this.functionCatalog,
+ new FunctionInvoker<>(this.functionProperties.getDefinition(), this.functionCatalog,
this.functionInspector, this.messageConverterFactory, this.errorChannel);
if (outputChannel != null) {
diff --git a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/function/FunctionProperties.java b/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/function/StreamFunctionProperties.java
similarity index 72%
rename from spring-cloud-stream/src/main/java/org/springframework/cloud/stream/function/FunctionProperties.java
rename to spring-cloud-stream/src/main/java/org/springframework/cloud/stream/function/StreamFunctionProperties.java
index 39488dae1..74677c86e 100644
--- a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/function/FunctionProperties.java
+++ b/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/function/StreamFunctionProperties.java
@@ -25,20 +25,21 @@ import org.springframework.boot.context.properties.ConfigurationProperties;
* @since 2.1
*/
@ConfigurationProperties("spring.cloud.stream.function")
-public class FunctionProperties {
+public class StreamFunctionProperties {
/**
- * Name of functions to bind. If several functions need to be composed into one, use pipes (e.g., 'fooFunc|barFunc')
+ * Definition of functions to bind. If several functions need to be composed
+ * into one, use pipes (e.g., 'fooFunc|barFunc')
*/
- private String name;
+ private String definition;
- public String getName() {
- return this.name;
+ public String getDefinition() {
+ return this.definition;
}
- public void setName(String name) {
- this.name = name;
+ public void setDefinition(String definition) {
+ this.definition = definition;
}
}
diff --git a/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/function/FunctionInvokerTests.java b/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/function/FunctionInvokerTests.java
index 6ec512a3d..e53e346ee 100644
--- a/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/function/FunctionInvokerTests.java
+++ b/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/function/FunctionInvokerTests.java
@@ -84,7 +84,7 @@ public class FunctionInvokerTests {
}
@Bean
- public Function messageToMessageNoType() {
+ public Function, Message>> messageToMessageNoType() {
return x -> MessageBuilder.withPayload(new Bar()).copyHeaders(x.getHeaders()).build();
}
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
new file mode 100644
index 000000000..458e9a637
--- /dev/null
+++ b/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/function/GreenfieldFunctionEnableBindingTests.java
@@ -0,0 +1,127 @@
+/*
+ * Copyright 2018 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.stream.function;
+
+import java.nio.charset.StandardCharsets;
+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.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.ConfigurableApplicationContext;
+import org.springframework.context.annotation.Bean;
+import org.springframework.integration.channel.QueueChannel;
+import org.springframework.messaging.Message;
+import org.springframework.messaging.PollableChannel;
+import org.springframework.messaging.support.GenericMessage;
+
+import static org.assertj.core.api.Assertions.assertThat;
+
+/**
+ * This test validates proper function binding for applications where EnableBinding is declared.
+ *
+ * @author Oleg Zhurakousky
+ */
+public class GreenfieldFunctionEnableBindingTests {
+
+ @Test
+ public void testSourceFromSupplier() {
+ try (ConfigurableApplicationContext context = new SpringApplicationBuilder(
+ TestChannelBinderConfiguration.getCompleteConfiguration(SourceFromSupplier.class)).web(
+ WebApplicationType.NONE).run("--spring.cloud.stream.function.definition=date", "--spring.jmx.enabled=false")) {
+
+ OutputDestination target = context.getBean(OutputDestination.class);
+ Message sourceMessage = target.receive(10000);
+ Date date = (Date) new CompositeMessageConverterFactory().getMessageConverterForAllRegistered().fromMessage(sourceMessage, Date.class);
+ assertThat(date).isEqualTo(new Date(12345L));
+ }
+ }
+
+ @Test
+ public void testProcessorFromFunction() {
+ try (ConfigurableApplicationContext context = new SpringApplicationBuilder(
+ TestChannelBinderConfiguration.getCompleteConfiguration(ProcessorFromFunction.class)).web(
+ WebApplicationType.NONE).run("--spring.cloud.stream.function.definition=toUpperCase", "--spring.jmx.enabled=false")) {
+
+ InputDestination source = context.getBean(InputDestination.class);
+ source.send(new GenericMessage("John Doe".getBytes()));
+ OutputDestination target = context.getBean(OutputDestination.class);
+ assertThat(target.receive(10000).getPayload()).isEqualTo("JOHN DOE".getBytes(StandardCharsets.UTF_8));
+ }
+ }
+
+ @Test
+ public void testSinkFromConsumer() {
+ try (ConfigurableApplicationContext context = new SpringApplicationBuilder(
+ TestChannelBinderConfiguration.getCompleteConfiguration(SinkFromConsumer.class)).web(
+ WebApplicationType.NONE).run("--spring.cloud.stream.function.definition=sink", "--spring.jmx.enabled=false")) {
+
+ InputDestination source = context.getBean(InputDestination.class);
+ PollableChannel result = context.getBean("result", PollableChannel.class);
+ source.send(new GenericMessage("John Doe".getBytes()));
+ assertThat(result.receive(10000).getPayload()).isEqualTo("John Doe");
+ }
+ }
+
+
+ @EnableAutoConfiguration
+ @EnableBinding(Source.class)
+ public static class SourceFromSupplier {
+ @Bean
+ public Supplier date() {
+ return () -> new Date(12345L);
+ }
+ }
+
+ @EnableAutoConfiguration
+ @EnableBinding(Processor.class)
+ public static class ProcessorFromFunction {
+ @Bean
+ public Function toUpperCase() {
+ return s -> s.toUpperCase();
+ }
+ }
+
+ @EnableAutoConfiguration
+ @EnableBinding(Sink.class)
+ public static class SinkFromConsumer {
+ @Bean
+ public PollableChannel result() {
+ return new QueueChannel();
+ }
+ @Bean
+ public Consumer sink(PollableChannel result) {
+ return s -> {
+ result.send(new GenericMessage(s));
+ System.out.println(s);
+ };
+ }
+ }
+}
diff --git a/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/function/NewSourceAsSupplierTests.java b/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/function/NewSourceAsSupplierTests.java
deleted file mode 100644
index 5cd6b3274..000000000
--- a/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/function/NewSourceAsSupplierTests.java
+++ /dev/null
@@ -1,132 +0,0 @@
-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 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("fopo".getBytes()));
- OutputDestination target = context.getBean(OutputDestination.class);
- Message 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("fopo".getBytes()));
-// OutputDestination target = context.getBean(OutputDestination.class);
-// Message 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("Hello No Binding".getBytes()));
-// OutputDestination target = context.getBean(OutputDestination.class);
-// Message 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() {
- return () -> new Date();
- }
- }
-
- @EnableAutoConfiguration
- @EnableBinding(Processor.class)
- public static class ProcessorFromFunction {
- @Bean
- public Function toUpperCase() {
- return s -> s.toUpperCase();
- }
- }
-
- @EnableAutoConfiguration
- @EnableBinding(Sink.class)
- public static class SinkFromConsumer {
- @Bean
- public Consumer sink() {
- return s -> System.out.println(s);
- }
- }
-
- @EnableAutoConfiguration
-// @EnableBinding(Sink.class)
- public static class SinkFromConsumerNoEnableBinding {
- @Bean
- public Consumer sink() {
- return s -> System.out.println("==> " + s);
- }
- }
-}
diff --git a/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/function/ProcessorToFunctionsSupportTests.java b/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/function/ProcessorToFunctionsSupportTests.java
index 0e738d8d2..d74602d66 100644
--- a/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/function/ProcessorToFunctionsSupportTests.java
+++ b/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/function/ProcessorToFunctionsSupportTests.java
@@ -74,7 +74,7 @@ public class ProcessorToFunctionsSupportTests {
new SpringApplicationBuilder(
TestChannelBinderConfiguration.getCompleteConfiguration(FunctionsConfiguration.class))
.web(WebApplicationType.NONE)
- .run("--spring.cloud.stream.function.name=toUpperCase", "--spring.jmx.enabled=false");
+ .run("--spring.cloud.stream.function.definition=toUpperCase", "--spring.jmx.enabled=false");
InputDestination source = context.getBean(InputDestination.class);
OutputDestination target = context.getBean(OutputDestination.class);
@@ -88,7 +88,7 @@ public class ProcessorToFunctionsSupportTests {
new SpringApplicationBuilder(
TestChannelBinderConfiguration.getCompleteConfiguration(FunctionsConfiguration.class))
.web(WebApplicationType.NONE)
- .run("--spring.cloud.stream.function.name=toUpperCase|concatWithSelf", "--spring.jmx.enabled=false");
+ .run("--spring.cloud.stream.function.definition=toUpperCase|concatWithSelf", "--spring.jmx.enabled=false");
InputDestination source = context.getBean(InputDestination.class);
OutputDestination target = context.getBean(OutputDestination.class);
@@ -102,7 +102,7 @@ public class ProcessorToFunctionsSupportTests {
new SpringApplicationBuilder(
TestChannelBinderConfiguration.getCompleteConfiguration(ConsumerConfiguration.class))
.web(WebApplicationType.NONE)
- .run("--spring.cloud.stream.function.name=log", "--spring.jmx.enabled=false");
+ .run("--spring.cloud.stream.function.definition=log", "--spring.jmx.enabled=false");
InputDestination source = context.getBean(InputDestination.class);
OutputDestination target = context.getBean(OutputDestination.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 47b40bbb7..ae4a70c4b 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
@@ -72,7 +72,7 @@ public class SourceToFunctionsSupportTests {
try (ConfigurableApplicationContext context = new SpringApplicationBuilder(
TestChannelBinderConfiguration.getCompleteConfiguration(FunctionsConfiguration.class)).web(
WebApplicationType.NONE)
- .run("--spring.cloud.stream.function.name=toUpperCase", "--spring.jmx.enabled=false")) {
+ .run("--spring.cloud.stream.function.definition=toUpperCase", "--spring.jmx.enabled=false")) {
OutputDestination target = context.getBean(OutputDestination.class);
assertThat(target.receive(1000).getPayload()).isEqualTo("HELLO FUNCTION".getBytes(StandardCharsets.UTF_8));
@@ -84,7 +84,7 @@ public class SourceToFunctionsSupportTests {
try (ConfigurableApplicationContext context = new SpringApplicationBuilder(
TestChannelBinderConfiguration.getCompleteConfiguration(FunctionsConfiguration.class)).web(
WebApplicationType.NONE)
- .run("--spring.cloud.stream.function.name=toUpperCase|concatWithSelf", "--spring.jmx.enabled=false")) {
+ .run("--spring.cloud.stream.function.definition=toUpperCase|concatWithSelf", "--spring.jmx.enabled=false")) {
OutputDestination target = context.getBean(OutputDestination.class);
assertThat(target.receive(1000).getPayload()).isEqualTo(
"HELLO FUNCTION:HELLO FUNCTION".getBytes(StandardCharsets.UTF_8));
@@ -97,7 +97,7 @@ public class SourceToFunctionsSupportTests {
new SpringApplicationBuilder(
TestChannelBinderConfiguration.getCompleteConfiguration(FunctionsConfigurationNoConversionPossible.class))
.web(WebApplicationType.NONE)
- .run("--spring.cloud.stream.function.name=toUpperCase|concatWithSelf",
+ .run("--spring.cloud.stream.function.definition=toUpperCase|concatWithSelf",
"--spring.jmx.enabled=false")) {
PollableChannel errorChannel = context.getBean("errorChannel", PollableChannel.class);
OutputDestination target = context.getBean(OutputDestination.class);
@@ -112,7 +112,7 @@ public class SourceToFunctionsSupportTests {
new SpringApplicationBuilder(
TestChannelBinderConfiguration.getCompleteConfiguration(FunctionsConfigurationNoConversionPossible.class))
.web(WebApplicationType.NONE)
- .run("--spring.cloud.stream.function.name=toUpperCase|concatWithSelf",
+ .run("--spring.cloud.stream.function.definition=toUpperCase|concatWithSelf",
"--spring.jmx.enabled=false")) {
PollableChannel errorChannel = context.getBean("errorChannel", PollableChannel.class);
OutputDestination target = context.getBean(OutputDestination.class);
@@ -125,7 +125,7 @@ public class SourceToFunctionsSupportTests {
public void testMessageSourceIsCreatedFromProvidedSupplier() {
try (ConfigurableApplicationContext context = new SpringApplicationBuilder(
TestChannelBinderConfiguration.getCompleteConfiguration(SupplierConfiguration.class)).web(
- WebApplicationType.NONE).run("--spring.cloud.stream.function.name=number", "--spring.jmx.enabled=false")) {
+ WebApplicationType.NONE).run("--spring.cloud.stream.function.definition=number", "--spring.jmx.enabled=false")) {
OutputDestination target = context.getBean(OutputDestination.class);
assertThat(target.receive(10000).getPayload()).isEqualTo("1".getBytes(StandardCharsets.UTF_8));
@@ -140,7 +140,7 @@ public class SourceToFunctionsSupportTests {
try (ConfigurableApplicationContext context = new SpringApplicationBuilder(
TestChannelBinderConfiguration.getCompleteConfiguration(SupplierConfiguration.class)).web(
WebApplicationType.NONE)
- .run("--spring.cloud.stream.function.name=number|concatWithSelf", "--spring.jmx.enabled=false")) {
+ .run("--spring.cloud.stream.function.definition=number|concatWithSelf", "--spring.jmx.enabled=false")) {
OutputDestination target = context.getBean(OutputDestination.class);
assertThat(target.receive(10000).getPayload()).isEqualTo("11".getBytes(StandardCharsets.UTF_8));
@@ -155,7 +155,7 @@ public class SourceToFunctionsSupportTests {
try (ConfigurableApplicationContext context = new SpringApplicationBuilder(
TestChannelBinderConfiguration.getCompleteConfiguration(SupplierConfiguration.class)).web(
WebApplicationType.NONE)
- .run("--spring.cloud.stream.function.name=number|concatWithSelf|multiplyByTwo",
+ .run("--spring.cloud.stream.function.definition=number|concatWithSelf|multiplyByTwo",
"--spring.jmx.enabled=false")) {
OutputDestination target = context.getBean(OutputDestination.class);
@@ -177,7 +177,7 @@ public class SourceToFunctionsSupportTests {
new SpringApplicationBuilder(
TestChannelBinderConfiguration.getCompleteConfiguration(SupplierConfiguration.class)).web(
WebApplicationType.NONE)
- .run("--spring.cloud.stream.function.name=doesNotExist", "--spring.jmx.enabled=false");
+ .run("--spring.cloud.stream.function.definition=doesNotExist", "--spring.jmx.enabled=false");
}
@EnableAutoConfiguration
@@ -297,11 +297,11 @@ public class SourceToFunctionsSupportTests {
private Source source;
@Autowired
- private FunctionProperties functionProperties;
+ private StreamFunctionProperties functionProperties;
@Bean
public IntegrationFlow messageSourceFlow(IntegrationFlowFunctionSupport functionSupport) {
- Assert.hasText(this.functionProperties.getName(), "Supplier name must be provided");
+ Assert.hasText(this.functionProperties.getDefinition(), "Supplier name must be provided");
return functionSupport.integrationFlowFromNamedSupplier().channel(this.source.output()).get();
}