From aa3589c2d63b65afce358f16268e97660ef1d4b1 Mon Sep 17 00:00:00 2001 From: Gary Russell Date: Mon, 20 Dec 2021 11:44:08 -0500 Subject: [PATCH] GH-1409: Fix Nacks for Async Replies Resolves https://github.com/spring-projects/spring-amqp/issues/1409 Normally, when message has a fatal exception (such as message conversion) `basicNack` with `multiple` true is used, to nack any previously unacked messages (e.g. when using batch size to limit the ack traffic). Even when using manual acks, fatal exceptions are nacked by the container because the user does not have access to the message. However, when using async replies, this has the side effect of nacking unprocessed messages. Detect whether async replies are being used and only nack individual records that cause fatal exceptions. Also, coerce the `AcknowledgeMode` to `MANUAL` for such listners. Add a test for both containers; send a good message followed by a bad one without actually completing the reply future. After the exception occurs and the container is stopped, there should be one messag in the queue. * Remove warning, deprecation; add docs. * Docs. **Cherry-pick to `2.3.x` & `2.2.x`** # Conflicts: # spring-rabbit/src/main/java/org/springframework/amqp/rabbit/listener/adapter/MessagingMessageListenerAdapter.java # Conflicts: # spring-rabbit/src/main/java/org/springframework/amqp/rabbit/listener/adapter/DelegatingInvocableHandler.java # spring-rabbit/src/main/java/org/springframework/amqp/rabbit/listener/adapter/HandlerAdapter.java --- .../amqp/core/MessageListener.java | 12 +- .../AbstractMessageListenerContainer.java | 10 +- .../listener/BlockingQueueConsumer.java | 38 +++- .../DirectMessageListenerContainer.java | 2 +- .../SimpleMessageListenerContainer.java | 7 +- .../AbstractAdaptableMessageListener.java | 19 +- .../adapter/DelegatingInvocableHandler.java | 25 +++ .../listener/adapter/HandlerAdapter.java | 27 +++ .../MessagingMessageListenerAdapter.java | 5 + .../rabbit/listener/adapter/MonoHandler.java | 44 +++++ .../rabbit/listener/AsyncReplyToTests.java | 178 ++++++++++++++++++ .../SimpleMessageListenerContainerTests.java | 1 + src/reference/asciidoc/amqp.adoc | 3 + 13 files changed, 340 insertions(+), 31 deletions(-) create mode 100644 spring-rabbit/src/main/java/org/springframework/amqp/rabbit/listener/adapter/MonoHandler.java create mode 100644 spring-rabbit/src/test/java/org/springframework/amqp/rabbit/listener/AsyncReplyToTests.java diff --git a/spring-amqp/src/main/java/org/springframework/amqp/core/MessageListener.java b/spring-amqp/src/main/java/org/springframework/amqp/core/MessageListener.java index e0dbe839..34f786dd 100644 --- a/spring-amqp/src/main/java/org/springframework/amqp/core/MessageListener.java +++ b/spring-amqp/src/main/java/org/springframework/amqp/core/MessageListener.java @@ -1,5 +1,5 @@ /* - * Copyright 2002-2019 the original author or authors. + * Copyright 2002-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. @@ -43,6 +43,16 @@ public interface MessageListener { // NOSONAR - empty } + /** + * Return true if this listener is request/reply and the replies are + * async. + * @return true for async replies. + * @since 2.2.21 + */ + default boolean isAsyncReplies() { + return false; + } + /** * Delivers a batch of messages. * @param messages the messages. 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 cc0409c0..b648e222 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 @@ -250,6 +250,8 @@ public abstract class AbstractMessageListenerContainer extends RabbitAccessor private volatile boolean lazyLoad; + private boolean asyncReplies; + @Override public void setApplicationEventPublisher(ApplicationEventPublisher applicationEventPublisher) { this.applicationEventPublisher = applicationEventPublisher; @@ -433,6 +435,7 @@ public abstract class AbstractMessageListenerContainer extends RabbitAccessor this.messageListener = messageListener; this.isBatchListener = messageListener instanceof BatchMessageListener || messageListener instanceof ChannelAwareBatchMessageListener; + this.asyncReplies = messageListener.isAsyncReplies(); } /** @@ -1009,10 +1012,12 @@ public abstract class AbstractMessageListenerContainer extends RabbitAccessor return this.possibleAuthenticationFailureFatal; } - protected boolean isPossibleAuthenticationFailureFatalSet() { return this.possibleAuthenticationFailureFatalSet; } + protected boolean isAsyncReplies() { + return this.asyncReplies; + } /** * Set to true to automatically declare elements (queues, exchanges, bindings) @@ -1193,6 +1198,9 @@ public abstract class AbstractMessageListenerContainer extends RabbitAccessor catch (IllegalStateException e) { this.logger.debug("Could not enable micrometer timers", e); } + if (this.isAsyncReplies() && !AcknowledgeMode.MANUAL.equals(this.acknowledgeMode)) { + this.acknowledgeMode = AcknowledgeMode.MANUAL; + } } @Override 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 fb766397..7baf68ab 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 @@ -758,6 +758,17 @@ public class BlockingQueueConsumer { * @param ex the thrown application exception or error */ public void rollbackOnExceptionIfNecessary(Throwable ex) { + rollbackOnExceptionIfNecessary(ex, -1); + } + + /** + * Perform a rollback, handling rollback exceptions properly. + * @param ex the thrown application exception or error + * @param tag delivery tag; when specified (greater than or equal to 0) only that + * message is nacked. + * @since 2.2.21. + */ + public void rollbackOnExceptionIfNecessary(Throwable ex, long tag) { boolean ackRequired = !this.acknowledgeMode.isAutoAck() && (!this.acknowledgeMode.isManual() || ContainerUtils.isRejectManual(ex)); @@ -769,14 +780,20 @@ public class BlockingQueueConsumer { RabbitUtils.rollbackIfNecessary(this.channel); } if (ackRequired) { - OptionalLong deliveryTag = this.deliveryTags.stream().mapToLong(l -> l).max(); - if (deliveryTag.isPresent()) { - this.channel.basicNack(deliveryTag.getAsLong(), true, - ContainerUtils.shouldRequeue(this.defaultRequeueRejected, ex, logger)); + if (tag < 0) { + OptionalLong deliveryTag = this.deliveryTags.stream().mapToLong(l -> l).max(); + if (deliveryTag.isPresent()) { + this.channel.basicNack(deliveryTag.getAsLong(), true, + ContainerUtils.shouldRequeue(this.defaultRequeueRejected, ex, logger)); + } + if (this.transactional) { + // Need to commit the reject (=nack) + RabbitUtils.commitIfNecessary(this.channel); + } } - if (this.transactional) { - // Need to commit the reject (=nack) - RabbitUtils.commitIfNecessary(this.channel); + else { + this.channel.basicNack(tag, false, + ContainerUtils.shouldRequeue(this.defaultRequeueRejected, ex, logger)); } } } @@ -785,7 +802,12 @@ public class BlockingQueueConsumer { throw RabbitExceptionTranslator.convertRabbitAccessException(e); // NOSONAR stack trace loss } finally { - this.deliveryTags.clear(); + if (tag < 0) { + this.deliveryTags.clear(); + } + else { + this.deliveryTags.remove(tag); + } } } 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 b46d1cd9..563b08c3 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 @@ -1196,7 +1196,7 @@ public class DirectMessageListenerContainer extends AbstractMessageListenerConta } } } - getChannel().basicNack(deliveryTag, true, + getChannel().basicNack(deliveryTag, !isAsyncReplies(), ContainerUtils.shouldRequeue(isDefaultRequeueRejected(), e, this.logger)); } catch (IOException e1) { 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 76ab289f..ddb45905 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 @@ -991,6 +991,9 @@ public class SimpleMessageListenerContainer extends AbstractMessageListenerConta } break; } + long tagToRollback = isAsyncReplies() + ? message.getMessageProperties().getDeliveryTag() + : -1; if (getTransactionManager() != null) { if (getTransactionAttribute().rollbackOn(ex)) { RabbitResourceHolder resourceHolder = (RabbitResourceHolder) TransactionSynchronizationManager @@ -1003,7 +1006,7 @@ public class SimpleMessageListenerContainer extends AbstractMessageListenerConta * If we don't actually have a transaction, we have to roll back * manually. See prepareHolderForRollback(). */ - consumer.rollbackOnExceptionIfNecessary(ex); + consumer.rollbackOnExceptionIfNecessary(ex, tagToRollback); } throw ex; // encompassing transaction will handle the rollback. } @@ -1015,7 +1018,7 @@ public class SimpleMessageListenerContainer extends AbstractMessageListenerConta } } else { - consumer.rollbackOnExceptionIfNecessary(ex); + consumer.rollbackOnExceptionIfNecessary(ex, tagToRollback); throw ex; } } diff --git a/spring-rabbit/src/main/java/org/springframework/amqp/rabbit/listener/adapter/AbstractAdaptableMessageListener.java b/spring-rabbit/src/main/java/org/springframework/amqp/rabbit/listener/adapter/AbstractAdaptableMessageListener.java index 34f0ae4f..2a3fae72 100644 --- a/spring-rabbit/src/main/java/org/springframework/amqp/rabbit/listener/adapter/AbstractAdaptableMessageListener.java +++ b/spring-rabbit/src/main/java/org/springframework/amqp/rabbit/listener/adapter/AbstractAdaptableMessageListener.java @@ -21,7 +21,6 @@ import java.lang.reflect.ParameterizedType; import java.lang.reflect.Type; import java.lang.reflect.WildcardType; import java.util.Arrays; -import java.util.function.Consumer; import org.apache.commons.logging.Log; import org.apache.commons.logging.LogFactory; @@ -56,7 +55,6 @@ import org.springframework.util.ClassUtils; import org.springframework.util.concurrent.ListenableFuture; import com.rabbitmq.client.Channel; -import reactor.core.publisher.Mono; /** * An abstract {@link org.springframework.amqp.core.MessageListener} adapter providing the @@ -81,7 +79,7 @@ public abstract class AbstractAdaptableMessageListener implements ChannelAwareMe private static final ParserContext PARSER_CONTEXT = new TemplateParserContext("!{", "}"); - private static final boolean monoPresent = // NOSONAR - lower case + static final boolean monoPresent = // NOSONAR - lower case, protected ClassUtils.isPresent("reactor.core.publisher.Mono", ChannelAwareMessageListener.class.getClassLoader()); /** @@ -635,19 +633,4 @@ public abstract class AbstractAdaptableMessageListener implements ChannelAwareMe } - private static class MonoHandler { // NOSONAR - pointless to name it ..Utils|Helper - - static boolean isMono(Object result) { - return result instanceof Mono; - } - - @SuppressWarnings("unchecked") - static void subscribe(Object returnValue, Consumer success, - Consumer failure, Runnable completeConsumer) { - - ((Mono) returnValue).subscribe(success, failure, completeConsumer); - } - - } - } diff --git a/spring-rabbit/src/main/java/org/springframework/amqp/rabbit/listener/adapter/DelegatingInvocableHandler.java b/spring-rabbit/src/main/java/org/springframework/amqp/rabbit/listener/adapter/DelegatingInvocableHandler.java index 65bd6452..68d9d9c0 100644 --- a/spring-rabbit/src/main/java/org/springframework/amqp/rabbit/listener/adapter/DelegatingInvocableHandler.java +++ b/spring-rabbit/src/main/java/org/springframework/amqp/rabbit/listener/adapter/DelegatingInvocableHandler.java @@ -20,6 +20,7 @@ import java.lang.annotation.Annotation; import java.lang.reflect.Method; import java.util.ArrayList; import java.util.Arrays; +import java.util.Iterator; import java.util.List; import java.util.Map; import java.util.concurrent.ConcurrentHashMap; @@ -41,6 +42,7 @@ import org.springframework.messaging.handler.annotation.Header; import org.springframework.messaging.handler.annotation.SendTo; import org.springframework.messaging.handler.invocation.InvocableHandlerMethod; import org.springframework.util.Assert; +import org.springframework.util.concurrent.ListenableFuture; /** @@ -75,6 +77,8 @@ public class DelegatingInvocableHandler { private final BeanExpressionContext beanExpressionContext; + private final boolean asyncReplies; + /** * Construct an instance with the supplied handlers for the bean. * @param handlers the handlers. @@ -106,6 +110,18 @@ public class DelegatingInvocableHandler { this.bean = bean; this.resolver = beanExpressionResolver; this.beanExpressionContext = beanExpressionContext; + boolean asyncReplies; + asyncReplies = defaultHandler != null && isAsyncReply(defaultHandler); + Iterator iterator = handlers.iterator(); + while (iterator.hasNext()) { + asyncReplies |= isAsyncReply(iterator.next()); + } + this.asyncReplies = asyncReplies; + } + + private boolean isAsyncReply(InvocableHandlerMethod method) { + return (AbstractAdaptableMessageListener.monoPresent && MonoHandler.isMono(method.getMethod().getReturnType())) + || ListenableFuture.class.isAssignableFrom(method.getMethod().getReturnType()); } /** @@ -115,6 +131,15 @@ public class DelegatingInvocableHandler { return this.bean; } + /** + * Return true if any handler method has an async reply type. + * @return the asyncReply. + * @since 2.2.21 + */ + public boolean isAsyncReplies() { + return this.asyncReplies; + } + /** * Invoke the method with the given message. * @param message the message. diff --git a/spring-rabbit/src/main/java/org/springframework/amqp/rabbit/listener/adapter/HandlerAdapter.java b/spring-rabbit/src/main/java/org/springframework/amqp/rabbit/listener/adapter/HandlerAdapter.java index 4ec8f724..308ca75b 100644 --- a/spring-rabbit/src/main/java/org/springframework/amqp/rabbit/listener/adapter/HandlerAdapter.java +++ b/spring-rabbit/src/main/java/org/springframework/amqp/rabbit/listener/adapter/HandlerAdapter.java @@ -22,6 +22,7 @@ import java.lang.reflect.Type; import org.springframework.lang.Nullable; import org.springframework.messaging.Message; import org.springframework.messaging.handler.invocation.InvocableHandlerMethod; +import org.springframework.util.concurrent.ListenableFuture; /** * A wrapper for either an {@link InvocableHandlerMethod} or @@ -38,14 +39,24 @@ public class HandlerAdapter { private final DelegatingInvocableHandler delegatingHandler; + private final boolean asyncReplies; + + /** + * Construct an instance with the provided method. + * @param invokerHandlerMethod the method. + */ public HandlerAdapter(InvocableHandlerMethod invokerHandlerMethod) { this.invokerHandlerMethod = invokerHandlerMethod; this.delegatingHandler = null; + this.asyncReplies = (AbstractAdaptableMessageListener.monoPresent + && MonoHandler.isMono(invokerHandlerMethod.getMethod().getReturnType())) + || ListenableFuture.class.isAssignableFrom(invokerHandlerMethod.getMethod().getReturnType()); } public HandlerAdapter(DelegatingInvocableHandler delegatingHandler) { this.invokerHandlerMethod = null; this.delegatingHandler = delegatingHandler; + this.asyncReplies = delegatingHandler.isAsyncReplies(); } public InvocationResult invoke(Message message, Object... providedArgs) throws Exception { // NOSONAR @@ -126,6 +137,22 @@ public class HandlerAdapter { } } + /** + * Return true if any handler method has an async reply type. + * @return the asyncReply. + * @since 2.2.21 + */ + public boolean isAsyncReplies() { + return this.asyncReplies; + } + + /** + * Build an {@link InvocationResult} for the result and inbound payload. + * @param result the result. + * @param inboundPayload the payload. + * @return the invocation result. + * @since 2.1.7 + */ @Nullable public InvocationResult getInvocationResultFor(Object result, Object inboundPayload) { if (this.invokerHandlerMethod != null) { diff --git a/spring-rabbit/src/main/java/org/springframework/amqp/rabbit/listener/adapter/MessagingMessageListenerAdapter.java b/spring-rabbit/src/main/java/org/springframework/amqp/rabbit/listener/adapter/MessagingMessageListenerAdapter.java index b350a828..b1a599ba 100644 --- a/spring-rabbit/src/main/java/org/springframework/amqp/rabbit/listener/adapter/MessagingMessageListenerAdapter.java +++ b/spring-rabbit/src/main/java/org/springframework/amqp/rabbit/listener/adapter/MessagingMessageListenerAdapter.java @@ -103,6 +103,11 @@ public class MessagingMessageListenerAdapter extends AbstractAdaptableMessageLis this.handlerAdapter = handlerAdapter; } + @Override + public boolean isAsyncReplies() { + return this.handlerAdapter.isAsyncReplies(); + } + /** * Set the {@link AmqpHeaderMapper} implementation to use to map the standard * AMQP headers. By default, a {@link org.springframework.amqp.support.SimpleAmqpHeaderMapper diff --git a/spring-rabbit/src/main/java/org/springframework/amqp/rabbit/listener/adapter/MonoHandler.java b/spring-rabbit/src/main/java/org/springframework/amqp/rabbit/listener/adapter/MonoHandler.java new file mode 100644 index 00000000..f594e44e --- /dev/null +++ b/spring-rabbit/src/main/java/org/springframework/amqp/rabbit/listener/adapter/MonoHandler.java @@ -0,0 +1,44 @@ +/* + * 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.listener.adapter; + +import java.util.function.Consumer; + +import reactor.core.publisher.Mono; + +/** + * Class to prevent direct links to {@link Mono}. + * @author Gary Russell + * @since 2.2.21 + */ +final class MonoHandler { // NOSONAR - pointless to name it ..Utils|Helper + + private MonoHandler() { + } + + static boolean isMono(Object result) { + return result instanceof Mono; + } + + @SuppressWarnings("unchecked") + static void subscribe(Object returnValue, Consumer success, + Consumer failure, Runnable completeConsumer) { + + ((Mono) returnValue).subscribe(success, failure, completeConsumer); + } + +} diff --git a/spring-rabbit/src/test/java/org/springframework/amqp/rabbit/listener/AsyncReplyToTests.java b/spring-rabbit/src/test/java/org/springframework/amqp/rabbit/listener/AsyncReplyToTests.java new file mode 100644 index 00000000..3bfab605 --- /dev/null +++ b/spring-rabbit/src/test/java/org/springframework/amqp/rabbit/listener/AsyncReplyToTests.java @@ -0,0 +1,178 @@ +/* + * 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.listener; + +import static org.assertj.core.api.Assertions.assertThat; + +import java.io.IOException; +import java.util.concurrent.CountDownLatch; +import java.util.concurrent.TimeUnit; +import java.util.concurrent.TimeoutException; + +import org.junit.jupiter.api.Test; + +import org.springframework.amqp.core.AcknowledgeMode; +import org.springframework.amqp.core.MessageBuilder; +import org.springframework.amqp.core.MessagePropertiesBuilder; +import org.springframework.amqp.rabbit.annotation.EnableRabbit; +import org.springframework.amqp.rabbit.annotation.RabbitListener; +import org.springframework.amqp.rabbit.config.DirectRabbitListenerContainerFactory; +import org.springframework.amqp.rabbit.config.SimpleRabbitListenerContainerFactory; +import org.springframework.amqp.rabbit.connection.CachingConnectionFactory; +import org.springframework.amqp.rabbit.connection.ConnectionFactory; +import org.springframework.amqp.rabbit.core.RabbitAdmin; +import org.springframework.amqp.rabbit.core.RabbitTemplate; +import org.springframework.amqp.rabbit.junit.RabbitAvailable; +import org.springframework.amqp.rabbit.junit.RabbitAvailableCondition; +import org.springframework.amqp.support.converter.Jackson2JsonMessageConverter; +import org.springframework.amqp.support.converter.MessageConverter; +import org.springframework.beans.factory.annotation.Autowired; +import org.springframework.context.annotation.Bean; +import org.springframework.context.annotation.Configuration; +import org.springframework.test.context.junit.jupiter.SpringJUnitConfig; +import org.springframework.util.concurrent.ListenableFuture; +import org.springframework.util.concurrent.SettableListenableFuture; + +import com.rabbitmq.client.Channel; + +/** + * @author Gary Russell + * @since 2.2.21 + * + */ +@SpringJUnitConfig +@RabbitAvailable(queues = { "async1", "async2" }) +public class AsyncReplyToTests { + + @Test + void ackSingleWhenFatalSMLC(@Autowired Config config, @Autowired RabbitListenerEndpointRegistry registry, + @Autowired RabbitTemplate template, @Autowired RabbitAdmin admin) throws IOException, InterruptedException { + + template.send("async1", MessageBuilder.withBody("\"foo\"".getBytes()).andProperties( + MessagePropertiesBuilder.newInstance() + .setContentType("application/json") + .setReplyTo("nowhere") + .build()) + .build()); + template.send("async1", MessageBuilder.withBody("junk".getBytes()).andProperties( + MessagePropertiesBuilder.newInstance() + .setContentType("application/json") + .setReplyTo("nowhere") + .build()) + .build()); + assertThat(config.smlcLatch.await(10, TimeUnit.SECONDS)).isTrue(); + registry.getListenerContainer("smlc").stop(); + assertThat(admin.getQueueInfo("async1").getMessageCount()).isEqualTo(1); + } + + @Test + void ackSingleWhenFatalDMLC(@Autowired Config config, @Autowired RabbitListenerEndpointRegistry registry, + @Autowired RabbitTemplate template, @Autowired RabbitAdmin admin) throws IOException, InterruptedException { + + template.send("async2", MessageBuilder.withBody("\"foo\"".getBytes()).andProperties( + MessagePropertiesBuilder.newInstance() + .setContentType("application/json") + .setReplyTo("nowhere") + .build()) + .build()); + template.send("async2", MessageBuilder.withBody("junk".getBytes()).andProperties( + MessagePropertiesBuilder.newInstance() + .setContentType("application/json") + .setReplyTo("nowhere") + .build()) + .build()); + assertThat(config.dmlcLatch.await(10, TimeUnit.SECONDS)).isTrue(); + registry.getListenerContainer("dmlc").stop(); + assertThat(admin.getQueueInfo("async2").getMessageCount()).isEqualTo(1); + } + + @Configuration + @EnableRabbit + static class Config { + + volatile CountDownLatch smlcLatch = new CountDownLatch(1); + + volatile CountDownLatch dmlcLatch = new CountDownLatch(1); + + @RabbitListener(id = "smlc", queues = "async1", containerFactory = "smlcf") + ListenableFuture listen1(String in, Channel channel) { + return new SettableListenableFuture<>(); + } + + @RabbitListener(id = "dmlc", queues = "async2", containerFactory = "dmlcf") + ListenableFuture listen2(String in, Channel channel) { + return new SettableListenableFuture<>(); + } + + @Bean + MessageConverter converter() { + return new Jackson2JsonMessageConverter(); + } + + @Bean + ConnectionFactory cf() throws IOException, TimeoutException { + return new CachingConnectionFactory(RabbitAvailableCondition.getBrokerRunning().getConnectionFactory()); + } + + @Bean + SimpleRabbitListenerContainerFactory smlcf(ConnectionFactory cf, MessageConverter converter) { + SimpleRabbitListenerContainerFactory factory = new SimpleRabbitListenerContainerFactory(); + factory.setConnectionFactory(cf); + factory.setAcknowledgeMode(AcknowledgeMode.MANUAL); + factory.setMessageConverter(converter); + factory.setErrorHandler(new ConditionalRejectingErrorHandler() { + + @Override + public void handleError(Throwable t) { + smlcLatch.countDown(); + super.handleError(t); + } + + }); + return factory; + } + + @Bean + DirectRabbitListenerContainerFactory dmlcf(ConnectionFactory cf, MessageConverter converter) { + DirectRabbitListenerContainerFactory factory = new DirectRabbitListenerContainerFactory(); + factory.setConnectionFactory(cf); + factory.setAcknowledgeMode(AcknowledgeMode.MANUAL); + factory.setMessageConverter(converter); + factory.setErrorHandler(new ConditionalRejectingErrorHandler() { + + @Override + public void handleError(Throwable t) { + dmlcLatch.countDown(); + super.handleError(t); + } + + }); + return factory; + } + + @Bean + RabbitTemplate template(ConnectionFactory cf) { + return new RabbitTemplate(cf); + } + + @Bean + RabbitAdmin admin(ConnectionFactory cf) { + return new RabbitAdmin(cf); + } + + } +} diff --git a/spring-rabbit/src/test/java/org/springframework/amqp/rabbit/listener/SimpleMessageListenerContainerTests.java b/spring-rabbit/src/test/java/org/springframework/amqp/rabbit/listener/SimpleMessageListenerContainerTests.java index 26ff440e..1c70c23d 100644 --- a/spring-rabbit/src/test/java/org/springframework/amqp/rabbit/listener/SimpleMessageListenerContainerTests.java +++ b/spring-rabbit/src/test/java/org/springframework/amqp/rabbit/listener/SimpleMessageListenerContainerTests.java @@ -692,6 +692,7 @@ public class SimpleMessageListenerContainerTests { verify(channel).basicAck(2, true); container.stop(); verify(listener).containerAckMode(AcknowledgeMode.AUTO); + verify(listener).isAsyncReplies(); verifyNoMoreInteractions(listener); } diff --git a/src/reference/asciidoc/amqp.adoc b/src/reference/asciidoc/amqp.adoc index abf36a85..48d1c0b3 100644 --- a/src/reference/asciidoc/amqp.adoc +++ b/src/reference/asciidoc/amqp.adoc @@ -3110,6 +3110,9 @@ If the async result is completed with an `AmqpRejectAndDontRequeueException`, th If the container's `defaultRequeueRejected` property is `false`, you can override that by setting the future's exception to a `ImmediateRequeueException` and the message will be requeued. If some exception occurs within the listener method that prevents creation of the async result object, you MUST catch that exception and return an appropriate return object that will cause the message to be acknowledged or requeued. +Starting with versions 2.2.21, 2.3.13, 2.4.1, the `AcknowledgeMode` will be automatically set the `MANUAL` when async return types are detected. +In addition, incoming messages with fatal exceptions will be negatively acknowledged individually, previously any prior unacknowledged message were also negatively acknowledged. + [[threading]] ===== Threading and Asynchronous Consumers