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 fe5d7030f..d43ea0562 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 @@ -208,6 +208,7 @@ public class FaultTolerantStepFactoryBean extends SimpleStepFactoryBean extends SimpleStepFactoryBean rollbackClassifier = new Classifier() { - public Boolean classify(Throwable classifiable) { - return getTransactionAttribute().rollbackOn(classifiable); - } - }; - - BatchRetryTemplate batchRetryTemplate = configureRetry(); - - 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); + SimpleChunkProvider chunkProvider = configureChunkProvider(); + SimpleChunkProcessor chunkProcessor = configureChunkProcessor(); registerExplicitItemListeners(chunkProvider, chunkProcessor); registerImplicitItemListeners(chunkProvider, chunkProcessor); @@ -255,8 +234,47 @@ public class FaultTolerantStepFactoryBean extends SimpleStepFactoryBean configureChunkProcessor() { + + SkipPolicy writeSkipPolicy = new LimitCheckingItemSkipPolicy(skipLimit, skippableExceptionClasses, + fatalExceptionClasses); + + Classifier rollbackClassifier = new Classifier() { + public Boolean classify(Throwable classifiable) { + return getTransactionAttribute().rollbackOn(classifiable); + } + }; + + BatchRetryTemplate batchRetryTemplate = configureRetry(); + + FaultTolerantChunkProcessor chunkProcessor = new FaultTolerantChunkProcessor(getItemProcessor(), + getItemWriter(), batchRetryTemplate); + chunkProcessor.setBuffering(!isReaderTransactionalQueue); + chunkProcessor.setWriteSkipPolicy(writeSkipPolicy); + chunkProcessor.setProcessSkipPolicy(writeSkipPolicy); + chunkProcessor.setRollbackClassifier(rollbackClassifier); + + return chunkProcessor; + } + + /** + * @return {@link ChunkProvider} configured for fault-tolerance. + */ + private FaultTolerantChunkProvider configureChunkProvider() { + + SkipPolicy readSkipPolicy = new LimitCheckingItemSkipPolicy(skipLimit, skippableExceptionClasses, + fatalExceptionClasses); + FaultTolerantChunkProvider chunkProvider = new FaultTolerantChunkProvider(getItemReader(), + getChunkOperations()); + chunkProvider.setSkipPolicy(readSkipPolicy); + + return chunkProvider; + } + + /** + * @return fully configured retry template for item processing phase. */ private BatchRetryTemplate configureRetry() { @@ -303,8 +321,6 @@ public class FaultTolerantStepFactoryBean extends SimpleStepFactoryBean() {