RESOLVED - issue BATCH-544: Add a step factory that allows simplified injection of a RetryPolicy
This commit is contained in:
@@ -20,6 +20,7 @@ import org.springframework.batch.retry.RetryContext;
|
||||
import org.springframework.batch.retry.RetryException;
|
||||
import org.springframework.batch.retry.RetryListener;
|
||||
import org.springframework.batch.retry.RetryOperations;
|
||||
import org.springframework.batch.retry.RetryPolicy;
|
||||
import org.springframework.batch.retry.backoff.BackOffPolicy;
|
||||
import org.springframework.batch.retry.callback.RecoveryRetryCallback;
|
||||
import org.springframework.batch.retry.policy.ExceptionClassifierRetryPolicy;
|
||||
@@ -40,8 +41,8 @@ import org.springframework.batch.support.SubclassExceptionClassifier;
|
||||
* to be exclusive.
|
||||
*
|
||||
* Skippable exceptions on write will by default cause transaction rollback - to
|
||||
* avoid rollback for specific exception class include it in the
|
||||
* transaction attribute as "no rollback for".
|
||||
* avoid rollback for specific exception class include it in the transaction
|
||||
* attribute as "no rollback for".
|
||||
*
|
||||
* @see SimpleStepFactoryBean
|
||||
*
|
||||
@@ -61,7 +62,7 @@ public class SkipLimitStepFactoryBean extends SimpleStepFactoryBean {
|
||||
|
||||
private int cacheCapacity = 0;
|
||||
|
||||
private int retryLimit;
|
||||
private int retryLimit = 0;
|
||||
|
||||
private Class[] retryableExceptionClasses = new Class[] {};
|
||||
|
||||
@@ -69,10 +70,26 @@ public class SkipLimitStepFactoryBean extends SimpleStepFactoryBean {
|
||||
|
||||
private RetryListener[] retryListeners;
|
||||
|
||||
private RetryPolicy retryPolicy;
|
||||
|
||||
/**
|
||||
* Setter for the retry policy. If this is specified the other retry
|
||||
* properties are ignored (retryLimit, backOffPolicy,
|
||||
* retryableExceptionClasses).
|
||||
*
|
||||
* @param retryPolicy
|
||||
* a stateless {@link RetryPolicy}
|
||||
*/
|
||||
public void setRetryPolicy(RetryPolicy retryPolicy) {
|
||||
this.retryPolicy = retryPolicy;
|
||||
}
|
||||
|
||||
/**
|
||||
* Public setter for the retry limit. Each item can be retried up to this
|
||||
* limit.
|
||||
* @param retryLimit the retry limit to set
|
||||
*
|
||||
* @param retryLimit
|
||||
* the retry limit to set
|
||||
*/
|
||||
public void setRetryLimit(int retryLimit) {
|
||||
this.retryLimit = retryLimit;
|
||||
@@ -93,7 +110,8 @@ public class SkipLimitStepFactoryBean extends SimpleStepFactoryBean {
|
||||
* this many failures in a single transaction. Defaults to the value in the
|
||||
* {@link MapRetryContextCache}.
|
||||
*
|
||||
* @param cacheCapacity the cacheCapacity to set
|
||||
* @param cacheCapacity
|
||||
* the cacheCapacity to set
|
||||
*/
|
||||
public void setCacheCapacity(int cacheCapacity) {
|
||||
this.cacheCapacity = cacheCapacity;
|
||||
@@ -101,7 +119,9 @@ public class SkipLimitStepFactoryBean extends SimpleStepFactoryBean {
|
||||
|
||||
/**
|
||||
* Public setter for the Class[].
|
||||
* @param retryableExceptionClasses the retryableExceptionClasses to set
|
||||
*
|
||||
* @param retryableExceptionClasses
|
||||
* the retryableExceptionClasses to set
|
||||
*/
|
||||
public void setRetryableExceptionClasses(Class[] retryableExceptionClasses) {
|
||||
this.retryableExceptionClasses = retryableExceptionClasses;
|
||||
@@ -109,7 +129,9 @@ public class SkipLimitStepFactoryBean extends SimpleStepFactoryBean {
|
||||
|
||||
/**
|
||||
* Public setter for the {@link BackOffPolicy}.
|
||||
* @param backOffPolicy the {@link BackOffPolicy} to set
|
||||
*
|
||||
* @param backOffPolicy
|
||||
* the {@link BackOffPolicy} to set
|
||||
*/
|
||||
public void setBackOffPolicy(BackOffPolicy backOffPolicy) {
|
||||
this.backOffPolicy = backOffPolicy;
|
||||
@@ -117,7 +139,9 @@ public class SkipLimitStepFactoryBean extends SimpleStepFactoryBean {
|
||||
|
||||
/**
|
||||
* Public setter for the {@link RetryListener}s.
|
||||
* @param retryListeners the {@link RetryListener}s to set
|
||||
*
|
||||
* @param retryListeners
|
||||
* the {@link RetryListener}s to set
|
||||
*/
|
||||
public void setRetryListeners(RetryListener[] retryListeners) {
|
||||
this.retryListeners = retryListeners;
|
||||
@@ -130,7 +154,8 @@ public class SkipLimitStepFactoryBean extends SimpleStepFactoryBean {
|
||||
* zero then all exceptions will be propagated from the chunk and cause the
|
||||
* step to abort.
|
||||
*
|
||||
* @param skipLimit the value to set. Default is 0 (never skip).
|
||||
* @param skipLimit
|
||||
* the value to set. Default is 0 (never skip).
|
||||
*/
|
||||
public void setSkipLimit(int skipLimit) {
|
||||
this.skipLimit = skipLimit;
|
||||
@@ -141,7 +166,8 @@ public class SkipLimitStepFactoryBean extends SimpleStepFactoryBean {
|
||||
* but will result in transaction rollback and the item which handling
|
||||
* caused the exception will be skipped.
|
||||
*
|
||||
* @param exceptionClasses defaults to <code>Exception</code>
|
||||
* @param exceptionClasses
|
||||
* defaults to <code>Exception</code>
|
||||
*/
|
||||
public void setSkippableExceptionClasses(Class[] exceptionClasses) {
|
||||
this.skippableExceptionClasses = exceptionClasses;
|
||||
@@ -150,7 +176,8 @@ public class SkipLimitStepFactoryBean extends SimpleStepFactoryBean {
|
||||
/**
|
||||
* Public setter for exception classes that should cause immediate failure.
|
||||
*
|
||||
* @param fatalExceptionClasses {@link Error} by default
|
||||
* @param fatalExceptionClasses
|
||||
* {@link Error} by default
|
||||
*/
|
||||
public void setFatalExceptionClasses(Class[] fatalExceptionClasses) {
|
||||
this.fatalExceptionClasses = fatalExceptionClasses;
|
||||
@@ -161,7 +188,8 @@ public class SkipLimitStepFactoryBean extends SimpleStepFactoryBean {
|
||||
* failed items so they can be skipped if encountered again, generally in
|
||||
* another transaction.
|
||||
*
|
||||
* @param itemKeyGenerator the {@link ItemKeyGenerator} to set.
|
||||
* @param itemKeyGenerator
|
||||
* the {@link ItemKeyGenerator} to set.
|
||||
*/
|
||||
public void setItemKeyGenerator(ItemKeyGenerator itemKeyGenerator) {
|
||||
this.itemKeyGenerator = itemKeyGenerator;
|
||||
@@ -174,43 +202,57 @@ public class SkipLimitStepFactoryBean extends SimpleStepFactoryBean {
|
||||
protected void applyConfiguration(ItemOrientedStep step) {
|
||||
super.applyConfiguration(step);
|
||||
|
||||
if (retryLimit > 0 || skipLimit > 0) {
|
||||
if (retryLimit > 0 || skipLimit > 0 || retryPolicy != null) {
|
||||
|
||||
addFatalExceptionIfMissing(SkipLimitExceededException.class);
|
||||
addFatalExceptionIfMissing(RetryException.class);
|
||||
|
||||
SimpleRetryPolicy simpleRetryPolicy = new SimpleRetryPolicy(retryLimit);
|
||||
if (retryableExceptionClasses.length > 0) { // otherwise we retry
|
||||
// all exceptions
|
||||
simpleRetryPolicy.setRetryableExceptionClasses(retryableExceptionClasses);
|
||||
}
|
||||
simpleRetryPolicy.setFatalExceptionClasses(fatalExceptionClasses);
|
||||
if (retryPolicy == null) {
|
||||
|
||||
SimpleRetryPolicy simpleRetryPolicy = new SimpleRetryPolicy(
|
||||
retryLimit);
|
||||
if (retryableExceptionClasses.length > 0) { // otherwise we
|
||||
// retry
|
||||
// all exceptions
|
||||
simpleRetryPolicy
|
||||
.setRetryableExceptionClasses(retryableExceptionClasses);
|
||||
}
|
||||
simpleRetryPolicy
|
||||
.setFatalExceptionClasses(fatalExceptionClasses);
|
||||
|
||||
ExceptionClassifierRetryPolicy classifierRetryPolicy = new ExceptionClassifierRetryPolicy();
|
||||
SubclassExceptionClassifier exceptionClassifier = new SubclassExceptionClassifier();
|
||||
HashMap exceptionTypeMap = new HashMap();
|
||||
for (int i = 0; i < retryableExceptionClasses.length; i++) {
|
||||
Class cls = retryableExceptionClasses[i];
|
||||
exceptionTypeMap.put(cls, "retry");
|
||||
}
|
||||
exceptionClassifier.setTypeMap(exceptionTypeMap);
|
||||
HashMap retryPolicyMap = new HashMap();
|
||||
retryPolicyMap.put("retry", simpleRetryPolicy);
|
||||
retryPolicyMap.put("default", new NeverRetryPolicy());
|
||||
classifierRetryPolicy.setPolicyMap(retryPolicyMap);
|
||||
classifierRetryPolicy
|
||||
.setExceptionClassifier(exceptionClassifier);
|
||||
retryPolicy = classifierRetryPolicy;
|
||||
|
||||
ExceptionClassifierRetryPolicy retryPolicy = new ExceptionClassifierRetryPolicy();
|
||||
SubclassExceptionClassifier exceptionClassifier = new SubclassExceptionClassifier();
|
||||
HashMap exceptionTypeMap = new HashMap();
|
||||
for (int i = 0; i < retryableExceptionClasses.length; i++) {
|
||||
Class cls = retryableExceptionClasses[i];
|
||||
exceptionTypeMap.put(cls, "retry");
|
||||
}
|
||||
exceptionClassifier.setTypeMap(exceptionTypeMap);
|
||||
HashMap retryPolicyMap = new HashMap();
|
||||
retryPolicyMap.put("retry", simpleRetryPolicy);
|
||||
retryPolicyMap.put("default", new NeverRetryPolicy());
|
||||
retryPolicy.setPolicyMap(retryPolicyMap);
|
||||
retryPolicy.setExceptionClassifier(exceptionClassifier);
|
||||
|
||||
// Co-ordinate the retry policy with the exception handler:
|
||||
getStepOperations().setExceptionHandler(
|
||||
new SimpleRetryExceptionHandler(retryPolicy, getExceptionHandler(), fatalExceptionClasses));
|
||||
new SimpleRetryExceptionHandler(retryPolicy,
|
||||
getExceptionHandler(), fatalExceptionClasses));
|
||||
|
||||
RecoveryCallbackRetryPolicy recoveryCallbackRetryPolicy = new RecoveryCallbackRetryPolicy(retryPolicy) {
|
||||
RecoveryCallbackRetryPolicy recoveryCallbackRetryPolicy = new RecoveryCallbackRetryPolicy(
|
||||
retryPolicy) {
|
||||
protected boolean recoverForException(Throwable ex) {
|
||||
return !getTransactionAttribute().rollbackOn(ex);
|
||||
}
|
||||
};
|
||||
if (cacheCapacity > 0) {
|
||||
recoveryCallbackRetryPolicy.setRetryContextCache(new MapRetryContextCache(cacheCapacity));
|
||||
recoveryCallbackRetryPolicy
|
||||
.setRetryContextCache(new MapRetryContextCache(
|
||||
cacheCapacity));
|
||||
}
|
||||
|
||||
RetryTemplate retryTemplate = new RetryTemplate();
|
||||
@@ -218,36 +260,41 @@ public class SkipLimitStepFactoryBean extends SimpleStepFactoryBean {
|
||||
retryTemplate.setListeners(retryListeners);
|
||||
}
|
||||
retryTemplate.setRetryPolicy(recoveryCallbackRetryPolicy);
|
||||
if (backOffPolicy != null) {
|
||||
if (retryPolicy == null && backOffPolicy != null) {
|
||||
retryTemplate.setBackOffPolicy(backOffPolicy);
|
||||
}
|
||||
|
||||
List exceptions = new ArrayList(Arrays.asList(skippableExceptionClasses));
|
||||
ItemSkipPolicy readSkipPolicy = new LimitCheckingItemSkipPolicy(skipLimit, exceptions, Arrays
|
||||
.asList(fatalExceptionClasses));
|
||||
List exceptions = new ArrayList(Arrays
|
||||
.asList(skippableExceptionClasses));
|
||||
ItemSkipPolicy readSkipPolicy = new LimitCheckingItemSkipPolicy(
|
||||
skipLimit, exceptions, Arrays.asList(fatalExceptionClasses));
|
||||
exceptions.addAll(Arrays.asList(retryableExceptionClasses));
|
||||
ItemSkipPolicy writeSkipPolicy = new LimitCheckingItemSkipPolicy(skipLimit, exceptions, Arrays
|
||||
.asList(fatalExceptionClasses));
|
||||
StatefulRetryItemHandler itemHandler = new StatefulRetryItemHandler(getItemReader(), getItemWriter(),
|
||||
retryTemplate, itemKeyGenerator, readSkipPolicy, writeSkipPolicy);
|
||||
itemHandler.setSkipListeners(BatchListenerFactoryHelper.getSkipListeners(getListeners()));
|
||||
ItemSkipPolicy writeSkipPolicy = new LimitCheckingItemSkipPolicy(
|
||||
skipLimit, exceptions, Arrays.asList(fatalExceptionClasses));
|
||||
StatefulRetryItemHandler itemHandler = new StatefulRetryItemHandler(
|
||||
getItemReader(), getItemWriter(), retryTemplate,
|
||||
itemKeyGenerator, readSkipPolicy, writeSkipPolicy);
|
||||
itemHandler.setSkipListeners(BatchListenerFactoryHelper
|
||||
.getSkipListeners(getListeners()));
|
||||
|
||||
step.setItemHandler(itemHandler);
|
||||
|
||||
}
|
||||
else {
|
||||
} else {
|
||||
// This is the default in ItemOrientedStep anyway...
|
||||
step.setItemHandler(new SimpleItemHandler(getItemReader(), getItemWriter()));
|
||||
step.setItemHandler(new SimpleItemHandler(getItemReader(),
|
||||
getItemWriter()));
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
public void addFatalExceptionIfMissing(Class cls) {
|
||||
List fatalExceptionList = new ArrayList(Arrays.asList(fatalExceptionClasses));
|
||||
List fatalExceptionList = new ArrayList(Arrays
|
||||
.asList(fatalExceptionClasses));
|
||||
if (!fatalExceptionList.contains(cls)) {
|
||||
fatalExceptionList.add(cls);
|
||||
}
|
||||
fatalExceptionClasses = (Class[]) fatalExceptionList.toArray(new Class[0]);
|
||||
fatalExceptionClasses = (Class[]) fatalExceptionList
|
||||
.toArray(new Class[0]);
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -281,8 +328,10 @@ public class SkipLimitStepFactoryBean extends SimpleStepFactoryBean {
|
||||
* @param retryTemplate
|
||||
* @param itemKeyGenerator
|
||||
*/
|
||||
public StatefulRetryItemHandler(ItemReader itemReader, ItemWriter itemWriter, RetryOperations retryTemplate,
|
||||
ItemKeyGenerator itemKeyGenerator, ItemSkipPolicy readSkipPolicy, ItemSkipPolicy writeSkipPolicy) {
|
||||
public StatefulRetryItemHandler(ItemReader itemReader,
|
||||
ItemWriter itemWriter, RetryOperations retryTemplate,
|
||||
ItemKeyGenerator itemKeyGenerator,
|
||||
ItemSkipPolicy readSkipPolicy, ItemSkipPolicy writeSkipPolicy) {
|
||||
super(itemReader, itemWriter);
|
||||
this.retryOperations = retryTemplate;
|
||||
this.itemKeyGenerator = itemKeyGenerator;
|
||||
@@ -307,7 +356,8 @@ public class SkipLimitStepFactoryBean extends SimpleStepFactoryBean {
|
||||
* Register a listener for callbacks at the appropriate stages in a skip
|
||||
* process.
|
||||
*
|
||||
* @param listener a {@link SkipListener}
|
||||
* @param listener
|
||||
* a {@link SkipListener}
|
||||
*/
|
||||
public void registerSkipListener(SkipListener listener) {
|
||||
this.listener.register(listener);
|
||||
@@ -317,8 +367,8 @@ public class SkipLimitStepFactoryBean extends SimpleStepFactoryBean {
|
||||
* Tries to read the item from the reader, in case of exception skip the
|
||||
* item if the skip policy allows, otherwise re-throw.
|
||||
*
|
||||
* @param contribution current StepContribution holding skipped items
|
||||
* count
|
||||
* @param contribution
|
||||
* current StepContribution holding skipped items count
|
||||
* @return next item for processing
|
||||
*/
|
||||
protected Object read(StepContribution contribution) throws Exception {
|
||||
@@ -326,22 +376,20 @@ public class SkipLimitStepFactoryBean extends SimpleStepFactoryBean {
|
||||
while (true) {
|
||||
try {
|
||||
return doRead();
|
||||
}
|
||||
catch (Exception e) {
|
||||
} catch (Exception e) {
|
||||
try {
|
||||
if (readSkipPolicy.shouldSkip(e, contribution.getStepSkipCount())) {
|
||||
if (readSkipPolicy.shouldSkip(e, contribution
|
||||
.getStepSkipCount())) {
|
||||
// increment skip count and try again
|
||||
contribution.incrementTemporaryReadSkipCount();
|
||||
listener.onSkipInRead(e);
|
||||
logger.debug("Skipping failed input", e);
|
||||
}
|
||||
else {
|
||||
} else {
|
||||
// re-throw only when the skip policy runs out of
|
||||
// patience
|
||||
throw e;
|
||||
}
|
||||
}
|
||||
catch (SkipLimitExceededException ex) {
|
||||
} catch (SkipLimitExceededException ex) {
|
||||
// we are headed for a abnormal ending so bake in the
|
||||
// skip count
|
||||
contribution.combineSkipCounts();
|
||||
@@ -361,19 +409,24 @@ public class SkipLimitStepFactoryBean extends SimpleStepFactoryBean {
|
||||
* the next transaction automatically.<br/>
|
||||
*
|
||||
* @see org.springframework.batch.core.step.item.SimpleItemHandler#write(java.lang.Object,
|
||||
* org.springframework.batch.core.StepContribution)
|
||||
* org.springframework.batch.core.StepContribution)
|
||||
*/
|
||||
protected void write(final Object item, final StepContribution contribution) throws Exception {
|
||||
RecoveryRetryCallback retryCallback = new RecoveryRetryCallback(item, new RetryCallback() {
|
||||
public Object doWithRetry(RetryContext context) throws Throwable {
|
||||
doWrite(item);
|
||||
return null;
|
||||
}
|
||||
}, itemKeyGenerator != null ? itemKeyGenerator.getKey(item) : item);
|
||||
protected void write(final Object item,
|
||||
final StepContribution contribution) throws Exception {
|
||||
RecoveryRetryCallback retryCallback = new RecoveryRetryCallback(
|
||||
item, new RetryCallback() {
|
||||
public Object doWithRetry(RetryContext context)
|
||||
throws Throwable {
|
||||
doWrite(item);
|
||||
return null;
|
||||
}
|
||||
}, itemKeyGenerator != null ? itemKeyGenerator.getKey(item)
|
||||
: item);
|
||||
retryCallback.setRecoveryCallback(new RecoveryCallback() {
|
||||
public Object recover(RetryContext context) {
|
||||
Throwable t = context.getLastThrowable();
|
||||
if (writeSkipPolicy.shouldSkip(t, contribution.getStepSkipCount())) {
|
||||
if (writeSkipPolicy.shouldSkip(t, contribution
|
||||
.getStepSkipCount())) {
|
||||
listener.onSkipInWrite(item, t);
|
||||
}
|
||||
contribution.incrementWriteSkipCount();
|
||||
|
||||
@@ -45,6 +45,7 @@ import org.springframework.batch.item.support.AbstractItemWriter;
|
||||
import org.springframework.batch.item.support.ListItemReader;
|
||||
import org.springframework.batch.retry.RetryException;
|
||||
import org.springframework.batch.retry.policy.RetryCacheCapacityExceededException;
|
||||
import org.springframework.batch.retry.policy.SimpleRetryPolicy;
|
||||
import org.springframework.batch.support.transaction.ResourcelessTransactionManager;
|
||||
import org.springframework.batch.support.transaction.TransactionAwareProxyFactory;
|
||||
import org.springframework.transaction.support.TransactionSynchronizationManager;
|
||||
@@ -259,6 +260,43 @@ public class StatefulRetryStepFactoryBeanTests extends TestCase {
|
||||
assertEquals(0, stepExecution.getItemCount().intValue());
|
||||
}
|
||||
|
||||
public void testRetryPolicy() throws Exception {
|
||||
factory.setRetryPolicy(new SimpleRetryPolicy(4));
|
||||
factory.setSkipLimit(0);
|
||||
List items = TransactionAwareProxyFactory.createTransactionalList();
|
||||
items.addAll(Arrays.asList(new String[] { "b" }));
|
||||
ItemReader provider = new ListItemReader(items) {
|
||||
public Object read() {
|
||||
Object item = super.read();
|
||||
count++;
|
||||
return item;
|
||||
}
|
||||
};
|
||||
ItemWriter itemWriter = new AbstractItemWriter() {
|
||||
public void write(Object item) throws Exception {
|
||||
logger.debug("Write Called! Item: [" + item + "]");
|
||||
throw new RuntimeException("Write error - planned but retryable.");
|
||||
}
|
||||
};
|
||||
factory.setItemReader(provider);
|
||||
factory.setItemWriter(itemWriter);
|
||||
AbstractStep step = (AbstractStep) factory.getObject();
|
||||
|
||||
StepExecution stepExecution = new StepExecution(step.getName(), jobExecution);
|
||||
try {
|
||||
step.execute(stepExecution);
|
||||
fail("Expected SkipLimitExceededException");
|
||||
}
|
||||
catch (SkipLimitExceededException e) {
|
||||
// expected
|
||||
}
|
||||
|
||||
assertEquals(0, stepExecution.getSkipCount());
|
||||
// b is processed 4 times plus the null at end
|
||||
assertEquals(5, count);
|
||||
assertEquals(0, stepExecution.getItemCount().intValue());
|
||||
}
|
||||
|
||||
public void testCacheLimitWithRetry() throws Exception {
|
||||
factory.setRetryableExceptionClasses(new Class[] { Exception.class });
|
||||
factory.setRetryLimit(2);
|
||||
|
||||
Reference in New Issue
Block a user