diff --git a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/function/StreamBridge.java b/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/function/StreamBridge.java index fd5fc9f95..fe4193e4e 100644 --- a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/function/StreamBridge.java +++ b/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/function/StreamBridge.java @@ -116,7 +116,9 @@ public final class StreamBridge implements SmartInitializingSingleton { /** * Sends 'data' to an output binding specified by 'bindingName' argument while * using default content type to deal with output type conversion (if necessary). - * @param bindingName the name of the output binding + * @param bindingName the name of the output binding. That said it requires a bit of clarification. + * When using bridge.send("foo"...), the 'foo' typically represents the binding name. However + * if such binding does not exist, the new binding will be created to support dynamic destinations. * @param data the data to send * @return true if data was sent successfully, otherwise false or throws an exception. */ @@ -133,7 +135,9 @@ public final class StreamBridge implements SmartInitializingSingleton { * provided via 'bindingName' does not have a corresponding binding such name will be * treated as dynamic destination. * - * @param bindingName the name of the output binding + * @param bindingName the name of the output binding. That said it requires a bit of clarification. + * When using bridge.send("foo"...), the 'foo' typically represents the binding name. However + * if such binding does not exist, the new binding will be created to support dynamic destinations. * @param data the data to send * @param outputContentType content type to be used to deal with output type conversion * @return true if data was sent successfully, otherwise false or throws an exception. diff --git a/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/function/StreamBridgeTests.java b/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/function/StreamBridgeTests.java index b72d86140..debca9553 100644 --- a/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/function/StreamBridgeTests.java +++ b/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/function/StreamBridgeTests.java @@ -17,6 +17,7 @@ package org.springframework.cloud.stream.function; import java.util.concurrent.atomic.AtomicBoolean; +import java.util.function.Consumer; import java.util.function.Function; import java.util.function.Supplier; @@ -30,6 +31,7 @@ import org.springframework.boot.builder.SpringApplicationBuilder; import org.springframework.cloud.stream.binder.test.OutputDestination; import org.springframework.cloud.stream.binder.test.TestChannelBinderConfiguration; import org.springframework.cloud.stream.binding.BinderAwareChannelResolver.NewDestinationBindingCallback; +import org.springframework.cloud.stream.config.BindingServiceProperties; import org.springframework.context.ConfigurableApplicationContext; import org.springframework.context.annotation.Bean; import org.springframework.integration.dsl.IntegrationFlow; @@ -53,6 +55,30 @@ public class StreamBridgeTests { System.clearProperty("spring.cloud.function.definition"); } + @Test + public void testBindingPropertiesAreHonored() { + try (ConfigurableApplicationContext context = new SpringApplicationBuilder(TestChannelBinderConfiguration + .getCompleteConfiguration(ConsumerConfiguration.class)) + .web(WebApplicationType.NONE).run( + "--spring.cloud.function.definition=consumer;function", + "--spring.jmx.enabled=false", + "--spring.cloud.stream.bindings.foo.destination=function-in-0", + "--spring.cloud.stream.bindings.foo.producer.partitionCount=5", + "--spring.cloud.stream.bindings.foo.consumer.concurrency=2")) { + + BindingServiceProperties bsProperties = context.getBean(BindingServiceProperties.class); + assertThat(bsProperties.getConsumerProperties("foo").getConcurrency()).isEqualTo(2); + assertThat(bsProperties.getProducerProperties("foo").getPartitionCount()).isEqualTo(5); + StreamBridge bridge = context.getBean(StreamBridge.class); + bridge.send("consumer-in-0", "hello foo"); + + OutputDestination outputDestination = context.getBean(OutputDestination.class); + Message message = outputDestination.receive(100, "function-out-0"); + assertThat(new String(message.getPayload())).isEqualTo("hello foo"); + assertThat(message.getHeaders().get("concurrency")).isEqualTo(2); + assertThat(message.getHeaders().get("partitionCount")).isEqualTo(5); + } + } //see https://github.com/spring-cloud/spring-cloud-function/issues/573 for more details @Test @@ -227,6 +253,29 @@ public class StreamBridgeTests { } + @EnableAutoConfiguration + public static class ConsumerConfiguration { + @Bean + public Consumer consumer(StreamBridge bridge, BindingServiceProperties properties) { + return v -> { + BindingServiceProperties p = properties; + bridge.send("foo", v); + }; + } + @Bean + public Function> function(StreamBridge bridge, BindingServiceProperties properties) { + return v -> { + int concurrency = properties.getConsumerProperties("foo").getConcurrency(); + int partitionCount = properties.getProducerProperties("foo").getPartitionCount(); + BindingServiceProperties p = properties; + return MessageBuilder.withPayload(v) + .setHeader("concurrency", concurrency) + .setHeader("partitionCount", partitionCount) + .build(); + }; + } + } + @EnableAutoConfiguration public static class TestConfiguration {