From 710d45a8a5431bf7e36e63e94850f77d6c9baaf8 Mon Sep 17 00:00:00 2001 From: Marius Bogoevici Date: Thu, 18 Aug 2016 12:28:59 -0400 Subject: [PATCH] Changes required by the refactoring of AbstractMessageChannelBinder --- .../binder/rabbit/RabbitMessageChannelBinder.java | 12 ++++++------ 1 file changed, 6 insertions(+), 6 deletions(-) diff --git a/spring-cloud-stream-binder-rabbit/src/main/java/org/springframework/cloud/stream/binder/rabbit/RabbitMessageChannelBinder.java b/spring-cloud-stream-binder-rabbit/src/main/java/org/springframework/cloud/stream/binder/rabbit/RabbitMessageChannelBinder.java index a88249e89..db856f436 100644 --- a/spring-cloud-stream-binder-rabbit/src/main/java/org/springframework/cloud/stream/binder/rabbit/RabbitMessageChannelBinder.java +++ b/spring-cloud-stream-binder-rabbit/src/main/java/org/springframework/cloud/stream/binder/rabbit/RabbitMessageChannelBinder.java @@ -83,7 +83,7 @@ import org.springframework.util.StringUtils; */ public class RabbitMessageChannelBinder extends AbstractMessageChannelBinder, - ExtendedProducerProperties, Queue> + ExtendedProducerProperties, Queue, TopicExchange> implements ExtendedPropertiesBinder { private static final AnonymousQueue.Base64UrlNamingStrategy ANONYMOUS_GROUP_NAME_GENERATOR @@ -328,21 +328,21 @@ public class RabbitMessageChannelBinder } @Override - protected void createProducerDestinationIfNecessary(String name, + protected TopicExchange createProducerDestinationIfNecessary(String name, ExtendedProducerProperties producerProperties) { String exchangeName = applyPrefix(producerProperties.getExtension().getPrefix(), name); TopicExchange exchange = new TopicExchange(exchangeName); declareExchange(exchangeName, exchange); + return exchange; } @Override - protected MessageHandler createProducerMessageHandler(final String destination, + protected MessageHandler createProducerMessageHandler(final TopicExchange exchange, ExtendedProducerProperties properties) throws Exception { String prefix = properties.getExtension().getPrefix(); - String exchangeName = applyPrefix(prefix, destination); - TopicExchange exchange = new TopicExchange(exchangeName); - declareExchange(exchangeName, exchange); + String exchangeName = exchange.getName(); + String destination = StringUtils.isEmpty(prefix) ? exchangeName : exchangeName.substring(prefix.length()); final AmqpOutboundEndpoint endpoint = new AmqpOutboundEndpoint(buildRabbitTemplate(properties.getExtension())); endpoint.setExchangeName(exchange.getName()); if (!properties.isPartitioned()) {