From ff4a979254fa9061562a4d91d86a4004930ff386 Mon Sep 17 00:00:00 2001 From: Mark Fisher Date: Fri, 9 Sep 2011 20:46:49 +0100 Subject: [PATCH] added support for 'exchange-name-expression' on AMQP outbound adapters --- .../AmqpOutboundChannelAdapterParser.java | 1 + .../config/AmqpOutboundGatewayParser.java | 3 +- .../amqp/outbound/AmqpOutboundEndpoint.java | 41 ++++++++++++++----- 3 files changed, 33 insertions(+), 12 deletions(-) diff --git a/spring-integration-amqp/src/main/java/org/springframework/integration/amqp/config/AmqpOutboundChannelAdapterParser.java b/spring-integration-amqp/src/main/java/org/springframework/integration/amqp/config/AmqpOutboundChannelAdapterParser.java index 064bccb6e5..c91e6c3bb9 100644 --- a/spring-integration-amqp/src/main/java/org/springframework/integration/amqp/config/AmqpOutboundChannelAdapterParser.java +++ b/spring-integration-amqp/src/main/java/org/springframework/integration/amqp/config/AmqpOutboundChannelAdapterParser.java @@ -42,6 +42,7 @@ public class AmqpOutboundChannelAdapterParser extends AbstractOutboundChannelAda } builder.addConstructorArgReference(amqpTemplateRef); IntegrationNamespaceUtils.setValueIfAttributeDefined(builder, element, "exchange-name"); + IntegrationNamespaceUtils.setValueIfAttributeDefined(builder, element, "exchange-name-expression"); IntegrationNamespaceUtils.setValueIfAttributeDefined(builder, element, "routing-key"); IntegrationNamespaceUtils.setValueIfAttributeDefined(builder, element, "routing-key-expression"); return builder.getBeanDefinition(); diff --git a/spring-integration-amqp/src/main/java/org/springframework/integration/amqp/config/AmqpOutboundGatewayParser.java b/spring-integration-amqp/src/main/java/org/springframework/integration/amqp/config/AmqpOutboundGatewayParser.java index c513e86194..c063d50b93 100644 --- a/spring-integration-amqp/src/main/java/org/springframework/integration/amqp/config/AmqpOutboundGatewayParser.java +++ b/spring-integration-amqp/src/main/java/org/springframework/integration/amqp/config/AmqpOutboundGatewayParser.java @@ -45,9 +45,10 @@ public class AmqpOutboundGatewayParser extends AbstractConsumerEndpointParser { builder.addConstructorArgReference(amqpTemplateRef); builder.addPropertyValue("expectReply", true); IntegrationNamespaceUtils.setValueIfAttributeDefined(builder, element, "exchange-name"); + IntegrationNamespaceUtils.setValueIfAttributeDefined(builder, element, "exchange-name-expression"); IntegrationNamespaceUtils.setValueIfAttributeDefined(builder, element, "routing-key"); - IntegrationNamespaceUtils.setReferenceIfAttributeDefined(builder, element, "reply-channel", "outputChannel"); IntegrationNamespaceUtils.setValueIfAttributeDefined(builder, element, "routing-key-expression"); + IntegrationNamespaceUtils.setReferenceIfAttributeDefined(builder, element, "reply-channel", "outputChannel"); return builder; } diff --git a/spring-integration-amqp/src/main/java/org/springframework/integration/amqp/outbound/AmqpOutboundEndpoint.java b/spring-integration-amqp/src/main/java/org/springframework/integration/amqp/outbound/AmqpOutboundEndpoint.java index b7ebf76aae..2723f94cc3 100644 --- a/spring-integration-amqp/src/main/java/org/springframework/integration/amqp/outbound/AmqpOutboundEndpoint.java +++ b/spring-integration-amqp/src/main/java/org/springframework/integration/amqp/outbound/AmqpOutboundEndpoint.java @@ -41,19 +41,24 @@ import org.springframework.util.Assert; */ public class AmqpOutboundEndpoint extends AbstractReplyProducingMessageHandler { + private static final ExpressionParser expressionParser = new SpelExpressionParser(new SpelParserConfiguration(true, true)); + + private final AmqpTemplate amqpTemplate; + private volatile boolean expectReply; + private volatile String exchangeName = ""; private volatile String routingKey = ""; - private volatile boolean expectReply; + private volatile String exchangeNameExpression; - private static final ExpressionParser expressionParser = new SpelExpressionParser(new SpelParserConfiguration(true, true)); + private volatile String routingKeyExpression; private volatile ExpressionEvaluatingMessageProcessor routingKeyGenerator; - private volatile String routingKeyExpression; + private volatile ExpressionEvaluatingMessageProcessor exchangeNameGenerator; private volatile AmqpHeaderMapper headerMapper = new DefaultAmqpHeaderMapper(); @@ -61,9 +66,15 @@ public class AmqpOutboundEndpoint extends AbstractReplyProducingMessageHandler { @Override protected void onInit() { super.onInit(); + Assert.state(exchangeNameExpression == null || "".equals(exchangeName), + "Either an exchangeName or an exchangeNameExpression can be provided, but not both"); + if (exchangeNameExpression != null) { + Expression expression = expressionParser.parseExpression(this.exchangeNameExpression); + this.exchangeNameGenerator = new ExpressionEvaluatingMessageProcessor(expression, String.class); + } Assert.state(routingKeyExpression == null || "".equals(routingKey), "Either a routingKey or a routingKeyExpression can be provided, but not both"); - if (routingKeyExpression!=null) { + if (routingKeyExpression != null) { Expression expression = expressionParser.parseExpression(this.routingKeyExpression); this.routingKeyGenerator = new ExpressionEvaluatingMessageProcessor(expression, String.class); } @@ -78,6 +89,10 @@ public class AmqpOutboundEndpoint extends AbstractReplyProducingMessageHandler { this.exchangeName = exchangeName; } + public void setExchangeNameExpression(String exchangeNameExpression) { + this.exchangeNameExpression = exchangeNameExpression; + } + public void setRoutingKey(String routingKey) { this.routingKey = routingKey; } @@ -97,21 +112,25 @@ public class AmqpOutboundEndpoint extends AbstractReplyProducingMessageHandler { @Override protected Object handleRequestMessage(Message requestMessage) { + String exchangeName = this.exchangeName; String routingKey = this.routingKey; - if (this.routingKeyGenerator!=null) { + if (this.exchangeNameGenerator != null) { + exchangeName = this.exchangeNameGenerator.processMessage(requestMessage); + } + if (this.routingKeyGenerator != null) { routingKey = this.routingKeyGenerator.processMessage(requestMessage); } if (this.expectReply) { - return this.sendAndReceive(requestMessage, routingKey); + return this.sendAndReceive(exchangeName, routingKey, requestMessage); } else { - this.send(requestMessage, routingKey); + this.send(exchangeName, routingKey, requestMessage); return null; } } - private void send(final Message requestMessage, String routingKey) { - this.amqpTemplate.convertAndSend(this.exchangeName, routingKey, requestMessage.getPayload(), + private void send(String exchangeName, String routingKey, final Message requestMessage) { + this.amqpTemplate.convertAndSend(exchangeName, routingKey, requestMessage.getPayload(), new MessagePostProcessor() { public org.springframework.amqp.core.Message postProcessMessage( org.springframework.amqp.core.Message message) throws AmqpException { @@ -121,14 +140,14 @@ public class AmqpOutboundEndpoint extends AbstractReplyProducingMessageHandler { }); } - private Message sendAndReceive(Message requestMessage, String routingKey) { + private Message sendAndReceive(String exchangeName, String routingKey, Message requestMessage) { // TODO: add a convertSendAndReceive method that accepts a MessagePostProcessor so we can map headers? Assert.isTrue(amqpTemplate instanceof RabbitTemplate, "RabbitTemplate implementation is required for send and receive"); MessageConverter converter = ((RabbitTemplate) this.amqpTemplate).getMessageConverter(); MessageProperties amqpMessageProperties = new MessageProperties(); this.headerMapper.fromHeaders(requestMessage.getHeaders(), amqpMessageProperties); org.springframework.amqp.core.Message amqpMessage = converter.toMessage(requestMessage.getPayload(), amqpMessageProperties); - org.springframework.amqp.core.Message amqpReplyMessage = this.amqpTemplate.sendAndReceive(this.exchangeName, routingKey, amqpMessage); + org.springframework.amqp.core.Message amqpReplyMessage = this.amqpTemplate.sendAndReceive(exchangeName, routingKey, amqpMessage); if (amqpReplyMessage == null) { return null; }