diff --git a/spring-rabbit/src/main/java/org/springframework/amqp/rabbit/connection/AfterCompletionFailedException.java b/spring-rabbit/src/main/java/org/springframework/amqp/rabbit/connection/AfterCompletionFailedException.java new file mode 100644 index 00000000..259f0d59 --- /dev/null +++ b/spring-rabbit/src/main/java/org/springframework/amqp/rabbit/connection/AfterCompletionFailedException.java @@ -0,0 +1,52 @@ +/* + * Copyright 2021 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 + * + * https://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 org.springframework.amqp.AmqpException; + +/** + * Represents a failure to commit or rollback when performing afterCompletion + * after the primary transaction completes. + * + * @author Gary Russell + * @since 2.4 + */ +public class AfterCompletionFailedException extends AmqpException { + + private static final long serialVersionUID = 1L; + + private final int syncStatus; + + /** + * Construct an instance with the provided properties. + * @param syncStatus the synchronization status. + * @param cause the cause. + */ + public AfterCompletionFailedException(int syncStatus, Throwable cause) { + super(cause); + this.syncStatus = syncStatus; + } + + /** + * Return the synchronization status. + * @return the status. + */ + public int getSyncStatus() { + return this.syncStatus; + } + +} 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 f2c66994..637c28e2 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 @@ -17,6 +17,7 @@ package org.springframework.amqp.rabbit.connection; import java.io.IOException; +import java.util.function.Consumer; import org.springframework.amqp.AmqpIOException; import org.springframework.lang.Nullable; @@ -42,6 +43,10 @@ import com.rabbitmq.client.Channel; */ public final class ConnectionFactoryUtils { + private static final ThreadLocal failures = new ThreadLocal<>(); + + private static boolean captureAfterCompletionExceptions; + private ConnectionFactoryUtils() { } @@ -181,11 +186,42 @@ public final class ConnectionFactoryUtils { resourceHolder.setSynchronizedWithTransaction(true); if (TransactionSynchronizationManager.isSynchronizationActive()) { TransactionSynchronizationManager.registerSynchronization(new RabbitResourceSynchronization(resourceHolder, - connectionFactory)); + connectionFactory, ConnectionFactoryUtils::completionFailed)); } return resourceHolder; } + private static void completionFailed(AfterCompletionFailedException ex) { + if (captureAfterCompletionExceptions) { + failures.set(ex); + } + } + + /** + * Call this method to enable capturing {@link AfterCompletionFailedException}s + * when using transaction synchronization. Exceptions are stored in a {@link ThreadLocal} + * which must be cleared by calling {@link #checkAfterCompletion()} after the transaction + * has completed. + * @param enable true to enable capture. + */ + public static void enableAfterCompletionFailureCapture(boolean enable) { + captureAfterCompletionExceptions = enable; + } + + /** + * When using transaction synchronization, call this method after the transaction commits to + * verify that the RabbitMQ transaction committed. + * @throws AfterCompletionFailedException if synchronization failed. + * @since 2.3.10 + */ + public static void checkAfterCompletion() { + AfterCompletionFailedException ex = failures.get(); + if (ex != null) { + failures.remove(); + throw ex; + } + } + public static void registerDeliveryTag(ConnectionFactory connectionFactory, Channel channel, Long tag) { Assert.notNull(connectionFactory, "ConnectionFactory must not be null"); @@ -318,9 +354,14 @@ public final class ConnectionFactoryUtils { private final RabbitResourceHolder resourceHolder; - RabbitResourceSynchronization(RabbitResourceHolder resourceHolder, Object resourceKey) { + private final Consumer afterCompletionCallback; + + RabbitResourceSynchronization(RabbitResourceHolder resourceHolder, Object resourceKey, + Consumer afterCompletionCallback) { + super(resourceHolder, resourceKey); this.resourceHolder = resourceHolder; + this.afterCompletionCallback = afterCompletionCallback; } @Override @@ -338,6 +379,9 @@ public final class ConnectionFactoryUtils { this.resourceHolder.rollbackAll(); } } + catch (RuntimeException ex) { + this.afterCompletionCallback.accept(new AfterCompletionFailedException(status, ex)); + } finally { if (this.resourceHolder.isReleaseAfterCompletion()) { this.resourceHolder.setSynchronizedWithTransaction(false); diff --git a/spring-rabbit/src/test/java/org/springframework/amqp/rabbit/core/RabbitTemplateTests.java b/spring-rabbit/src/test/java/org/springframework/amqp/rabbit/core/RabbitTemplateTests.java index 7208ba7d..c56bf2fc 100644 --- a/spring-rabbit/src/test/java/org/springframework/amqp/rabbit/core/RabbitTemplateTests.java +++ b/spring-rabbit/src/test/java/org/springframework/amqp/rabbit/core/RabbitTemplateTests.java @@ -17,6 +17,7 @@ package org.springframework.amqp.rabbit.core; import static org.assertj.core.api.Assertions.assertThat; +import static org.assertj.core.api.Assertions.assertThatExceptionOfType; import static org.assertj.core.api.Assertions.assertThatIllegalStateException; import static org.assertj.core.api.Assertions.assertThatThrownBy; import static org.assertj.core.api.Assertions.fail; @@ -62,8 +63,10 @@ import org.springframework.amqp.core.Queue; import org.springframework.amqp.core.ReceiveAndReplyCallback; import org.springframework.amqp.core.ReturnedMessage; import org.springframework.amqp.rabbit.connection.AbstractRoutingConnectionFactory; +import org.springframework.amqp.rabbit.connection.AfterCompletionFailedException; import org.springframework.amqp.rabbit.connection.CachingConnectionFactory; import org.springframework.amqp.rabbit.connection.ChannelProxy; +import org.springframework.amqp.rabbit.connection.ConnectionFactoryUtils; import org.springframework.amqp.rabbit.connection.PublisherCallbackChannel; import org.springframework.amqp.rabbit.connection.RabbitUtils; import org.springframework.amqp.rabbit.connection.SimpleRoutingConnectionFactory; @@ -616,6 +619,7 @@ public class RabbitTemplateTests { @Test void resourcesClearedAfterTxFailsWithSync() throws IOException, TimeoutException { + ConnectionFactoryUtils.enableAfterCompletionFailureCapture(true); ConnectionFactory mockConnectionFactory = mock(ConnectionFactory.class); Connection mockConnection = mock(Connection.class); Channel mockChannel = mock(Channel.class); @@ -638,6 +642,9 @@ public class RabbitTemplateTests { assertThatIllegalStateException() .isThrownBy(() -> (TransactionSynchronizationManager.getSynchronizations()).isEmpty()) .withMessage("Transaction synchronization is not active"); + assertThatExceptionOfType(AfterCompletionFailedException.class) + .isThrownBy(() -> ConnectionFactoryUtils.checkAfterCompletion()); + ConnectionFactoryUtils.enableAfterCompletionFailureCapture(false); } @SuppressWarnings("serial") diff --git a/src/reference/asciidoc/amqp.adoc b/src/reference/asciidoc/amqp.adoc index 445560ab..3fed22b8 100644 --- a/src/reference/asciidoc/amqp.adoc +++ b/src/reference/asciidoc/amqp.adoc @@ -5673,6 +5673,17 @@ If you prefer XML configuration, you can declare the following bean in your XML ---- ==== +[[tx-sync]] +===== Transaction Synchronization + +Synchronizing a RabbitMQ transaction with some other (e.g. DBMS) transaction provides "Best Effort One Phase Commit" semantics. +It is possible that the RabbitMQ transaction fails to commit during the after completion phase of transaction synchronization. +This is logged by the `spring-tx` infrastructure as an error, but no exception is thrown to the calling code. +Starting with version 2.3.10, you can call `ConnectionUtils.checkAfterCompletion()` after the transaction has committed on the same thread that processed the transaction. +It will simply return if no exception occurred; otherwise it will throw an `AfterCompletionFailedException` which will have a property representing the synchronization status of the completion. + +Enable this feature by calling `ConnectionFactoryUtils.enableAfterCompletionFailureCapture(true)`; this is a global flag and applies to all threads. + [[containerAttributes]] ==== Message Listener Container Configuration