From dbf4c2122f64f8df11632aa33957d655d15d0ebd Mon Sep 17 00:00:00 2001 From: Mark Fisher Date: Fri, 9 Sep 2011 20:36:02 +0100 Subject: [PATCH 1/3] changed default connection-factory value to 'rabbitConnectionFactory' to be consistent with the Spring AMQP listener-container --- .../amqp/config/AbstractAmqpInboundAdapterParser.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/spring-integration-amqp/src/main/java/org/springframework/integration/amqp/config/AbstractAmqpInboundAdapterParser.java b/spring-integration-amqp/src/main/java/org/springframework/integration/amqp/config/AbstractAmqpInboundAdapterParser.java index 3ba3b0dd4d..ccf9691bb9 100644 --- a/spring-integration-amqp/src/main/java/org/springframework/integration/amqp/config/AbstractAmqpInboundAdapterParser.java +++ b/spring-integration-amqp/src/main/java/org/springframework/integration/amqp/config/AbstractAmqpInboundAdapterParser.java @@ -111,7 +111,7 @@ abstract class AbstractAmqpInboundAdapterParser extends AbstractSingleBeanDefini "org.springframework.amqp.rabbit.listener.SimpleMessageListenerContainer"); String connectionFactoryRef = element.getAttribute("connection-factory"); if (!StringUtils.hasText(connectionFactoryRef)) { - connectionFactoryRef = "connectionFactory"; + connectionFactoryRef = "rabbitConnectionFactory"; } builder.addConstructorArgReference(connectionFactoryRef); for (String attributeName : CONTAINER_VALUE_ATTRIBUTES) { From ff4a979254fa9061562a4d91d86a4004930ff386 Mon Sep 17 00:00:00 2001 From: Mark Fisher Date: Fri, 9 Sep 2011 20:46:49 +0100 Subject: [PATCH 2/3] 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; } From 84375120ad302fb2d7bad78796c607df1fab56c9 Mon Sep 17 00:00:00 2001 From: Mark Fisher Date: Fri, 9 Sep 2011 21:00:56 +0100 Subject: [PATCH 3/3] added 'exchange-name-expression' attributes to the schema --- .../amqp/config/spring-integration-amqp-2.1.xsd | 16 +++++++++++++++- 1 file changed, 15 insertions(+), 1 deletion(-) diff --git a/spring-integration-amqp/src/main/resources/org/springframework/integration/amqp/config/spring-integration-amqp-2.1.xsd b/spring-integration-amqp/src/main/resources/org/springframework/integration/amqp/config/spring-integration-amqp-2.1.xsd index c2f08a2d8d..0833fce007 100644 --- a/spring-integration-amqp/src/main/resources/org/springframework/integration/amqp/config/spring-integration-amqp-2.1.xsd +++ b/spring-integration-amqp/src/main/resources/org/springframework/integration/amqp/config/spring-integration-amqp-2.1.xsd @@ -43,7 +43,14 @@ - The name of the AMQP Exchange to which Messages should be sent. If not provided, Messages will be sent to the default, no-name Exchange. + The fixed name of the AMQP Exchange to which Messages should be sent. If not provided, Messages will be sent to the default, no-name Exchange. + + + + + + + The exchange name to use when sending Messages evaluated as an expression on the message (e.g. 'headers.exchange'). By default, this will be an emtpy String. @@ -156,6 +163,13 @@ + + + + The exchange name to use when sending Messages evaluated as an expression on the message (e.g. 'headers.exchange'). By default, this will be an emtpy String. + + +