diff --git a/core/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/function/StreamBridge.java b/core/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/function/StreamBridge.java index 9a695c6a6..f6000de56 100644 --- a/core/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/function/StreamBridge.java +++ b/core/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/function/StreamBridge.java @@ -218,7 +218,8 @@ public final class StreamBridge implements SmartInitializingSingleton { functionToInvoke = new PartitionAwareFunctionWrapper(functionToInvoke, this.applicationContext, producerProperties); } - String targetType = this.resolveBinderTargetType(bindingName, MessageChannel.class, this.applicationContext.getBean(BinderFactory.class)); + String targetType = this.resolveBinderTargetType(bindingName, binderName, MessageChannel.class, + this.applicationContext.getBean(BinderFactory.class)); Message messageToSend = data instanceof Message ? MessageBuilder.fromMessage((Message) data).setHeaderIfAbsent(MessageUtils.TARGET_PROTOCOL, targetType).build() @@ -299,10 +300,10 @@ public final class StreamBridge implements SmartInitializingSingleton { return messageChannel; } - private String resolveBinderTargetType(String channelName, Class bindableType, BinderFactory binderFactory) { - String binderConfigurationName = this.bindingServiceProperties + private String resolveBinderTargetType(String channelName, String binderName, Class bindableType, BinderFactory binderFactory) { + String binderConfigurationName = binderName != null ? binderName : this.bindingServiceProperties .getBinder(channelName); - Binder binder = binderFactory.getBinder(binderConfigurationName, bindableType); + Binder binder = binderFactory.getBinder(binderConfigurationName, bindableType); String targetProtocol = binder.getClass().getSimpleName().startsWith("Rabbit") ? "amqp" : "kafka"; return targetProtocol; }