diff --git a/spring-rabbit/src/main/java/org/springframework/amqp/rabbit/config/StatefulRetryOperationsInterceptorFactoryBean.java b/spring-rabbit/src/main/java/org/springframework/amqp/rabbit/config/StatefulRetryOperationsInterceptorFactoryBean.java index efed12c4..98791e7a 100644 --- a/spring-rabbit/src/main/java/org/springframework/amqp/rabbit/config/StatefulRetryOperationsInterceptorFactoryBean.java +++ b/spring-rabbit/src/main/java/org/springframework/amqp/rabbit/config/StatefulRetryOperationsInterceptorFactoryBean.java @@ -84,9 +84,9 @@ public class StatefulRetryOperationsInterceptorFactoryBean extends AbstractRetry public Void recover(Object[] args, Throwable cause) { Message message = (Message) args[1]; if (messageRecoverer == null) { - logger.warn("Message dropped on recovery: " + message); + logger.warn("Message dropped on recovery: " + message, cause); } else { - messageRecoverer.recover(message); + messageRecoverer.recover(message, cause); } return null; } diff --git a/spring-rabbit/src/main/java/org/springframework/amqp/rabbit/config/StatelessRetryOperationsInterceptorFactoryBean.java b/spring-rabbit/src/main/java/org/springframework/amqp/rabbit/config/StatelessRetryOperationsInterceptorFactoryBean.java index a2aa7890..ad2183fb 100644 --- a/spring-rabbit/src/main/java/org/springframework/amqp/rabbit/config/StatelessRetryOperationsInterceptorFactoryBean.java +++ b/spring-rabbit/src/main/java/org/springframework/amqp/rabbit/config/StatelessRetryOperationsInterceptorFactoryBean.java @@ -12,6 +12,8 @@ */ package org.springframework.amqp.rabbit.config; +import org.apache.commons.logging.Log; +import org.apache.commons.logging.LogFactory; import org.springframework.amqp.core.Message; import org.springframework.amqp.rabbit.retry.MessageRecoverer; import org.springframework.retry.RetryOperations; @@ -36,6 +38,8 @@ import org.springframework.retry.support.RetryTemplate; */ public class StatelessRetryOperationsInterceptorFactoryBean extends AbstractRetryOperationsInterceptorFactoryBean { + private static Log logger = LogFactory.getLog(StatelessRetryOperationsInterceptorFactoryBean.class); + public RetryOperationsInterceptor getObject() { RetryOperationsInterceptor retryInterceptor = new RetryOperationsInterceptor(); @@ -46,15 +50,17 @@ public class StatelessRetryOperationsInterceptorFactoryBean extends AbstractRetr retryInterceptor.setRetryOperations(retryTemplate); final MessageRecoverer messageRecoverer = getMessageRecoverer(); - if (messageRecoverer != null) { - retryInterceptor.setRecoverer(new MethodInvocationRecoverer() { - public Void recover(Object[] args, Throwable cause) { - Message message = (Message) args[1]; - messageRecoverer.recover(message); - return null; + retryInterceptor.setRecoverer(new MethodInvocationRecoverer() { + public Void recover(Object[] args, Throwable cause) { + Message message = (Message) args[1]; + if (messageRecoverer == null) { + logger.warn("Message dropped on recovery: " + message, cause); + } else { + messageRecoverer.recover(message, cause); } - }); - } + return null; + } + }); return retryInterceptor; diff --git a/spring-rabbit/src/main/java/org/springframework/amqp/rabbit/retry/MessageRecoverer.java b/spring-rabbit/src/main/java/org/springframework/amqp/rabbit/retry/MessageRecoverer.java index 672307b8..33465c5e 100644 --- a/spring-rabbit/src/main/java/org/springframework/amqp/rabbit/retry/MessageRecoverer.java +++ b/spring-rabbit/src/main/java/org/springframework/amqp/rabbit/retry/MessageRecoverer.java @@ -24,7 +24,8 @@ public interface MessageRecoverer { * Callback for message that was consumed but failed all retry attempts. * * @param message the message to recover + * @param cause the cause of the error */ - void recover(Message message); + void recover(Message message, Throwable cause); } diff --git a/spring-rabbit/src/test/java/org/springframework/amqp/rabbit/listener/MessageListenerContainerRetryIntegrationTests.java b/spring-rabbit/src/test/java/org/springframework/amqp/rabbit/listener/MessageListenerContainerRetryIntegrationTests.java index f18b9dbc..ae07ecb7 100644 --- a/spring-rabbit/src/test/java/org/springframework/amqp/rabbit/listener/MessageListenerContainerRetryIntegrationTests.java +++ b/spring-rabbit/src/test/java/org/springframework/amqp/rabbit/listener/MessageListenerContainerRetryIntegrationTests.java @@ -143,7 +143,7 @@ public class MessageListenerContainerRetryIntegrationTests { factory = new StatelessRetryOperationsInterceptorFactoryBean(); } factory.setMessageRecoverer(new MessageRecoverer() { - public void recover(Message message) { + public void recover(Message message, Throwable cause) { logger.info("Recovered: " + message); latch.countDown(); }