diff --git a/spring-batch-core/src/main/java/org/springframework/batch/core/step/item/AbstractFaultTolerantChunkOrientedTasklet.java b/spring-batch-core/src/main/java/org/springframework/batch/core/step/item/AbstractFaultTolerantChunkOrientedTasklet.java index 54a978abd..70e1e8c7f 100644 --- a/spring-batch-core/src/main/java/org/springframework/batch/core/step/item/AbstractFaultTolerantChunkOrientedTasklet.java +++ b/spring-batch-core/src/main/java/org/springframework/batch/core/step/item/AbstractFaultTolerantChunkOrientedTasklet.java @@ -36,6 +36,12 @@ import org.springframework.core.AttributeAccessor; * @author Robert Kasanicky */ public abstract class AbstractFaultTolerantChunkOrientedTasklet extends AbstractItemOrientedTasklet { + + final static protected String SKIPPED_INPUTS_KEY = "SKIPPED_INPUTS_KEY"; + + final static protected String SKIPPED_OUTPUTS_KEY = "SKIPPED_OUTPUTS_KEY"; + + final static protected String SKIPPED_READS_KEY = "SKIPPED_READS_KEY"; final private RetryOperations retryOperations; @@ -43,20 +49,27 @@ public abstract class AbstractFaultTolerantChunkOrientedTasklet extends Ab final private ItemSkipPolicy processSkipPolicy; + final private ItemSkipPolicy readSkipPolicy; + final private Classifier rollbackClassifier; public AbstractFaultTolerantChunkOrientedTasklet(ItemReader itemReader, ItemProcessor itemProcessor, ItemWriter itemWriter, - RetryOperations retryOperations, ItemSkipPolicy processSkipPolicy, ItemSkipPolicy writeSkipPolicy, - Classifier rollbackClassifier) { + RetryOperations retryOperations, ItemSkipPolicy readSkipPolicy, ItemSkipPolicy processSkipPolicy, + ItemSkipPolicy writeSkipPolicy, Classifier rollbackClassifier) { super(itemReader, itemProcessor, itemWriter); this.retryOperations = retryOperations; + this.readSkipPolicy = readSkipPolicy; this.processSkipPolicy = processSkipPolicy; this.writeSkipPolicy = writeSkipPolicy; this.rollbackClassifier = rollbackClassifier; } + protected ItemSkipPolicy getReadSkipPolicy() { + return readSkipPolicy; + } + /** * Call all skip listeners in read-process-write order * @param skippedReads read exceptions diff --git a/spring-batch-core/src/main/java/org/springframework/batch/core/step/item/FaultTolerantChunkOrientedTasklet.java b/spring-batch-core/src/main/java/org/springframework/batch/core/step/item/FaultTolerantChunkOrientedTasklet.java index 27b9aaca5..6b0e21552 100644 --- a/spring-batch-core/src/main/java/org/springframework/batch/core/step/item/FaultTolerantChunkOrientedTasklet.java +++ b/spring-batch-core/src/main/java/org/springframework/batch/core/step/item/FaultTolerantChunkOrientedTasklet.java @@ -42,7 +42,7 @@ import org.springframework.core.AttributeAccessor; * listener is invoked and the skip count incremented. A retryable exception is * thus also effectively also implicitly skippable. * - * ItemProcessor is assumed to be transactional. In case of rollback caused by + * ItemProcessor is assumed to be transactional. In case of rollback caused by * error on write the processing phase will be repeated. * * @author Dave Syer @@ -50,17 +50,9 @@ import org.springframework.core.AttributeAccessor; */ public class FaultTolerantChunkOrientedTasklet extends AbstractFaultTolerantChunkOrientedTasklet { - private static final String INPUT_BUFFER_KEY = "INPUT_BUFFER_KEY"; + final static private String INPUT_BUFFER_KEY = "INPUT_BUFFER_KEY"; - private final RepeatOperations repeatOperations; - - final private ItemSkipPolicy readSkipPolicy; - - private static final String SKIPPED_OUTPUTS_KEY = "SKIPPED_OUTPUTS_BUFFER_KEY"; - - private static final String SKIPPED_INPUTS_KEY = "SKIPPED_INPUTS_BUFFER_KEY"; - - private static final String SKIPPED_READS_KEY = "SKIPPED_READS_BUFFER_KEY"; + final private RepeatOperations repeatOperations; public FaultTolerantChunkOrientedTasklet(ItemReader itemReader, ItemProcessor itemProcessor, ItemWriter itemWriter, @@ -68,10 +60,9 @@ public class FaultTolerantChunkOrientedTasklet extends AbstractFaultTolera Classifier rollbackClassifier, ItemSkipPolicy readSkipPolicy, ItemSkipPolicy writeSkipPolicy, ItemSkipPolicy processSkipPolicy) { - super(itemReader, itemProcessor, itemWriter, retryTemplate, processSkipPolicy, writeSkipPolicy, + super(itemReader, itemProcessor, itemWriter, retryTemplate, readSkipPolicy, processSkipPolicy, writeSkipPolicy, rollbackClassifier); this.repeatOperations = chunkOperations; - this.readSkipPolicy = readSkipPolicy; } /** @@ -156,7 +147,7 @@ public class FaultTolerantChunkOrientedTasklet extends AbstractFaultTolera } catch (Exception e) { - if (readSkipPolicy.shouldSkip(e, contribution.getStepSkipCount())) { + if (getReadSkipPolicy().shouldSkip(e, contribution.getStepSkipCount())) { // increment skip count and try again contribution.incrementReadSkipCount(); skippedReads.add(e); diff --git a/spring-batch-core/src/main/java/org/springframework/batch/core/step/item/NonbufferingFaultTolerantChunkOrientedTasklet.java b/spring-batch-core/src/main/java/org/springframework/batch/core/step/item/NonbufferingFaultTolerantChunkOrientedTasklet.java index bde7f0f46..116b1e273 100644 --- a/spring-batch-core/src/main/java/org/springframework/batch/core/step/item/NonbufferingFaultTolerantChunkOrientedTasklet.java +++ b/spring-batch-core/src/main/java/org/springframework/batch/core/step/item/NonbufferingFaultTolerantChunkOrientedTasklet.java @@ -34,15 +34,7 @@ import org.springframework.core.AttributeAccessor; public class NonbufferingFaultTolerantChunkOrientedTasklet extends AbstractFaultTolerantChunkOrientedTasklet { - private static final String SKIPPED_INPUTS_KEY = "SKIPPED_INPUTS_KEY"; - - private static final String SKIPPED_OUTPUTS_KEY = "SKIPPED_OUTPUTS_KEY"; - - private static final String SKIPPED_READS_KEY = "SKIPPED_READS_KEY"; - - private final RepeatOperations repeatOperations; - - private final ItemSkipPolicy readSkipPolicy; + final private RepeatOperations repeatOperations; public NonbufferingFaultTolerantChunkOrientedTasklet(ItemReader itemReader, ItemProcessor itemProcessor, ItemWriter itemWriter, @@ -50,10 +42,9 @@ public class NonbufferingFaultTolerantChunkOrientedTasklet extends Classifier rollbackClassifier, ItemSkipPolicy readSkipPolicy, ItemSkipPolicy writeSkipPolicy, ItemSkipPolicy processSkipPolicy) { - super(itemReader, itemProcessor, itemWriter, retryTemplate, processSkipPolicy, writeSkipPolicy, + super(itemReader, itemProcessor, itemWriter, retryTemplate, readSkipPolicy, processSkipPolicy, writeSkipPolicy, rollbackClassifier); this.repeatOperations = chunkOperations; - this.readSkipPolicy = readSkipPolicy; } /** @@ -118,7 +109,7 @@ public class NonbufferingFaultTolerantChunkOrientedTasklet extends } catch (Exception e) { - if (readSkipPolicy.shouldSkip(e, contribution.getStepSkipCount())) { + if (getReadSkipPolicy().shouldSkip(e, contribution.getStepSkipCount())) { // increment skip count and try again contribution.incrementReadSkipCount(); skipped.add(e); diff --git a/spring-batch-core/src/test/java/org/springframework/batch/core/step/item/FaultTolerantChunkOrientedTaskletTests.java b/spring-batch-core/src/test/java/org/springframework/batch/core/step/item/FaultTolerantChunkOrientedTaskletTests.java index 379e25c35..656ef10dd 100644 --- a/spring-batch-core/src/test/java/org/springframework/batch/core/step/item/FaultTolerantChunkOrientedTaskletTests.java +++ b/spring-batch-core/src/test/java/org/springframework/batch/core/step/item/FaultTolerantChunkOrientedTaskletTests.java @@ -161,7 +161,7 @@ public class FaultTolerantChunkOrientedTaskletTests { catch (Exception e) { assertEquals("Barf!", e.getMessage()); } - assertTrue(attributes.hasAttribute("SKIPPED_OUTPUTS_BUFFER_KEY")); + assertTrue(attributes.hasAttribute("SKIPPED_OUTPUTS_KEY")); tasklet.execute(contribution, attributes); assertEquals(1, contribution.getReadCount()); assertEquals(1, contribution.getWriteSkipCount()); @@ -191,10 +191,10 @@ public class FaultTolerantChunkOrientedTaskletTests { catch (Exception e) { assertEquals("Barf!", e.getMessage()); } - assertTrue(attributes.hasAttribute("SKIPPED_OUTPUTS_BUFFER_KEY")); + assertTrue(attributes.hasAttribute("SKIPPED_OUTPUTS_KEY")); } @SuppressWarnings("unchecked") - Map skips = (Map) attributes.getAttribute("SKIPPED_OUTPUTS_BUFFER_KEY"); + Map skips = (Map) attributes.getAttribute("SKIPPED_OUTPUTS_KEY"); assertEquals(1, skips.size()); // The last recovery for this chunk... tasklet.execute(contribution, attributes); @@ -214,7 +214,7 @@ public class FaultTolerantChunkOrientedTaskletTests { catch (SkipLimitExceededException e) { // expected } - assertTrue(attributes.hasAttribute("SKIPPED_OUTPUTS_BUFFER_KEY")); + assertTrue(attributes.hasAttribute("SKIPPED_OUTPUTS_KEY")); assertEquals(3, contribution.getReadCount()); assertEquals(0, contribution.getFilterCount()); assertEquals(2, contribution.getWriteSkipCount()); @@ -285,7 +285,7 @@ public class FaultTolerantChunkOrientedTaskletTests { assertTrue(attributes.hasAttribute("INPUT_BUFFER_KEY")); } @SuppressWarnings("unchecked") - Map skips = (Map) attributes.getAttribute("SKIPPED_INPUTS_BUFFER_KEY"); + Map skips = (Map) attributes.getAttribute("SKIPPED_INPUTS_KEY"); assertEquals(1, skips.size()); // The last recovery for this chunk...