diff --git a/spring-rabbit/src/main/java/org/springframework/amqp/rabbit/connection/ConnectionFactoryUtils.java b/spring-rabbit/src/main/java/org/springframework/amqp/rabbit/connection/ConnectionFactoryUtils.java index 03ec2551..e0a549f2 100644 --- a/spring-rabbit/src/main/java/org/springframework/amqp/rabbit/connection/ConnectionFactoryUtils.java +++ b/spring-rabbit/src/main/java/org/springframework/amqp/rabbit/connection/ConnectionFactoryUtils.java @@ -1,5 +1,5 @@ /* - * Copyright 2002-2016 the original author or authors. + * Copyright 2002-2017 the original author or authors. * * Licensed under the Apache License, Version 2.0 (the "License"); * you may not use this file except in compliance with the License. @@ -167,11 +167,11 @@ public final class ConnectionFactoryUtils { RabbitUtils.closeConnection(resourceHolder.getConnection()); } - public static void bindResourceToTransaction(RabbitResourceHolder resourceHolder, + public static RabbitResourceHolder bindResourceToTransaction(RabbitResourceHolder resourceHolder, ConnectionFactory connectionFactory, boolean synched) { if (TransactionSynchronizationManager.hasResource(connectionFactory) || !TransactionSynchronizationManager.isActualTransactionActive() || !synched) { - return; + return (RabbitResourceHolder) TransactionSynchronizationManager.getResource(connectionFactory); } TransactionSynchronizationManager.bindResource(connectionFactory, resourceHolder); resourceHolder.setSynchronizedWithTransaction(true); @@ -179,6 +179,7 @@ public final class ConnectionFactoryUtils { TransactionSynchronizationManager.registerSynchronization(new RabbitResourceSynchronization(resourceHolder, connectionFactory)); } + return resourceHolder; } public static void registerDeliveryTag(ConnectionFactory connectionFactory, Channel channel, Long tag) { diff --git a/spring-rabbit/src/main/java/org/springframework/amqp/rabbit/connection/RabbitResourceHolder.java b/spring-rabbit/src/main/java/org/springframework/amqp/rabbit/connection/RabbitResourceHolder.java index d74f2ca1..6ff8b811 100644 --- a/spring-rabbit/src/main/java/org/springframework/amqp/rabbit/connection/RabbitResourceHolder.java +++ b/spring-rabbit/src/main/java/org/springframework/amqp/rabbit/connection/RabbitResourceHolder.java @@ -1,5 +1,5 @@ /* - * Copyright 2002-2016 the original author or authors. + * Copyright 2002-2017 the original author or authors. * * Licensed under the Apache License, Version 2.0 (the "License"); * you may not use this file except in compliance with the License. @@ -64,22 +64,23 @@ public class RabbitResourceHolder extends ResourceHolderSupport { private final MultiValueMap deliveryTags = new LinkedMultiValueMap(); - private boolean transactional; + private final boolean releaseAfterCompletion; - private boolean releaseAfterCompletion = true; + private boolean requeueOnRollback = true; // No need for volatile written/read on the same thread. /** * Create a new RabbitResourceHolder that is open for resources to be added. */ public RabbitResourceHolder() { + this.releaseAfterCompletion = true; } /** + * Construct an instance for the channel. * @param channel a channel to add * @param releaseAfterCompletion true if the channel should be released after completion. */ public RabbitResourceHolder(Channel channel, boolean releaseAfterCompletion) { - this(); addChannel(channel); this.releaseAfterCompletion = releaseAfterCompletion; } @@ -99,6 +100,15 @@ public class RabbitResourceHolder extends ResourceHolderSupport { return this.releaseAfterCompletion; } + /** + * Set to true to requeue a message on rollback; default true. + * @param requeueOnRollback true to requeue + * @since 1.7.1 + */ + public void setRequeueOnRollback(boolean requeueOnRollback) { + this.requeueOnRollback = requeueOnRollback; + } + public final void addConnection(Connection connection) { Assert.isTrue(!this.frozen, "Cannot add Connection because RabbitResourceHolder is frozen"); Assert.notNull(connection, "Connection must not be null"); @@ -196,7 +206,7 @@ public class RabbitResourceHolder extends ResourceHolderSupport { if (this.deliveryTags.containsKey(channel)) { for (Long deliveryTag : this.deliveryTags.get(channel)) { try { - channel.basicReject(deliveryTag, true); + channel.basicReject(deliveryTag, this.requeueOnRollback); } catch (IOException ex) { throw new AmqpIOException(ex); @@ -209,10 +219,13 @@ public class RabbitResourceHolder extends ResourceHolderSupport { } /** + * Invalid - always returned false. * @return true if the channels in this holder are transactional + * @deprecated Not used */ + @Deprecated public boolean isChannelTransactional() { - return this.transactional; + return false; } } diff --git a/spring-rabbit/src/main/java/org/springframework/amqp/rabbit/listener/AbstractMessageListenerContainer.java b/spring-rabbit/src/main/java/org/springframework/amqp/rabbit/listener/AbstractMessageListenerContainer.java index 72f8ce94..914e009a 100644 --- a/spring-rabbit/src/main/java/org/springframework/amqp/rabbit/listener/AbstractMessageListenerContainer.java +++ b/spring-rabbit/src/main/java/org/springframework/amqp/rabbit/listener/AbstractMessageListenerContainer.java @@ -195,6 +195,7 @@ public abstract class AbstractMessageListenerContainer extends RabbitAccessor private ConditionalExceptionLogger exclusiveConsumerExceptionLogger = new DefaultExclusiveConsumerLogger(); + private boolean alwaysRequeueWithTxManagerRollback; /** * {@inheritDoc} @@ -878,6 +879,28 @@ public abstract class AbstractMessageListenerContainer extends RabbitAccessor return this.exclusiveConsumerExceptionLogger; } + /** + * Set to true to always requeue on transaction rollback with an external + * {@link #setTransactionManager(PlatformTransactionManager) TransactionManager}. + * With earlier releases, when a transaction manager was configured, a transaction + * rollback always requeued the message. This was inconsistent with local transactions + * where the normal {@link #setDefaultRequeueRejected(boolean) defaultRequeueRejected} + * and {@link AmqpRejectAndDontRequeueException} logic was honored to determine whether + * the message was requeued. RabbitMQ does not consider the message delivery to be part + * of the transaction. + * This boolean was introduced in 1.7.1, set to true by default, to be consistent with + * previous behavior. Starting with version 2.0, it is false by default. + * @param alwaysRequeueWithTxManagerRollback true to always requeue on rollback. + * @since 1.7.1. + */ + public void setAlwaysRequeueWithTxManagerRollback(boolean alwaysRequeueWithTxManagerRollback) { + this.alwaysRequeueWithTxManagerRollback = alwaysRequeueWithTxManagerRollback; + } + + protected boolean isAlwaysRequeueWithTxManagerRollback() { + return this.alwaysRequeueWithTxManagerRollback; + } + /** * Delegates to {@link #validateConfiguration()} and {@link #initialize()}. */ @@ -1519,6 +1542,20 @@ public abstract class AbstractMessageListenerContainer extends RabbitAccessor return false; } + /** + * A null resource holder is rare, but possible if the transaction attribute caused no + * transaction to be started (e.g. {@code TransactionDefinition.PROPAGATION_NONE}). In + * that case the delivery tags will have been processed manually. + * @param resourceHolder the bound resource holder (if a transaction is active). + * @param exception the exception. + */ + protected void prepareHolderForRollback(RabbitResourceHolder resourceHolder, RuntimeException exception) { + if (resourceHolder != null) { + resourceHolder.setRequeueOnRollback(isAlwaysRequeueWithTxManagerRollback() || + RabbitUtils.shouldRequeue(isDefaultRequeueRejected(), exception, logger)); + } + } + @FunctionalInterface private interface ContainerDelegate { diff --git a/spring-rabbit/src/main/java/org/springframework/amqp/rabbit/listener/BlockingQueueConsumer.java b/spring-rabbit/src/main/java/org/springframework/amqp/rabbit/listener/BlockingQueueConsumer.java index 61da0727..1b2857fe 100644 --- a/spring-rabbit/src/main/java/org/springframework/amqp/rabbit/listener/BlockingQueueConsumer.java +++ b/spring-rabbit/src/main/java/org/springframework/amqp/rabbit/listener/BlockingQueueConsumer.java @@ -55,6 +55,7 @@ import org.springframework.amqp.rabbit.support.Delivery; import org.springframework.amqp.rabbit.support.MessagePropertiesConverter; import org.springframework.amqp.rabbit.support.RabbitExceptionTranslator; import org.springframework.amqp.support.ConsumerTagStrategy; +import org.springframework.transaction.support.TransactionSynchronizationManager; import org.springframework.util.ObjectUtils; import org.springframework.util.backoff.BackOffExecution; @@ -706,18 +707,24 @@ public class BlockingQueueConsumer { return false; } + /* + * If we have a TX Manager, but no TX, act like we are locally transacted. + */ + boolean isLocallyTransacted = locallyTransacted + || (this.transactional + && TransactionSynchronizationManager.getResource(this.connectionFactory) == null); try { boolean ackRequired = !this.acknowledgeMode.isAutoAck() && !this.acknowledgeMode.isManual(); if (ackRequired) { - if (!this.transactional || locallyTransacted) { + if (!this.transactional || isLocallyTransacted) { long deliveryTag = new ArrayList(this.deliveryTags).get(this.deliveryTags.size() - 1); this.channel.basicAck(deliveryTag, true); } } - if (locallyTransacted) { + if (isLocallyTransacted) { // For manual acks we still need to commit RabbitUtils.commitIfNecessary(this.channel); } diff --git a/spring-rabbit/src/main/java/org/springframework/amqp/rabbit/listener/DirectMessageListenerContainer.java b/spring-rabbit/src/main/java/org/springframework/amqp/rabbit/listener/DirectMessageListenerContainer.java index cbf84707..d879b3fe 100644 --- a/spring-rabbit/src/main/java/org/springframework/amqp/rabbit/listener/DirectMessageListenerContainer.java +++ b/spring-rabbit/src/main/java/org/springframework/amqp/rabbit/listener/DirectMessageListenerContainer.java @@ -1,5 +1,5 @@ /* - * Copyright 2016 the original author or authors. + * Copyright 2016-2017 the original author or authors. * * Licensed under the Apache License, Version 2.0 (the "License"); * you may not use this file except in compliance with the License. @@ -53,6 +53,7 @@ import org.springframework.scheduling.TaskScheduler; import org.springframework.scheduling.concurrent.ThreadPoolTaskScheduler; import org.springframework.transaction.PlatformTransactionManager; import org.springframework.transaction.interceptor.TransactionAttribute; +import org.springframework.transaction.support.TransactionSynchronizationManager; import org.springframework.transaction.support.TransactionTemplate; import org.springframework.util.Assert; import org.springframework.util.CollectionUtils; @@ -698,15 +699,17 @@ public class DirectMessageListenerContainer extends AbstractMessageListenerConta new TransactionTemplate(this.transactionManager, this.transactionAttribute); } this.transactionTemplate.execute(s -> { - RabbitResourceHolder resourceHolder = new RabbitResourceHolder(getChannel(), false); - resourceHolder.addDeliveryTag(getChannel(), deliveryTag); - ConnectionFactoryUtils.bindResourceToTransaction(resourceHolder, - this.connectionFactory, true); + RabbitResourceHolder resourceHolder = ConnectionFactoryUtils.bindResourceToTransaction( + new RabbitResourceHolder(getChannel(), false), this.connectionFactory, true); + if (resourceHolder != null) { + resourceHolder.addDeliveryTag(getChannel(), deliveryTag); + } // unbound in ResourceHolderSynchronization.beforeCompletion() try { callExecuteListener(message, deliveryTag); } catch (RuntimeException e1) { + prepareHolderForRollback(resourceHolder, e1); throw e1; } catch (Throwable e2) { //NOSONAR ok to catch Throwable here because we re-throw it below @@ -761,6 +764,15 @@ public class DirectMessageListenerContainer extends AbstractMessageListenerConta this.logger.error("Failed to invoke listener", e); if (this.transactionManager != null) { if (this.transactionAttribute.rollbackOn(e)) { + RabbitResourceHolder resourceHolder = (RabbitResourceHolder) TransactionSynchronizationManager + .getResource(getConnectionFactory()); + if (resourceHolder == null) { + /* + * If we don't actually have a transaction, we have to roll back + * manually. See prepareHolderForRollback(). + */ + rollback(deliveryTag, e); + } throw e; // encompassing transaction will handle the rollback. } else { @@ -778,13 +790,19 @@ public class DirectMessageListenerContainer extends AbstractMessageListenerConta } private void handleAck(long deliveryTag, boolean channelLocallyTransacted) throws IOException { + /* + * If we have a TX Manager, but no TX, act like we are locally transacted. + */ + boolean isLocallyTransacted = channelLocallyTransacted + || (isChannelTransacted() + && TransactionSynchronizationManager.getResource(this.connectionFactory) == null); try { if (this.ackRequired) { - if (!isChannelTransacted() || channelLocallyTransacted) { + if (!isChannelTransacted() || isLocallyTransacted) { getChannel().basicAck(deliveryTag, false); } } - if (channelLocallyTransacted) { + if (isLocallyTransacted) { RabbitUtils.commitIfNecessary(getChannel()); } } diff --git a/spring-rabbit/src/main/java/org/springframework/amqp/rabbit/listener/SimpleMessageListenerContainer.java b/spring-rabbit/src/main/java/org/springframework/amqp/rabbit/listener/SimpleMessageListenerContainer.java index 948c0ae5..5b56f7ce 100644 --- a/spring-rabbit/src/main/java/org/springframework/amqp/rabbit/listener/SimpleMessageListenerContainer.java +++ b/spring-rabbit/src/main/java/org/springframework/amqp/rabbit/listener/SimpleMessageListenerContainer.java @@ -51,6 +51,7 @@ import org.springframework.beans.BeansException; import org.springframework.context.ApplicationContext; import org.springframework.jmx.export.annotation.ManagedMetric; import org.springframework.jmx.support.MetricType; +import org.springframework.transaction.support.TransactionSynchronizationManager; import org.springframework.transaction.support.TransactionTemplate; import org.springframework.util.Assert; import org.springframework.util.StringUtils; @@ -727,7 +728,7 @@ public class SimpleMessageListenerContainer extends AbstractMessageListenerConta } return this.transactionTemplate .execute(status -> { - ConnectionFactoryUtils.bindResourceToTransaction( + RabbitResourceHolder resourceHolder = ConnectionFactoryUtils.bindResourceToTransaction( new RabbitResourceHolder(consumer.getChannel(), false), getConnectionFactory(), true); // unbound in ResourceHolderSynchronization.beforeCompletion() @@ -735,6 +736,7 @@ public class SimpleMessageListenerContainer extends AbstractMessageListenerConta return doReceiveAndExecute(consumer); } catch (RuntimeException e1) { + prepareHolderForRollback(resourceHolder, e1); throw e1; } catch (Throwable e2) { //NOSONAR @@ -783,7 +785,18 @@ public class SimpleMessageListenerContainer extends AbstractMessageListenerConta } if (getTransactionManager() != null) { if (getTransactionAttribute().rollbackOn(ex)) { - consumer.clearDeliveryTags(); + RabbitResourceHolder resourceHolder = (RabbitResourceHolder) TransactionSynchronizationManager + .getResource(getConnectionFactory()); + if (resourceHolder != null) { + consumer.clearDeliveryTags(); + } + else { + /* + * If we don't actually have a transaction, we have to roll back + * manually. See prepareHolderForRollback(). + */ + consumer.rollbackOnExceptionIfNecessary(ex); + } throw ex; // encompassing transaction will handle the rollback. } else { diff --git a/spring-rabbit/src/test/java/org/springframework/amqp/rabbit/connection/ConnectionFactoryUtilsTests.java b/spring-rabbit/src/test/java/org/springframework/amqp/rabbit/connection/ConnectionFactoryUtilsTests.java new file mode 100644 index 00000000..9a63b088 --- /dev/null +++ b/spring-rabbit/src/test/java/org/springframework/amqp/rabbit/connection/ConnectionFactoryUtilsTests.java @@ -0,0 +1,44 @@ +/* + * Copyright 2017 the original author or authors. + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package org.springframework.amqp.rabbit.connection; + +import static org.junit.Assert.assertSame; +import static org.mockito.Mockito.mock; + +import org.junit.Test; + +import org.springframework.transaction.support.TransactionSynchronizationManager; + +/** + * @author Gary Russell + * @since 1.7.1 + * + */ +public class ConnectionFactoryUtilsTests { + + @Test + public void testResourceHolder() { + RabbitResourceHolder h1 = new RabbitResourceHolder(); + RabbitResourceHolder h2 = new RabbitResourceHolder(); + ConnectionFactory connectionFactory = mock(ConnectionFactory.class); + TransactionSynchronizationManager.setActualTransactionActive(true); + ConnectionFactoryUtils.bindResourceToTransaction(h1, connectionFactory, true); + assertSame(h1, ConnectionFactoryUtils.bindResourceToTransaction(h2, connectionFactory, true)); + TransactionSynchronizationManager.clear(); + } + +} diff --git a/spring-rabbit/src/test/java/org/springframework/amqp/rabbit/listener/ExternalTxManagerTests.java b/spring-rabbit/src/test/java/org/springframework/amqp/rabbit/listener/ExternalTxManagerTests.java index c7fc9d21..5288d460 100644 --- a/spring-rabbit/src/test/java/org/springframework/amqp/rabbit/listener/ExternalTxManagerTests.java +++ b/spring-rabbit/src/test/java/org/springframework/amqp/rabbit/listener/ExternalTxManagerTests.java @@ -24,6 +24,7 @@ import static org.mockito.ArgumentMatchers.anyBoolean; import static org.mockito.ArgumentMatchers.anyLong; import static org.mockito.ArgumentMatchers.anyMap; import static org.mockito.ArgumentMatchers.anyString; +import static org.mockito.ArgumentMatchers.eq; import static org.mockito.BDDMockito.given; import static org.mockito.BDDMockito.willAnswer; import static org.mockito.Mockito.mock; @@ -35,11 +36,11 @@ import java.util.List; import java.util.concurrent.CountDownLatch; import java.util.concurrent.ExecutorService; import java.util.concurrent.TimeUnit; +import java.util.concurrent.atomic.AtomicBoolean; import java.util.concurrent.atomic.AtomicReference; import org.junit.Test; import org.mockito.Mockito; -import org.mockito.invocation.InvocationOnMock; import org.mockito.stubbing.Answer; import org.springframework.amqp.AmqpRejectAndDontRequeueException; @@ -55,6 +56,7 @@ import org.springframework.amqp.rabbit.transaction.RabbitTransactionManager; import org.springframework.beans.DirectFieldAccessor; import org.springframework.transaction.TransactionDefinition; import org.springframework.transaction.TransactionException; +import org.springframework.transaction.interceptor.DefaultTransactionAttribute; import org.springframework.transaction.interceptor.NoRollbackRuleAttribute; import org.springframework.transaction.interceptor.RollbackRuleAttribute; import org.springframework.transaction.interceptor.RuleBasedTransactionAttribute; @@ -94,20 +96,7 @@ public abstract class ExternalTxManagerTests { final AtomicReference tooManyChannels = new AtomicReference(); - willAnswer(new Answer() { - boolean done; - @Override - public Channel answer(InvocationOnMock invocation) throws Throwable { - if (!done) { - done = true; - return onlyChannel; - } - tooManyChannels.set(new Exception("More than one channel requested")); - Channel channel = mock(Channel.class); - given(channel.isOpen()).willReturn(true); - return channel; - } - }).given(mockConnection).createChannel(); + willAnswer(ensureOneChannelAnswer(onlyChannel, tooManyChannels)).given(mockConnection).createChannel(); final AtomicReference consumer = new AtomicReference(); final CountDownLatch consumerLatch = new CountDownLatch(1); @@ -242,11 +231,30 @@ public abstract class ExternalTxManagerTests { container.stop(); } + @Test + public void testMessageListenerRollback() throws Exception { + testMessageListenerRollbackGuts(true, TransactionDefinition.PROPAGATION_REQUIRED); + } + + @Test + public void testMessageListenerRollbackDontRequeue() throws Exception { + testMessageListenerRollbackGuts(false, TransactionDefinition.PROPAGATION_REQUIRED); + } + + @Test + public void testMessageListenerRollbackNoBoundTransaction() throws Exception { + testMessageListenerRollbackGuts(true, TransactionDefinition.PROPAGATION_NEVER); + } + + @Test + public void testMessageListenerRollbackDontRequeueNoBoundTransaction() throws Exception { + testMessageListenerRollbackGuts(false, TransactionDefinition.PROPAGATION_NEVER); + } + /** * Verifies that the channel is rolled back after an exception. */ - @Test - public void testMessageListenerRollback() throws Exception { + private void testMessageListenerRollbackGuts(boolean expectRequeue, int propagation) throws Exception { ConnectionFactory mockConnectionFactory = mock(ConnectionFactory.class); Connection mockConnection = mock(Connection.class); final Channel channel = mock(Channel.class); @@ -259,6 +267,7 @@ public abstract class ExternalTxManagerTests { given(mockConnection.isOpen()).willReturn(true); final AtomicReference tooManyChannels = new AtomicReference(); + willAnswer(ensureOneChannelAnswer(channel, tooManyChannels)).given(mockConnection).createChannel(); willAnswer(invocation -> channel).given(mockConnection).createChannel(); @@ -279,11 +288,25 @@ public abstract class ExternalTxManagerTests { return null; }).given(channel).txRollback(); + final CountDownLatch rejectLatch = new CountDownLatch(1); + willAnswer(invocation -> { + rejectLatch.countDown(); + return null; + }).given(channel).basicReject(anyLong(), anyBoolean()); + willAnswer(invocation -> { + rejectLatch.countDown(); + return null; + }).given(channel).basicNack(anyLong(), anyBoolean(), anyBoolean()); + + final CountDownLatch latch = new CountDownLatch(1); AbstractMessageListenerContainer container = createContainer(cachingConnectionFactory); + container.setTransactionAttribute(new DefaultTransactionAttribute(propagation)); container.setMessageListener(message -> { latch.countDown(); - throw new RuntimeException("force rollback"); + throw expectRequeue + ? new RuntimeException("force rollback") + : new AmqpRejectAndDontRequeueException("force rollback"); }); container.setQueueNames("queue"); container.setChannelTransacted(true); @@ -304,7 +327,99 @@ public abstract class ExternalTxManagerTests { } verify(mockConnection, times(1)).createChannel(); + assertTrue(rejectLatch.await(10, TimeUnit.SECONDS)); + assertTrue(rollbackLatch.await(10, TimeUnit.SECONDS)); + if (propagation != TransactionDefinition.PROPAGATION_NEVER) { + verify(channel).basicReject(anyLong(), eq(expectRequeue)); + } + else { + verify(channel).basicNack(anyLong(), eq(Boolean.TRUE), eq(expectRequeue)); + } + container.stop(); + } + + @Test + public void testMessageListenerCommit() throws Exception { + testMessageListenerCommitGuts(TransactionDefinition.PROPAGATION_REQUIRED); + } + + @Test + public void testMessageListenerCommitNoBoundTransaction() throws Exception { + testMessageListenerCommitGuts(TransactionDefinition.PROPAGATION_NEVER); + } + + /** + * Verifies that the channel is committed. + */ + private void testMessageListenerCommitGuts(int propagation) throws Exception { + ConnectionFactory mockConnectionFactory = mock(ConnectionFactory.class); + Connection mockConnection = mock(Connection.class); + final Channel channel = mock(Channel.class); + given(channel.isOpen()).willReturn(true); + + final CachingConnectionFactory cachingConnectionFactory = new CachingConnectionFactory(mockConnectionFactory); + cachingConnectionFactory.setExecutor(mock(ExecutorService.class)); + + given(mockConnectionFactory.newConnection(any(ExecutorService.class), anyString())).willReturn(mockConnection); + given(mockConnection.isOpen()).willReturn(true); + + final AtomicReference tooManyChannels = new AtomicReference(); + willAnswer(ensureOneChannelAnswer(channel, tooManyChannels)).given(mockConnection).createChannel(); + + willAnswer(invocation -> channel).given(mockConnection).createChannel(); + + final AtomicReference consumer = new AtomicReference(); + final CountDownLatch consumerLatch = new CountDownLatch(1); + + willAnswer(invocation -> { + consumer.set(invocation.getArgument(6)); + consumerLatch.countDown(); + return "consumerTag"; + }).given(channel) + .basicConsume(anyString(), anyBoolean(), anyString(), anyBoolean(), anyBoolean(), anyMap(), + any(Consumer.class)); + + final CountDownLatch commitLatch = new CountDownLatch(1); + willAnswer(invocation -> { + commitLatch.countDown(); + return null; + }).given(channel).txCommit(); + final CountDownLatch ackLatch = new CountDownLatch(1); + willAnswer(invocation -> { + ackLatch.countDown(); + return null; + }).given(channel).basicAck(anyLong(), anyBoolean()); + + final CountDownLatch latch = new CountDownLatch(1); + AbstractMessageListenerContainer container = createContainer(cachingConnectionFactory); + container.setTransactionAttribute(new DefaultTransactionAttribute(propagation)); + container.setMessageListener(message -> { + latch.countDown(); + }); + container.setQueueNames("queue"); + container.setChannelTransacted(true); + container.setShutdownTimeout(100); + container.setTransactionManager(new DummyTxManager()); + container.afterPropertiesSet(); + container.start(); + assertTrue(consumerLatch.await(10, TimeUnit.SECONDS)); + + consumer.get().handleDelivery("qux", new Envelope(1, false, "foo", "bar"), new BasicProperties(), + new byte[] { 0 }); + + assertTrue(latch.await(10, TimeUnit.SECONDS)); + + Exception e = tooManyChannels.get(); + if (e != null) { + throw e; + } + + verify(mockConnection, times(1)).createChannel(); + + assertTrue(ackLatch.await(10, TimeUnit.SECONDS)); + assertTrue(commitLatch.await(10, TimeUnit.SECONDS)); + verify(channel).basicAck(anyLong(), anyBoolean()); container.stop(); } @@ -341,20 +456,7 @@ public abstract class ExternalTxManagerTests { final AtomicReference tooManyChannels = new AtomicReference(); - willAnswer(new Answer() { - boolean done; - @Override - public Channel answer(InvocationOnMock invocation) throws Throwable { - if (!done) { - done = true; - return listenerChannel; - } - tooManyChannels.set(new Exception("More than one channel requested")); - Channel channel = mock(Channel.class); - given(channel.isOpen()).willReturn(true); - return channel; - } - }).given(listenerConnection).createChannel(); + willAnswer(ensureOneChannelAnswer(listenerChannel, tooManyChannels)).given(listenerConnection).createChannel(); final AtomicReference consumer = new AtomicReference(); final CountDownLatch consumerLatch = new CountDownLatch(1); @@ -438,20 +540,7 @@ public abstract class ExternalTxManagerTests { final AtomicReference tooManyChannels = new AtomicReference(); - willAnswer(new Answer() { - boolean done; - @Override - public Channel answer(InvocationOnMock invocation) throws Throwable { - if (!done) { - done = true; - return onlyChannel; - } - tooManyChannels.set(new Exception("More than one channel requested")); - Channel channel = mock(Channel.class); - given(channel.isOpen()).willReturn(true); - return channel; - } - }).given(mockConnection).createChannel(); + willAnswer(ensureOneChannelAnswer(onlyChannel, tooManyChannels)).given(mockConnection).createChannel(); final AtomicReference consumer = new AtomicReference(); final CountDownLatch consumerLatch = new CountDownLatch(1); @@ -531,20 +620,7 @@ public abstract class ExternalTxManagerTests { final AtomicReference tooManyChannels = new AtomicReference(); - willAnswer(new Answer() { - boolean done; - @Override - public Channel answer(InvocationOnMock invocation) throws Throwable { - if (!done) { - done = true; - return onlyChannel; - } - tooManyChannels.set(new Exception("More than one channel requested")); - Channel channel = mock(Channel.class); - given(channel.isOpen()).willReturn(true); - return channel; - } - }).given(mockConnection).createChannel(); + willAnswer(ensureOneChannelAnswer(onlyChannel, tooManyChannels)).given(mockConnection).createChannel(); final AtomicReference consumer = new AtomicReference(); final CountDownLatch consumerLatch = new CountDownLatch(1); @@ -625,20 +701,7 @@ public abstract class ExternalTxManagerTests { final AtomicReference tooManyChannels = new AtomicReference(); - willAnswer(new Answer() { - boolean done; - @Override - public Channel answer(InvocationOnMock invocation) throws Throwable { - if (!done) { - done = true; - return onlyChannel; - } - tooManyChannels.set(new Exception("More than one channel requested")); - Channel channel = mock(Channel.class); - given(channel.isOpen()).willReturn(true); - return channel; - } - }).given(mockConnection).createChannel(); + willAnswer(ensureOneChannelAnswer(onlyChannel, tooManyChannels)).given(mockConnection).createChannel(); final AtomicReference consumer = new AtomicReference(); final CountDownLatch consumerLatch = new CountDownLatch(1); @@ -697,6 +760,21 @@ public abstract class ExternalTxManagerTests { container.stop(); } + private Answer ensureOneChannelAnswer(final Channel onlyChannel, + final AtomicReference tooManyChannels) { + final AtomicBoolean done = new AtomicBoolean(); + return invocation -> { + if (!done.get()) { + done.set(true); + return onlyChannel; + } + tooManyChannels.set(new Exception("More than one channel requested")); + Channel channel = mock(Channel.class); + given(channel.isOpen()).willReturn(true); + return channel; + }; + } + protected abstract AbstractMessageListenerContainer createContainer(AbstractConnectionFactory connectionFactory); @SuppressWarnings("serial") diff --git a/src/reference/asciidoc/amqp.adoc b/src/reference/asciidoc/amqp.adoc index 339d69a2..b54d3765 100644 --- a/src/reference/asciidoc/amqp.adoc +++ b/src/reference/asciidoc/amqp.adoc @@ -3870,6 +3870,7 @@ The default `FatalExceptionStrategy` logs a warning message when an exception is Since _version 1.6.3_ a convenient way to add user exceptions to the fatal list is to subclass `ConditionalRejectingErrorHandler.DefaultExceptionStrategy` and override the method `isUserCauseFatal(Throwable cause)` to return true for fatal exceptions. +[[transactions]] ==== Transactions ===== Introduction @@ -3961,6 +3962,7 @@ public AbstractMessageListenerContainer container() { } ---- +[[transaction-rollback]] ===== A note on Rollback of Received Messages AMQP transactions only apply to messages and acks sent to the broker, so when there is a rollback of a Spring transaction and a message has been received, what Spring AMQP has to do is not just rollback the transaction, but also manually reject the message (sort of a nack, but that's not what the specification calls it). @@ -3971,6 +3973,15 @@ For more information about RabbitMQ transactions, and their limitations, refer t NOTE: Prior to *RabbitMQ 2.7.0*, such messages (and any that are unacked when a channel is closed or aborts) went to the back of the queue on a Rabbit broker, since 2.7.0, rejected messages go to the front of the queue, in a similar manner to JMS rolled back messages. +[NOTE] +==== +Previously, message requeue on transaction rollback was inconsistent between local transactions and when a `TransactionManager` was provided. +In the former case, the normal requeue logic (`AmqpRejectAndDontRequeueException` or `defaultRequeueRejected=false`) applied (see <>); with a transaction manager, the message was unconditionally requeued on rollback. +Starting with __version 2.0__, the behavior is consistent and the normal requeue logic is applied in both cases. +To revert to the previous behavior, set the container's `alwaysRequeueWithTxManagerRollback` property to `true`. +See <>. +==== + ===== Using the RabbitTransactionManager The http://static.springsource.org/spring-amqp/docs/latest_ga/api/org/springframework/amqp/rabbit/transaction/RabbitTransactionManager.html[RabbitTransactionManager] is an alternative to executing Rabbit operations within, and synchronized with, external transactions. @@ -4457,6 +4468,15 @@ Set this to false, to discard (or route to a dead-letter queue) such messages. a| image::images/tickmark.png[] a| image::images/tickmark.png[] +| alwaysRequeueWithTx +ManagerRollback +(N/A) + +| Set to `true` to always requeue messages on rollback when a transaction manager is configured. + +a| image::images/tickmark.png[] +a| image::images/tickmark.png[] + |=== [[listener-concurrency]] diff --git a/src/reference/asciidoc/whats-new.adoc b/src/reference/asciidoc/whats-new.adoc index 052ee871..f910fcb5 100644 --- a/src/reference/asciidoc/whats-new.adoc +++ b/src/reference/asciidoc/whats-new.adoc @@ -48,9 +48,16 @@ See <> for more information. ===== Listener Container Changes +====== Message Count + Previously, `MessageProperties.getMessageCount()` returned `0` for messages emitted by the container. This property only applies when using `basicGet` (e.g. from `RabbitTemplate.receive()` methods) and is now initialized to `null` for container messages. +====== Transaction Rollback behavior + +Message requeue on transaction rollback is now consistent, regardless of whether or not a transaction manager is configured. +See <> for more information. + ===== Connection Factory Changes The connection and channel listener interfaces now provide a mechanism to obtain information about exceptions.