From d43b89dedcd3247b211bc78e84b56bd2167f0b80 Mon Sep 17 00:00:00 2001 From: Artem Bilan Date: Fri, 1 Aug 2014 15:59:49 +0300 Subject: [PATCH] INT-2791 temp-channel-transacted for AMQP Channel JIRA: https://jira.spring.io/browse/INT-2791 * Add `template-channel-transacted` attribute for the ``s to separate configuration for container and RabbitTemplate * Make container's `channelTransacted` as `false` by default * Upgrade to `S-AMQP-1.4.0` * Add `AmqpHeaders.PUBLISH_CONFIRM_NACK_CAUSE` according to the new `RabbitTemplate.ConfirmCallback#confirm` signature Conflicts: src/reference/docbook/whats-new.xml Conflicts: src/reference/docbook/whats-new.xml INT-2791: Polishing Upgrade to S-AMQP-2.0 Fix for mock publisher confirm/return tests, since `channel` is now physically closed, if the `ConnectionFactory` isn't `CachingConnectionFactory` --- build.gradle | 2 +- .../integration/amqp/AmqpHeaders.java | 2 ++ .../amqp/config/AmqpChannelFactoryBean.java | 17 ++++++++------ .../amqp/config/AmqpChannelParser.java | 1 + .../amqp/outbound/AmqpOutboundEndpoint.java | 13 +++++++++-- .../config/spring-integration-amqp-4.1.xsd | 23 +++++++++++++------ .../config/AmqpChannelParserTests-context.xml | 3 ++- .../amqp/config/AmqpChannelParserTests.java | 4 ++++ ...AmqpOutboundChannelAdapterParserTests.java | 13 +++++++---- src/reference/docbook/amqp.xml | 16 ++++++++++++- src/reference/docbook/whats-new.xml | 8 +++++++ 11 files changed, 78 insertions(+), 24 deletions(-) diff --git a/build.gradle b/build.gradle index 9741434238..1f7abe9be8 100644 --- a/build.gradle +++ b/build.gradle @@ -112,7 +112,7 @@ subprojects { subproject -> slf4jVersion = "1.7.6" smack3Version = '3.2.1' smackVersion = '4.0.0' - springAmqpVersion = project.hasProperty('springAmqpVersion') ? project.springAmqpVersion : '1.3.5.RELEASE' + springAmqpVersion = project.hasProperty('springAmqpVersion') ? project.springAmqpVersion : '2.0.0.BUILD-SNAPSHOT' springDataMongoVersion = '1.5.0.RELEASE' springDataRedisVersion = '1.3.0.RELEASE' springGemfireVersion = '1.4.0.RELEASE' diff --git a/spring-integration-amqp/src/main/java/org/springframework/integration/amqp/AmqpHeaders.java b/spring-integration-amqp/src/main/java/org/springframework/integration/amqp/AmqpHeaders.java index 8a59c41c6e..3bffbdd9ff 100644 --- a/spring-integration-amqp/src/main/java/org/springframework/integration/amqp/AmqpHeaders.java +++ b/spring-integration-amqp/src/main/java/org/springframework/integration/amqp/AmqpHeaders.java @@ -79,6 +79,8 @@ public abstract class AmqpHeaders { public static final String PUBLISH_CONFIRM = PREFIX + "publishConfirm"; + public static final String PUBLISH_CONFIRM_NACK_CAUSE = PREFIX + "publishConfirmNackCause"; + public static final String RETURN_REPLY_CODE = PREFIX + "returnReplyCode"; public static final String RETURN_REPLY_TEXT = PREFIX + "returnReplyText"; diff --git a/spring-integration-amqp/src/main/java/org/springframework/integration/amqp/config/AmqpChannelFactoryBean.java b/spring-integration-amqp/src/main/java/org/springframework/integration/amqp/config/AmqpChannelFactoryBean.java index 3a1b79df8c..c1f21a5080 100644 --- a/spring-integration-amqp/src/main/java/org/springframework/integration/amqp/config/AmqpChannelFactoryBean.java +++ b/spring-integration-amqp/src/main/java/org/springframework/integration/amqp/config/AmqpChannelFactoryBean.java @@ -109,10 +109,7 @@ public class AmqpChannelFactoryBean extends AbstractFactoryBean headers = new HashMap(); + headers.put(AmqpHeaders.PUBLISH_CONFIRM, ack); + if (!ack && StringUtils.hasText(cause)) { + headers.put(AmqpHeaders.PUBLISH_CONFIRM_NACK_CAUSE, cause); + } + Message confirmMessage = this.getMessageBuilderFactory().withPayload(userCorrelationData) - .setHeader(AmqpHeaders.PUBLISH_CONFIRM, ack) + .copyHeaders(headers) .build(); if (ack && this.confirmAckChannel != null) { this.confirmAckChannel.send(confirmMessage); diff --git a/spring-integration-amqp/src/main/resources/org/springframework/integration/amqp/config/spring-integration-amqp-4.1.xsd b/spring-integration-amqp/src/main/resources/org/springframework/integration/amqp/config/spring-integration-amqp-4.1.xsd index d259a4a55a..788eb556a5 100644 --- a/spring-integration-amqp/src/main/resources/org/springframework/integration/amqp/config/spring-integration-amqp-4.1.xsd +++ b/spring-integration-amqp/src/main/resources/org/springframework/integration/amqp/config/spring-integration-amqp-4.1.xsd @@ -621,13 +621,6 @@ standard headers to also be mapped. To map all non-standard headers the 'NON_STA - - - - Flag to indicate that channels created by this component will be transactional. - - - @@ -663,6 +656,14 @@ standard headers to also be mapped. To map all non-standard headers the 'NON_STA SimpleMessageListenerContainer, such as channelTransacted, connectionFactory, and messagePropertiesConverter. + + + + Flag to indicate that channels created by this component will be transactional. + Only applies to messages sent to this channel, or when 'message-driven' is 'false'. + + + @@ -693,6 +694,14 @@ standard headers to also be mapped. To map all non-standard headers the 'NON_STA connectionFactory, and messsagePropertiesConverter. + + + + Flag to indicate that channels created by this component will be transactional. + Only applies to outbound messages when 'message-driven' is 'true'. + + + diff --git a/spring-integration-amqp/src/test/java/org/springframework/integration/amqp/config/AmqpChannelParserTests-context.xml b/spring-integration-amqp/src/test/java/org/springframework/integration/amqp/config/AmqpChannelParserTests-context.xml index 0855ea5a6e..0bfed1e2f0 100644 --- a/spring-integration-amqp/src/test/java/org/springframework/integration/amqp/config/AmqpChannelParserTests-context.xml +++ b/spring-integration-amqp/src/test/java/org/springframework/integration/amqp/config/AmqpChannelParserTests-context.xml @@ -17,7 +17,8 @@ - + diff --git a/spring-integration-amqp/src/test/java/org/springframework/integration/amqp/config/AmqpChannelParserTests.java b/spring-integration-amqp/src/test/java/org/springframework/integration/amqp/config/AmqpChannelParserTests.java index 4cd44986a7..3049c2fbe8 100644 --- a/spring-integration-amqp/src/test/java/org/springframework/integration/amqp/config/AmqpChannelParserTests.java +++ b/spring-integration-amqp/src/test/java/org/springframework/integration/amqp/config/AmqpChannelParserTests.java @@ -60,6 +60,8 @@ public class AmqpChannelParserTests { assertSame(mbf, TestUtils.getPropertyValue(channel, "dispatcher.messageBuilderFactory")); assertSame(mbf, TestUtils.getPropertyValue(channel, "container.messageListener.messageBuilderFactory")); assertTrue(TestUtils.getPropertyValue(channel, "container.missingQueuesFatal", Boolean.class)); + assertFalse(TestUtils.getPropertyValue(channel, "container.transactional", Boolean.class)); + assertFalse(TestUtils.getPropertyValue(channel, "amqpTemplate.transactional", Boolean.class)); } @Test @@ -68,6 +70,8 @@ public class AmqpChannelParserTests { assertEquals(1, TestUtils.getPropertyValue( TestUtils.getPropertyValue(channel, "dispatcher"), "maxSubscribers", Integer.class).intValue()); assertFalse(TestUtils.getPropertyValue(channel, "container.missingQueuesFatal", Boolean.class)); + assertFalse(TestUtils.getPropertyValue(channel, "container.transactional", Boolean.class)); + assertTrue(TestUtils.getPropertyValue(channel, "amqpTemplate.transactional", Boolean.class)); } diff --git a/spring-integration-amqp/src/test/java/org/springframework/integration/amqp/config/AmqpOutboundChannelAdapterParserTests.java b/spring-integration-amqp/src/test/java/org/springframework/integration/amqp/config/AmqpOutboundChannelAdapterParserTests.java index b21004892f..a192309c6f 100644 --- a/spring-integration-amqp/src/test/java/org/springframework/integration/amqp/config/AmqpOutboundChannelAdapterParserTests.java +++ b/spring-integration-amqp/src/test/java/org/springframework/integration/amqp/config/AmqpOutboundChannelAdapterParserTests.java @@ -24,6 +24,7 @@ import static org.junit.Assert.assertSame; import static org.junit.Assert.assertTrue; import static org.junit.Assert.fail; import static org.mockito.Matchers.any; +import static org.mockito.Mockito.doAnswer; import static org.mockito.Mockito.doThrow; import static org.mockito.Mockito.mock; import static org.mockito.Mockito.never; @@ -38,11 +39,14 @@ import java.util.List; import java.util.Map; import java.util.concurrent.atomic.AtomicBoolean; +import com.rabbitmq.client.AMQP.BasicProperties; +import com.rabbitmq.client.Channel; import org.apache.commons.logging.Log; import org.junit.Test; import org.junit.runner.RunWith; import org.mockito.Matchers; import org.mockito.Mockito; +import org.mockito.internal.stubbing.answers.DoesNothing; import org.mockito.invocation.InvocationOnMock; import org.mockito.stubbing.Answer; @@ -82,9 +86,6 @@ import org.springframework.test.context.ContextConfiguration; import org.springframework.test.context.junit4.SpringJUnit4ClassRunner; import org.springframework.util.ReflectionUtils; -import com.rabbitmq.client.AMQP.BasicProperties; -import com.rabbitmq.client.Channel; - /** * @author Mark Fisher * @author Oleg Zhurakousky @@ -186,7 +187,8 @@ public class AmqpOutboundChannelAdapterParserTests { Channel mockChannel = mock(Channel.class); when(connectionFactory.createConnection()).thenReturn(mockConnection); - PublisherCallbackChannelImpl publisherCallbackChannel = new PublisherCallbackChannelImpl(mockChannel); + PublisherCallbackChannelImpl publisherCallbackChannel = spy(new PublisherCallbackChannelImpl(mockChannel)); + doAnswer(new DoesNothing()).when(publisherCallbackChannel).close(); when(mockConnection.createChannel(false)).thenReturn(publisherCallbackChannel); MessageChannel requestChannel = context.getBean("pcRequestChannel", MessageChannel.class); @@ -246,7 +248,8 @@ public class AmqpOutboundChannelAdapterParserTests { Channel mockChannel = mock(Channel.class); when(connectionFactory.createConnection()).thenReturn(mockConnection); - PublisherCallbackChannelImpl publisherCallbackChannel = new PublisherCallbackChannelImpl(mockChannel); + PublisherCallbackChannelImpl publisherCallbackChannel = spy(new PublisherCallbackChannelImpl(mockChannel)); + doAnswer(new DoesNothing()).when(publisherCallbackChannel).close(); when(mockConnection.createChannel(false)).thenReturn(publisherCallbackChannel); MessageChannel requestChannel = context.getBean("returnRequestChannel", MessageChannel.class); diff --git a/src/reference/docbook/amqp.xml b/src/reference/docbook/amqp.xml index f63c71d7b0..5ca9da1b29 100644 --- a/src/reference/docbook/amqp.xml +++ b/src/reference/docbook/amqp.xml @@ -489,6 +489,11 @@ public Object handle(@Payload String payload, @Header(AmqpHeaders.CHANNEL) Chann defined by this expression and the message will have a header 'amqp_publishConfirm' set to true (ack) or false (nack). Examples: "headers['myCorrelationData']", "payload". Optional. + + Starting with version 4.1 the new amqp_publishConfirmNackCause + message header has been added. It contains a cause message of 'nack' for publisher + confirm. + The channel to which positive (ack) publisher confirms are sent; payload is @@ -671,7 +676,7 @@ public Object handle(@Payload String payload, @Header(AmqpHeaders.CHANNEL) Chann -
+
AMQP Backed Message Channels @@ -697,6 +702,14 @@ public Object handle(@Payload String payload, @Header(AmqpHeaders.CHANNEL) Chann and bind that to the Fanout Exchange while registering a consumer on that Queue to receive Messages. There is no "pollable" option for a publish-subscribe-channel; it must be message-driven. + + Starting with version 4.1 AMQP Backed Message Channels, alongside with + channel-transacted, support template-channel-transacted to separate + transactional configuration for the AbstractMessageListenerContainer + and for the RabbitTemplate. + Note, previously, the channel-transacted was true by default, now it changed to + false as standard default value for the AbstractMessageListenerContainer. +
AMQP Message Headers @@ -760,6 +773,7 @@ public Object handle(@Payload String payload, @Header(AmqpHeaders.CHANNEL) Chann amqp_springReplyCorrelation amqp_springReplyToStack amqp_publishConfirm + amqp_publishConfirmNackCause amqp_returnReplyCode amqp_returnReplyText amqp_returnExchange diff --git a/src/reference/docbook/whats-new.xml b/src/reference/docbook/whats-new.xml index e15ec49615..0aa753ddb2 100644 --- a/src/reference/docbook/whats-new.xml +++ b/src/reference/docbook/whats-new.xml @@ -154,5 +154,13 @@ See for more information.
+
+ AMQP Channels: template-channel-transacted + + The new template-channel-transacted attribute has been introduced for AMQP + MessageChannels. + See for more information. + +