diff --git a/spring-batch-core/src/main/java/org/springframework/batch/core/step/item/FaultTolerantStepFactoryBean.java b/spring-batch-core/src/main/java/org/springframework/batch/core/step/item/FaultTolerantStepFactoryBean.java index 31d7b0192..2402ffee2 100644 --- a/spring-batch-core/src/main/java/org/springframework/batch/core/step/item/FaultTolerantStepFactoryBean.java +++ b/spring-batch-core/src/main/java/org/springframework/batch/core/step/item/FaultTolerantStepFactoryBean.java @@ -172,7 +172,7 @@ public class FaultTolerantStepFactoryBean extends SimpleStepFactoryBean extends SimpleStepFactoryBean 1 || skipLimit > 0 || retryPolicy != null) { + if (!(retryLimit > 1 || skipLimit > 0 || retryPolicy != null)) { + // zero fault-tolerance, just use the parent's simple config + return; + } - addFatalExceptionIfMissing(SkipLimitExceededException.class); - addFatalExceptionIfMissing(NonSkippableReadException.class); - addFatalExceptionIfMissing(SkipListenerFailedException.class); - addFatalExceptionIfMissing(RetryException.class); + addFatalExceptionIfMissing(SkipLimitExceededException.class, NonSkippableReadException.class, + SkipListenerFailedException.class, RetryException.class); - if (retryPolicy == null) { - - SimpleRetryPolicy simpleRetryPolicy = new SimpleRetryPolicy(retryLimit); - if (!retryableExceptionClasses.isEmpty()) { - // otherwise we retry all exceptions - simpleRetryPolicy.setRetryableExceptionClasses(retryableExceptionClasses); - } - - retryPolicy = simpleRetryPolicy; + if (retryPolicy == null) { + SimpleRetryPolicy simpleRetryPolicy = new SimpleRetryPolicy(retryLimit); + if (!retryableExceptionClasses.isEmpty()) { + // otherwise we retry all exceptions + simpleRetryPolicy.setRetryableExceptionClasses(retryableExceptionClasses); } - // wrapper of the injected retry policy takes care of fatal - // exceptions (never retried) - final NeverRetryPolicy neverRetryPolicy = new NeverRetryPolicy(); - ExceptionClassifierRetryPolicy retryPolicyWrapper = new ExceptionClassifierRetryPolicy(); - retryPolicyWrapper.setExceptionClassifier(new Classifier() { - - public RetryPolicy classify(Throwable classifiable) { - - for (Class fatal : fatalExceptionClasses) { - if (fatal.isAssignableFrom(classifiable.getClass())) { - return neverRetryPolicy; - } - } - return retryPolicy; - } - }); - BatchRetryTemplate batchRetryTemplate = new BatchRetryTemplate(); - if (backOffPolicy != null) { - batchRetryTemplate.setBackOffPolicy(backOffPolicy); - } - batchRetryTemplate.setRetryPolicy(retryPolicyWrapper); - - // Co-ordinate the retry policy with the exception handler: - RepeatOperations stepOperations = getStepOperations(); - if (stepOperations instanceof RepeatTemplate) { - SimpleRetryExceptionHandler exceptionHandler = new SimpleRetryExceptionHandler(retryPolicyWrapper, - getExceptionHandler(), fatalExceptionClasses); - ((RepeatTemplate) stepOperations).setExceptionHandler(exceptionHandler); - } - - if (retryContextCache == null) { - if (cacheCapacity > 0) { - batchRetryTemplate.setRetryContextCache(new MapRetryContextCache(cacheCapacity)); - } - } - else { - batchRetryTemplate.setRetryContextCache(retryContextCache); - } - - if (retryListeners != null) { - batchRetryTemplate.setListeners(retryListeners); - } - - List> exceptions = new ArrayList>( - skippableExceptionClasses); - SkipPolicy readSkipPolicy = new LimitCheckingItemSkipPolicy(skipLimit, skippableExceptionClasses, - new ArrayList>(fatalExceptionClasses)); - exceptions.addAll(new ArrayList>(retryableExceptionClasses)); - SkipPolicy writeSkipPolicy = new LimitCheckingItemSkipPolicy(skipLimit, exceptions, - new ArrayList>(fatalExceptionClasses)); - - Classifier rollbackClassifier = new Classifier() { - public Boolean classify(Throwable classifiable) { - return getTransactionAttribute().rollbackOn(classifiable); - } - }; - - FaultTolerantChunkProvider chunkProvider = new FaultTolerantChunkProvider(getItemReader(), - getChunkOperations()); - chunkProvider.setSkipPolicy(readSkipPolicy); - chunkProvider.setListeners(BatchListenerFactoryHelper.> getListeners(getListeners(), - ItemReadListener.class)); - chunkProvider.setListeners(BatchListenerFactoryHelper.> getListeners(getListeners(), - SkipListener.class)); - - FaultTolerantChunkProcessor chunkProcessor = new FaultTolerantChunkProcessor( - getItemProcessor(), getItemWriter(), batchRetryTemplate); - chunkProcessor.setBuffering(!isReaderTransactionalQueue); - chunkProcessor.setWriteSkipPolicy(writeSkipPolicy); - chunkProcessor.setProcessSkipPolicy(writeSkipPolicy); - chunkProcessor.setRollbackClassifier(rollbackClassifier); - chunkProcessor.setListeners(BatchListenerFactoryHelper.> getListeners( - getListeners(), ItemProcessListener.class)); - chunkProcessor.setListeners(BatchListenerFactoryHelper.> getListeners(getListeners(), - ItemWriteListener.class)); - chunkProcessor.setListeners(BatchListenerFactoryHelper.> getListeners(getListeners(), - SkipListener.class)); - - - for (Object itemHandler : new Object[] { getItemReader(), getItemWriter(), getItemProcessor() }) { - - if (itemHandler instanceof SkipListener) { - chunkProvider.registerListener((StepListener) itemHandler); - chunkProcessor.registerListener((StepListener) itemHandler); - // already registered with both so avoid double-registering - continue; - } - if (itemHandler instanceof ItemReadListener) { - chunkProvider.registerListener((StepListener) itemHandler); - } - if (itemHandler instanceof ItemProcessListener || itemHandler instanceof ItemWriteListener) { - chunkProcessor.registerListener((StepListener) itemHandler); - } - } - ChunkOrientedTasklet tasklet = new ChunkOrientedTasklet(chunkProvider, chunkProcessor); - tasklet.setBuffering(!isReaderTransactionalQueue); - - step.setTasklet(tasklet); + retryPolicy = simpleRetryPolicy; } + // wrapper of the injected retry policy takes care of fatal + // exceptions (never retried) + final NeverRetryPolicy neverRetryPolicy = new NeverRetryPolicy(); + ExceptionClassifierRetryPolicy retryPolicyWrapper = new ExceptionClassifierRetryPolicy(); + retryPolicyWrapper.setExceptionClassifier(new Classifier() { + + public RetryPolicy classify(Throwable classifiable) { + + for (Class fatal : fatalExceptionClasses) { + if (fatal.isAssignableFrom(classifiable.getClass())) { + return neverRetryPolicy; + } + } + return retryPolicy; + } + }); + BatchRetryTemplate batchRetryTemplate = new BatchRetryTemplate(); + if (backOffPolicy != null) { + batchRetryTemplate.setBackOffPolicy(backOffPolicy); + } + batchRetryTemplate.setRetryPolicy(retryPolicyWrapper); + + // Co-ordinate the retry policy with the exception handler: + RepeatOperations stepOperations = getStepOperations(); + if (stepOperations instanceof RepeatTemplate) { + SimpleRetryExceptionHandler exceptionHandler = new SimpleRetryExceptionHandler(retryPolicyWrapper, + getExceptionHandler(), fatalExceptionClasses); + ((RepeatTemplate) stepOperations).setExceptionHandler(exceptionHandler); + } + + if (retryContextCache == null) { + if (cacheCapacity > 0) { + batchRetryTemplate.setRetryContextCache(new MapRetryContextCache(cacheCapacity)); + } + } + else { + batchRetryTemplate.setRetryContextCache(retryContextCache); + } + + if (retryListeners != null) { + batchRetryTemplate.setListeners(retryListeners); + } + + List> exceptions = new ArrayList>( + skippableExceptionClasses); + SkipPolicy readSkipPolicy = new LimitCheckingItemSkipPolicy(skipLimit, skippableExceptionClasses, + new ArrayList>(fatalExceptionClasses)); + exceptions.addAll(new ArrayList>(retryableExceptionClasses)); + SkipPolicy writeSkipPolicy = new LimitCheckingItemSkipPolicy(skipLimit, exceptions, + new ArrayList>(fatalExceptionClasses)); + + Classifier rollbackClassifier = new Classifier() { + public Boolean classify(Throwable classifiable) { + return getTransactionAttribute().rollbackOn(classifiable); + } + }; + + FaultTolerantChunkProvider chunkProvider = new FaultTolerantChunkProvider(getItemReader(), + getChunkOperations()); + chunkProvider.setSkipPolicy(readSkipPolicy); + FaultTolerantChunkProcessor chunkProcessor = new FaultTolerantChunkProcessor(getItemProcessor(), + getItemWriter(), batchRetryTemplate); + chunkProcessor.setBuffering(!isReaderTransactionalQueue); + chunkProcessor.setWriteSkipPolicy(writeSkipPolicy); + chunkProcessor.setProcessSkipPolicy(writeSkipPolicy); + chunkProcessor.setRollbackClassifier(rollbackClassifier); + + registerItemListeners(chunkProvider, chunkProcessor); + autoRegisterItemListeners(chunkProvider, chunkProcessor); + + ChunkOrientedTasklet tasklet = new ChunkOrientedTasklet(chunkProvider, chunkProcessor); + tasklet.setBuffering(!isReaderTransactionalQueue); + + step.setTasklet(tasklet); + } - public void addFatalExceptionIfMissing(Class cls) { - List> fatalExceptionList = new ArrayList>(); - for (Class exceptionClass : fatalExceptionClasses) { + /** + * Register injected item listeners. + */ + private void registerItemListeners(SimpleChunkProvider chunkProvider, SimpleChunkProcessor chunkProcessor) { + + chunkProvider.setListeners(BatchListenerFactoryHelper.> getListeners(getListeners(), + ItemReadListener.class)); + chunkProvider.setListeners(BatchListenerFactoryHelper.> getListeners(getListeners(), + SkipListener.class)); + + chunkProcessor.setListeners(BatchListenerFactoryHelper.> getListeners(getListeners(), + ItemProcessListener.class)); + chunkProcessor.setListeners(BatchListenerFactoryHelper.> getListeners(getListeners(), + ItemWriteListener.class)); + chunkProcessor.setListeners(BatchListenerFactoryHelper.> getListeners(getListeners(), + SkipListener.class)); + } + + /** + * Auto-register reader, processor and writer as item listeners if applicable + */ + private void autoRegisterItemListeners(SimpleChunkProvider chunkProvider, + SimpleChunkProcessor chunkProcessor) { + for (Object itemHandler : new Object[] { getItemReader(), getItemWriter(), getItemProcessor() }) { + + if (itemHandler instanceof SkipListener) { + chunkProvider.registerListener((StepListener) itemHandler); + chunkProcessor.registerListener((StepListener) itemHandler); + // already registered with both so avoid double-registering + continue; + } + if (itemHandler instanceof ItemReadListener) { + chunkProvider.registerListener((StepListener) itemHandler); + } + if (itemHandler instanceof ItemProcessListener || itemHandler instanceof ItemWriteListener) { + chunkProcessor.registerListener((StepListener) itemHandler); + } + } + } + + @SuppressWarnings("unchecked") + private void addFatalExceptionIfMissing(Class... cls) { + List fatalExceptionList = new ArrayList>(); + for (Class exceptionClass : fatalExceptionClasses) { fatalExceptionList.add(exceptionClass); } - if (!fatalExceptionList.contains(cls)) { - fatalExceptionList.add(cls); + for (Class fatal : cls) { + if (!fatalExceptionList.contains(fatal)) { + fatalExceptionList.add(fatal); + } } fatalExceptionClasses = fatalExceptionList; }