diff --git a/spring-batch-core/src/main/java/org/springframework/batch/core/step/item/FaultTolerantChunkProcessor.java b/spring-batch-core/src/main/java/org/springframework/batch/core/step/item/FaultTolerantChunkProcessor.java index 2356fe6ab..ab2cdccc8 100755 --- a/spring-batch-core/src/main/java/org/springframework/batch/core/step/item/FaultTolerantChunkProcessor.java +++ b/spring-batch-core/src/main/java/org/springframework/batch/core/step/item/FaultTolerantChunkProcessor.java @@ -19,6 +19,7 @@ package org.springframework.batch.core.step.item; import java.util.ArrayList; import java.util.Collections; import java.util.List; +import java.util.concurrent.atomic.AtomicInteger; import org.apache.commons.logging.Log; import org.apache.commons.logging.LogFactory; @@ -39,7 +40,7 @@ import org.springframework.batch.retry.support.DefaultRetryState; /** * FaultTolerant implementation of the {@link ChunkProcessor} interface, that - * allows for skipping or retry of items that cause exceptions during writing. + * allows for skipping or retry of items that cause exceptions during writing. * */ public class FaultTolerantChunkProcessor extends SimpleChunkProcessor { @@ -124,6 +125,11 @@ public class FaultTolerantChunkProcessor extends SimpleChunkProcessor transform(final StepContribution contribution, Chunk inputs) throws Exception { Chunk outputs = new Chunk(); + Object userData = inputs.getUserData(); + @SuppressWarnings("unchecked") + final Chunk cache = (userData instanceof Chunk) ? (Chunk) userData : null; + final Chunk.ChunkIterator cacheIterator = (cache != null) ? cache.iterator() : null; + final AtomicInteger count = new AtomicInteger(0); for (final Chunk.ChunkIterator iterator = inputs.iterator(); iterator.hasNext();) { @@ -134,7 +140,21 @@ public class FaultTolerantChunkProcessor extends SimpleChunkProcessor 1) { + /* + * If there is a cached chunk then we must be + * scanning for errors in the writer, in which case + * only the first one will be written, and for the + * rest we need to fill in the output from the + * cache. + */ + output = cached; + } + else { + output = doProcess(item); + } } catch (Exception e) { if (rollbackClassifier.classify(e)) { @@ -259,10 +279,10 @@ public class FaultTolerantChunkProcessor extends SimpleChunkProcessor.ChunkIterator inputIterator = inputs.iterator(); Chunk.ChunkIterator outputIterator = outputs.iterator(); checkSkipPolicy(inputIterator, outputIterator, le, contribution); diff --git a/spring-batch-core/src/main/java/org/springframework/batch/core/step/item/SimpleChunkProcessor.java b/spring-batch-core/src/main/java/org/springframework/batch/core/step/item/SimpleChunkProcessor.java index 799c0aa1d..56fa6e964 100755 --- a/spring-batch-core/src/main/java/org/springframework/batch/core/step/item/SimpleChunkProcessor.java +++ b/spring-batch-core/src/main/java/org/springframework/batch/core/step/item/SimpleChunkProcessor.java @@ -183,7 +183,9 @@ public class SimpleChunkProcessor implements ChunkProcessor, Initializi contribution.incrementFilterCount(inputsSize - outputs.size() - inputs.getSkips().size()); + boolean busy = skips.isBusy(); outputs = new Chunk(outputs.getItems(), skips.getSkips()); + outputs.setBusy(busy); // Remember for next time if there are skips accumulating inputs.setUserData(outputs); diff --git a/spring-batch-core/src/test/java/org/springframework/batch/core/step/item/FaultTolerantStepFactoryBeanRetryTests.java b/spring-batch-core/src/test/java/org/springframework/batch/core/step/item/FaultTolerantStepFactoryBeanRetryTests.java index bd374a968..185489245 100644 --- a/spring-batch-core/src/test/java/org/springframework/batch/core/step/item/FaultTolerantStepFactoryBeanRetryTests.java +++ b/spring-batch-core/src/test/java/org/springframework/batch/core/step/item/FaultTolerantStepFactoryBeanRetryTests.java @@ -317,6 +317,7 @@ public class FaultTolerantStepFactoryBeanRetryTests { // [a, b, c, d, e, f, null] assertEquals(7, provided.size()); // [a, b, b, b, b, b, c, d, d, d, d, d, e, f] + System.err.println(processed); assertEquals(14, processed.size()); // [b, d] assertEquals(2, recovered.size()); diff --git a/spring-batch-core/src/test/java/org/springframework/batch/core/step/item/FaultTolerantStepFactoryBeanRollbackTests.java b/spring-batch-core/src/test/java/org/springframework/batch/core/step/item/FaultTolerantStepFactoryBeanRollbackTests.java index 2826cd257..4dcbea7ab 100644 --- a/spring-batch-core/src/test/java/org/springframework/batch/core/step/item/FaultTolerantStepFactoryBeanRollbackTests.java +++ b/spring-batch-core/src/test/java/org/springframework/batch/core/step/item/FaultTolerantStepFactoryBeanRollbackTests.java @@ -313,11 +313,11 @@ public class FaultTolerantStepFactoryBeanRollbackTests { step.execute(stepExecution); assertEquals(BatchStatus.COMPLETED, stepExecution.getStatus()); - // TODO: Fix this with BATCH-1259? - assertEquals("[1, 2, 1, 3, 4, 1, 3, 5]", processor.getProcessed().toString()); assertEquals("[1, 3, 5]", processor.getCommitted().toString()); assertEquals("[1, 3, 5]", writer.getWritten().toString()); assertEquals("[1, 3, 5]", writer.getCommitted().toString()); + // TODO: Fix this with BATCH-1259? + assertEquals("[1, 2, 1, 3, 4, 1, 3, 5]", processor.getProcessed().toString()); } @Test @@ -347,12 +347,12 @@ public class FaultTolerantStepFactoryBeanRollbackTests { step.execute(stepExecution); assertEquals(BatchStatus.COMPLETED, stepExecution.getStatus()); - // TODO: Fix this with BATCH-1256 - assertEquals("[1, 2, 3, 4, 5, 1, 2, 3, 4, 5, 2, 3, 4, 5, 3, 4, 5, 4, 5, 5]", processor.getProcessed() - .toString()); - assertEquals("[1, 2, 3, 4, 5, 2, 3, 4, 5, 3, 4, 5, 5]", processor.getCommitted().toString()); - assertEquals("[1, 2, 3, 4, 1, 2, 3, 4, 5]", writer.getWritten().toString()); + assertEquals("[1, 2, 3, 5]", processor.getCommitted().toString()); assertEquals("[1, 2, 3, 5]", writer.getCommitted().toString()); + assertEquals("[1, 2, 3, 4, 1, 2, 3, 4, 5]", writer.getWritten().toString()); + // TODO: Fix this with BATCH-1259? + assertEquals("[1, 2, 3, 4, 5, 1, 2, 3, 4, 5]", processor.getProcessed() + .toString()); } @Test @@ -365,12 +365,12 @@ public class FaultTolerantStepFactoryBeanRollbackTests { step.execute(stepExecution); assertEquals(BatchStatus.COMPLETED, stepExecution.getStatus()); - // TODO: Fix this with BATCH-1256 - assertEquals("[1, 2, 3, 4, 5, 1, 2, 3, 4, 5, 2, 3, 4, 5, 3, 4, 5, 4, 5, 5]", processor.getProcessed() - .toString()); - assertEquals("[1, 2, 3, 4, 5, 3, 4, 5, 5]", processor.getCommitted().toString()); - assertEquals("[1, 2, 1, 2, 3, 4, 5]", writer.getWritten().toString()); assertEquals("[1, 3, 5]", writer.getCommitted().toString()); + assertEquals("[1, 2, 1, 2, 3, 4, 5]", writer.getWritten().toString()); + assertEquals("[1, 3, 5]", processor.getCommitted().toString()); + // TODO: Fix this with BATCH-1259? + assertEquals("[1, 2, 3, 4, 5, 1, 2, 3, 4, 5]", processor.getProcessed() + .toString()); } @SuppressWarnings("unchecked") diff --git a/spring-batch-core/src/test/java/org/springframework/batch/core/step/item/FaultTolerantStepFactoryBeanTests.java b/spring-batch-core/src/test/java/org/springframework/batch/core/step/item/FaultTolerantStepFactoryBeanTests.java index d07672549..b1ab20c4d 100644 --- a/spring-batch-core/src/test/java/org/springframework/batch/core/step/item/FaultTolerantStepFactoryBeanTests.java +++ b/spring-batch-core/src/test/java/org/springframework/batch/core/step/item/FaultTolerantStepFactoryBeanTests.java @@ -540,10 +540,9 @@ public class FaultTolerantStepFactoryBeanTests { assertEquals(1, stepExecution.getSkipCount()); assertEquals(2, stepExecution.getRollbackCount()); - // 1,2,3,4,3,4,4 - two re-processing attempts until the item is - // identified and finally skipped on the third attempt - assertEquals(7, processor.getProcessed().size()); - assertEquals("[1, 2, 3, 4, 3, 4, 4]", processor.getProcessed().toString()); + // 1,2,3,4,3,4 - two re-processing attempts until the item is + // identified and finally skipped on the second attempt + assertEquals("[1, 2, 3, 4, 3, 4]", processor.getProcessed().toString()); assertStepExecutionsAreEqual(stepExecution, repository.getLastStepExecution(jobExecution.getJobInstance(), step .getName())); }