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 b86cc6e1a..73239fdab 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 @@ -272,7 +272,7 @@ public final class StreamBridge implements SmartInitializingSingleton { binder = binderFactory.getBinder(binderName, messageChannel.getClass()); } - if (producerProperties.isPartitioned()) { + if (producerProperties != null && producerProperties.isPartitioned()) { BindingProperties bindingProperties = this.bindingServiceProperties.getBindingProperties(destinationName); ((AbstractMessageChannel) messageChannel) .addInterceptor(new DefaultPartitioningInterceptor(bindingProperties, this.applicationContext.getBeanFactory())); 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 bffadc8e1..4dd9dec6b 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 @@ -35,6 +35,7 @@ import org.springframework.boot.WebApplicationType; import org.springframework.boot.autoconfigure.EnableAutoConfiguration; import org.springframework.boot.builder.SpringApplicationBuilder; import org.springframework.cloud.function.context.catalog.SimpleFunctionRegistry.FunctionInvocationWrapper; +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.binding.BinderAwareChannelResolver.NewDestinationBindingCallback; @@ -55,12 +56,15 @@ import org.springframework.messaging.MessageHeaders; import org.springframework.messaging.converter.AbstractMessageConverter; import org.springframework.messaging.converter.MessageConverter; import org.springframework.messaging.support.ChannelInterceptor; +import org.springframework.messaging.support.GenericMessage; import org.springframework.messaging.support.MessageBuilder; import org.springframework.util.MimeType; import org.springframework.util.MimeTypeUtils; import org.springframework.util.ReflectionUtils; import static org.assertj.core.api.Assertions.assertThat; +import static org.junit.Assert.assertEquals; +import static org.junit.Assert.assertNotNull; import static org.junit.Assert.fail; /** @@ -467,6 +471,35 @@ public class StreamBridgeTests { assertThat(context.getBean("callbackVerifier", AtomicBoolean.class)).isTrue(); } } + @EnableAutoConfiguration + public static class DynamicProducerConfig { + @Bean + public Function, Message> uppercase() { + return msg -> MessageBuilder.withPayload(msg.getPayload().toUpperCase()) + .setHeader("spring.cloud.stream.sendto.destination", "dynamicTopic").build(); + } + } + + @Test + public void testDynamicProducerDestination() { + System.clearProperty("spring.cloud.function.definition"); + ConfigurableApplicationContext context = new SpringApplicationBuilder( + TestChannelBinderConfiguration.getCompleteConfiguration(DynamicProducerConfig.class)) + .web(WebApplicationType.NONE) + .run("--spring.jmx.enabled=false", + "--spring.cloud.function.definition=uppercase", + "--spring.cloud.stream.bindings.uppercase-in-0.destination=upper" + ); + + InputDestination source = context.getBean(InputDestination.class); + source.send(new GenericMessage<>("John Doe".getBytes()), "upper"); + + OutputDestination target = context.getBean(OutputDestination.class); + Message message = target.receive(5, "dynamicTopic"); + + assertNotNull(message); + assertEquals(new String(message.getPayload()), "JOHN DOE"); + } @EnableAutoConfiguration public static class EmptyConfiguration {