INT-1435 added support for runtime priority from SI Message header on JMS outbound gateway

This commit is contained in:
Mark Fisher
2011-02-02 19:40:50 -05:00
parent da6174c8f0
commit 0cb10b5942
3 changed files with 81 additions and 23 deletions

View File

@@ -370,15 +370,19 @@ public class JmsOutboundGateway extends AbstractReplyProducingMessageHandler {
jmsRequest.setJMSReplyTo(replyTo);
connection.start();
Integer priority = requestMessage.getHeaders().getPriority();
if (priority == null) {
priority = this.priority;
}
javax.jms.Message replyMessage = null;
if (this.correlationKey != null) {
replyMessage = this.doSendAndReceiveWithGeneratedCorrelationId(jmsRequest, replyTo, session);
replyMessage = this.doSendAndReceiveWithGeneratedCorrelationId(jmsRequest, replyTo, session, priority);
}
else if (replyTo instanceof TemporaryQueue || replyTo instanceof TemporaryTopic) {
replyMessage = this.doSendAndReceiveWithTemporaryReplyToDestination(jmsRequest, replyTo, session);
replyMessage = this.doSendAndReceiveWithTemporaryReplyToDestination(jmsRequest, replyTo, session, priority);
}
else {
replyMessage = this.doSendAndReceiveWithMessageIdCorrelation(jmsRequest, replyTo, session);
replyMessage = this.doSendAndReceiveWithMessageIdCorrelation(jmsRequest, replyTo, session, priority);
}
return replyMessage;
}
@@ -392,7 +396,7 @@ public class JmsOutboundGateway extends AbstractReplyProducingMessageHandler {
/**
* Creates the MessageConsumer before sending the request Message since we are generating our own correlationId value for the MessageSelector.
*/
private javax.jms.Message doSendAndReceiveWithGeneratedCorrelationId(javax.jms.Message jmsRequest, Destination replyTo, Session session) throws JMSException {
private javax.jms.Message doSendAndReceiveWithGeneratedCorrelationId(javax.jms.Message jmsRequest, Destination replyTo, Session session, int priority) throws JMSException {
MessageProducer messageProducer = null;
MessageConsumer messageConsumer = null;
try {
@@ -409,7 +413,7 @@ public class JmsOutboundGateway extends AbstractReplyProducingMessageHandler {
messageSelector = this.correlationKey + " = '" + correlationId + "'";
}
messageConsumer = session.createConsumer(replyTo, messageSelector);
this.sendRequestMessage(jmsRequest, messageProducer);
this.sendRequestMessage(jmsRequest, messageProducer, priority);
return this.receiveReplyMessage(messageConsumer);
}
finally {
@@ -421,13 +425,13 @@ public class JmsOutboundGateway extends AbstractReplyProducingMessageHandler {
/**
* Creates the MessageConsumer before sending the request Message since we do not need any correlation.
*/
private javax.jms.Message doSendAndReceiveWithTemporaryReplyToDestination(javax.jms.Message jmsRequest, Destination replyTo, Session session) throws JMSException {
private javax.jms.Message doSendAndReceiveWithTemporaryReplyToDestination(javax.jms.Message jmsRequest, Destination replyTo, Session session, int priority) throws JMSException {
MessageProducer messageProducer = null;
MessageConsumer messageConsumer = null;
try {
messageProducer = session.createProducer(this.getRequestDestination(session));
messageConsumer = session.createConsumer(replyTo);
this.sendRequestMessage(jmsRequest, messageProducer);
this.sendRequestMessage(jmsRequest, messageProducer, priority);
return this.receiveReplyMessage(messageConsumer);
}
finally {
@@ -439,7 +443,7 @@ public class JmsOutboundGateway extends AbstractReplyProducingMessageHandler {
/**
* Creates the MessageConsumer after sending the request Message since we need the MessageID for correlation with a MessageSelector.
*/
private javax.jms.Message doSendAndReceiveWithMessageIdCorrelation(javax.jms.Message jmsRequest, Destination replyTo, Session session) throws JMSException {
private javax.jms.Message doSendAndReceiveWithMessageIdCorrelation(javax.jms.Message jmsRequest, Destination replyTo, Session session, int priority) throws JMSException {
if (replyTo instanceof Topic && logger.isWarnEnabled()) {
logger.warn("Relying on the MessageID for correlation is not recommended when using a Topic as the replyTo Destination " +
"because that ID can only be provided to a MessageSelector after the reuqest Message has been sent thereby " +
@@ -451,7 +455,7 @@ public class JmsOutboundGateway extends AbstractReplyProducingMessageHandler {
MessageConsumer messageConsumer = null;
try {
messageProducer = session.createProducer(this.getRequestDestination(session));
this.sendRequestMessage(jmsRequest, messageProducer);
this.sendRequestMessage(jmsRequest, messageProducer, priority);
String messageId = jmsRequest.getJMSMessageID().replaceAll("'", "''");
String messageSelector = "JMSCorrelationID = '" + messageId + "'";
messageConsumer = session.createConsumer(replyTo, messageSelector);
@@ -463,9 +467,9 @@ public class JmsOutboundGateway extends AbstractReplyProducingMessageHandler {
}
}
private void sendRequestMessage(javax.jms.Message jmsRequest, MessageProducer messageProducer) throws JMSException {
private void sendRequestMessage(javax.jms.Message jmsRequest, MessageProducer messageProducer, int priority) throws JMSException {
if (this.explicitQosEnabled) {
messageProducer.send(jmsRequest, this.deliveryMode, this.priority, this.timeToLive);
messageProducer.send(jmsRequest, this.deliveryMode, priority, this.timeToLive);
}
else {
messageProducer.send(jmsRequest);

View File

@@ -1,17 +1,23 @@
<?xml version="1.0" encoding="UTF-8"?>
<beans xmlns="http://www.springframework.org/schema/beans"
xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance"
xmlns:jms="http://www.springframework.org/schema/jms"
xmlns:int="http://www.springframework.org/schema/integration"
xmlns:int-jms="http://www.springframework.org/schema/integration/jms"
xsi:schemaLocation="http://www.springframework.org/schema/beans http://www.springframework.org/schema/beans/spring-beans-3.0.xsd
http://www.springframework.org/schema/jms http://www.springframework.org/schema/jms/spring-jms-3.0.xsd
http://www.springframework.org/schema/integration http://www.springframework.org/schema/integration/spring-integration-2.0.xsd
http://www.springframework.org/schema/integration/jms http://www.springframework.org/schema/integration/jms/spring-integration-jms-2.0.xsd">
<int-jms:outbound-channel-adapter id="outbound" destination-name="queue.test.priority" priority="3" explicit-qos-enabled="true"/>
<int-jms:outbound-channel-adapter id="channelAdapterChannel" destination-name="queue.test.priority.channelAdapter" priority="3" explicit-qos-enabled="true"/>
<int-jms:message-driven-channel-adapter channel="results" destination-name="queue.test.priority" extract-payload="false"/>
<int:channel id="results">
<int-jms:message-driven-channel-adapter channel="channelAdapterResults" destination-name="queue.test.priority.channelAdapter" extract-payload="false"/>
<int:channel id="gatewayChannel"/>
<int-jms:outbound-gateway request-channel="gatewayChannel" request-destination-name="queue.test.priority.gateway" priority="2" explicit-qos-enabled="true"/>
<int:channel id="channelAdapterResults">
<int:queue capacity="2"/>
</int:channel>
@@ -25,4 +31,10 @@
<property name="cacheProducers" value="false"/>
</bean>
<jms:listener-container>
<jms:listener destination="queue.test.priority.gateway" ref="priorityReader"/>
</jms:listener-container>
<bean id="priorityReader" class="org.springframework.integration.jms.config.JmsPriorityTests$PriorityReader"/>
</beans>

View File

@@ -20,6 +20,11 @@ import static org.junit.Assert.assertEquals;
import static org.junit.Assert.assertNotNull;
import static org.junit.Assert.assertTrue;
import javax.jms.JMSException;
import javax.jms.MessageProducer;
import javax.jms.Session;
import javax.jms.TextMessage;
import org.junit.Before;
import org.junit.Test;
import org.junit.runner.RunWith;
@@ -27,8 +32,10 @@ import org.junit.runner.RunWith;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.integration.Message;
import org.springframework.integration.MessageChannel;
import org.springframework.integration.channel.QueueChannel;
import org.springframework.integration.core.PollableChannel;
import org.springframework.integration.support.MessageBuilder;
import org.springframework.jms.listener.SessionAwareMessageListener;
import org.springframework.test.context.ContextConfiguration;
import org.springframework.test.context.junit4.SpringJUnit4ClassRunner;
@@ -41,10 +48,14 @@ import org.springframework.test.context.junit4.SpringJUnit4ClassRunner;
public class JmsPriorityTests {
@Autowired
private MessageChannel outbound;
private MessageChannel channelAdapterChannel;
@Autowired
private PollableChannel results;
private PollableChannel channelAdapterResults;
@Autowired
private MessageChannel gatewayChannel;
@Before
public void prepareActiveMq() {
@@ -52,10 +63,10 @@ public class JmsPriorityTests {
}
@Test
public void verifyPrioritySettingOnAdapterUsedAsJmsPriorityIfNoHeader() throws Exception {
public void verifyPrioritySettingOnChannelAdapterUsedAsJmsPriorityIfNoHeader() throws Exception {
Message<?> message = MessageBuilder.withPayload("test").build();
outbound.send(message);
Message<?> result = results.receive(5000);
channelAdapterChannel.send(message);
Message<?> result = channelAdapterResults.receive(5000);
assertNotNull(result);
assertTrue(result.getPayload() instanceof javax.jms.Message);
javax.jms.Message jmsMessage = (javax.jms.Message) result.getPayload();
@@ -63,14 +74,45 @@ public class JmsPriorityTests {
}
@Test
public void verifyPriorityHeaderUsedAsJmsPriority() throws Exception {
public void verifyPriorityHeaderUsedAsJmsPriorityWithChannelAdapter() throws Exception {
Message<?> message = MessageBuilder.withPayload("test").setPriority(7).build();
outbound.send(message);
Message<?> result = results.receive(5000);
channelAdapterChannel.send(message);
Message<?> result = channelAdapterResults.receive(5000);
assertNotNull(result);
assertTrue(result.getPayload() instanceof javax.jms.Message);
javax.jms.Message jmsMessage = (javax.jms.Message) result.getPayload();
assertEquals(7, jmsMessage.getJMSPriority());
}
@Test
public void verifyPrioritySettingOnGatewayUsedAsJmsPriorityIfNoHeader() throws Exception {
QueueChannel replyChannel = new QueueChannel();
Message<?> message = MessageBuilder.withPayload("test").setReplyChannel(replyChannel).build();
gatewayChannel.send(message);
Message<?> result = replyChannel.receive(5000);
assertNotNull(result);
assertEquals("priority=2", result.getPayload());
}
@Test
public void verifyPriorityHeaderUsedAsJmsPriorityWithGateway() throws Exception {
QueueChannel replyChannel = new QueueChannel();
Message<?> message = MessageBuilder.withPayload("test").setPriority(8).setReplyChannel(replyChannel).build();
gatewayChannel.send(message);
Message<?> result = replyChannel.receive(5000);
assertNotNull(result);
assertEquals("priority=8", result.getPayload());
}
public static class PriorityReader implements SessionAwareMessageListener<javax.jms.Message> {
public void onMessage(javax.jms.Message request, Session session) throws JMSException {
String text = "priority=" + request.getJMSPriority();
TextMessage reply = session.createTextMessage(text);
MessageProducer producer = session.createProducer(request.getJMSReplyTo());
producer.send(reply);
}
}
}