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 940de41d5..d96989415 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 @@ -167,6 +167,9 @@ public final class StreamBridge implements SmartInitializingSingleton { SubscribableChannel resolveDestination(String destinationName, ProducerProperties producerProperties) { SubscribableChannel messageChannel = this.channelCache.get(destinationName); + if (messageChannel == null && this.applicationContext.containsBean(destinationName)) { + messageChannel = this.applicationContext.getBean(destinationName, SubscribableChannel.class); + } if (messageChannel == null) { messageChannel = new DirectWithAttributesChannel(); this.bindingService.bindProducer(messageChannel, destinationName, false); 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 c83f9fa6e..7c74e2bff 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 @@ -29,6 +29,8 @@ import org.springframework.cloud.stream.binder.test.OutputDestination; import org.springframework.cloud.stream.binder.test.TestChannelBinderConfiguration; import org.springframework.context.ConfigurableApplicationContext; import org.springframework.context.annotation.Bean; +import org.springframework.integration.dsl.IntegrationFlow; +import org.springframework.integration.dsl.IntegrationFlows; import org.springframework.messaging.Message; import org.springframework.messaging.support.MessageBuilder; @@ -170,6 +172,21 @@ public class StreamBridgeTests { } } + @Test + public void testWithIntegrationFlowBecauseMarcinSaidSo() { + try (ConfigurableApplicationContext context = new SpringApplicationBuilder(TestChannelBinderConfiguration + .getCompleteConfiguration(IntegrationFlowConfiguration.class)) + .web(WebApplicationType.NONE).run("--spring.jmx.enabled=false")) { + + StreamBridge bridge = context.getBean(StreamBridge.class); + bridge.send("foo", "blah"); + + OutputDestination outputDestination = context.getBean(OutputDestination.class); + Message message = outputDestination.receive(100, "output"); + assertThat(new String(message.getPayload())).isEqualTo("BLAH"); + } + } + @EnableAutoConfiguration public static class EmptyConfiguration { @@ -183,4 +200,18 @@ public class StreamBridgeTests { return () -> "hello"; } } + + @EnableAutoConfiguration + public static class IntegrationFlowConfiguration { + + @Bean + public IntegrationFlow transform(StreamBridge bridge) { + return IntegrationFlows.from("foo").transform(v -> { + String s = new String((byte[]) v); + return s.toUpperCase(); + }) + .handle(v -> bridge.send("output", v)) + .get(); + } + } }