AMQP-155: add throwable
This commit is contained in:
@@ -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;
|
||||
}
|
||||
|
||||
@@ -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<Void>() {
|
||||
public Void recover(Object[] args, Throwable cause) {
|
||||
Message message = (Message) args[1];
|
||||
messageRecoverer.recover(message);
|
||||
return null;
|
||||
retryInterceptor.setRecoverer(new MethodInvocationRecoverer<Void>() {
|
||||
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;
|
||||
|
||||
|
||||
@@ -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);
|
||||
|
||||
}
|
||||
|
||||
@@ -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();
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user