diff --git a/spring-batch-core/src/main/java/org/springframework/batch/core/step/item/SkipLimitStepFactoryBean.java b/spring-batch-core/src/main/java/org/springframework/batch/core/step/item/SkipLimitStepFactoryBean.java index f8135cdf1..be9eb726d 100644 --- a/spring-batch-core/src/main/java/org/springframework/batch/core/step/item/SkipLimitStepFactoryBean.java +++ b/spring-batch-core/src/main/java/org/springframework/batch/core/step/item/SkipLimitStepFactoryBean.java @@ -33,6 +33,7 @@ import org.springframework.batch.retry.policy.MapRetryContextCache; import org.springframework.batch.retry.policy.RetryContextCache; import org.springframework.batch.retry.policy.SimpleRetryPolicy; import org.springframework.batch.retry.support.RetryTemplate; +import org.springframework.batch.support.Classifier; /** * Factory bean for step that provides options for configuring skip behaviour. @@ -225,6 +226,12 @@ public class SkipLimitStepFactoryBean extends SimpleStepFactoryBean retryTemplate.setBackOffPolicy(backOffPolicy); } retryTemplate.setRetryPolicy(retryPolicy); + Classifier rollbackClassifier = new Classifier() { + public Boolean classify(Throwable classifiable) { + return getTransactionAttribute().rollbackOn(classifiable); + } + }; + retryTemplate.setRollbackClassifier(rollbackClassifier); // Co-ordinate the retry policy with the exception handler: RepeatOperations stepOperations = getStepOperations(); @@ -254,7 +261,7 @@ public class SkipLimitStepFactoryBean extends SimpleStepFactoryBean ItemSkipPolicy writeSkipPolicy = new LimitCheckingItemSkipPolicy(skipLimit, exceptions, new ArrayList>(fatalExceptionClasses)); ChunkOrientedTasklet tasklet = new StatefulRetryTasklet(getItemReader(), getItemProcessor(), - getItemWriter(), getChunkOperations(), retryTemplate, readSkipPolicy, writeSkipPolicy, writeSkipPolicy); + getItemWriter(), getChunkOperations(), retryTemplate, rollbackClassifier, readSkipPolicy, writeSkipPolicy, writeSkipPolicy); tasklet.setListeners(getListeners()); step.setTasklet(tasklet); @@ -297,6 +304,8 @@ public class SkipLimitStepFactoryBean extends SimpleStepFactoryBean final private ItemSkipPolicy processSkipPolicy; + final private Classifier rollbackClassifier; + /** * @param itemReader * @param itemWriter @@ -304,10 +313,11 @@ public class SkipLimitStepFactoryBean extends SimpleStepFactoryBean */ public StatefulRetryTasklet(ItemReader itemReader, ItemProcessor itemProcessor, ItemWriter itemWriter, - RepeatOperations chunkOperations, RetryOperations retryTemplate, ItemSkipPolicy readSkipPolicy, + RepeatOperations chunkOperations, RetryOperations retryTemplate, Classifier rollbackClassifier, ItemSkipPolicy readSkipPolicy, ItemSkipPolicy writeSkipPolicy, ItemSkipPolicy processSkipPolicy) { super(itemReader, itemProcessor, itemWriter, chunkOperations); this.retryOperations = retryTemplate; + this.rollbackClassifier = rollbackClassifier; this.readSkipPolicy = readSkipPolicy; this.writeSkipPolicy = writeSkipPolicy; this.processSkipPolicy = processSkipPolicy; @@ -472,7 +482,9 @@ public class SkipLimitStepFactoryBean extends SimpleStepFactoryBean } catch (Exception e) { checkSkipPolicy(contribution, iterator, e); - throw e; + if (rollbackClassifier.classify(e)) { + throw e; + } } } diff --git a/spring-batch-core/src/test/java/org/springframework/batch/core/repository/dao/AbstractExecutionContextDaoTests.java b/spring-batch-core/src/test/java/org/springframework/batch/core/repository/dao/AbstractExecutionContextDaoTests.java index a287766d2..7babacaf2 100644 --- a/spring-batch-core/src/test/java/org/springframework/batch/core/repository/dao/AbstractExecutionContextDaoTests.java +++ b/spring-batch-core/src/test/java/org/springframework/batch/core/repository/dao/AbstractExecutionContextDaoTests.java @@ -3,8 +3,6 @@ package org.springframework.batch.core.repository.dao; import static org.junit.Assert.assertEquals; import java.util.Collections; -import java.util.HashMap; - import org.junit.Before; import org.junit.Test; import org.springframework.batch.core.JobExecution; diff --git a/spring-batch-core/src/test/java/org/springframework/batch/core/step/item/StatefulRetryTaskletTests.java b/spring-batch-core/src/test/java/org/springframework/batch/core/step/item/StatefulRetryTaskletTests.java index ac58da53f..ed54a7775 100644 --- a/spring-batch-core/src/test/java/org/springframework/batch/core/step/item/StatefulRetryTaskletTests.java +++ b/spring-batch-core/src/test/java/org/springframework/batch/core/step/item/StatefulRetryTaskletTests.java @@ -42,6 +42,7 @@ import org.springframework.batch.repeat.policy.SimpleCompletionPolicy; import org.springframework.batch.repeat.support.RepeatTemplate; import org.springframework.batch.retry.policy.NeverRetryPolicy; import org.springframework.batch.retry.support.RetryTemplate; +import org.springframework.batch.support.Classifier; /** * @author Dave Syer @@ -84,6 +85,12 @@ public class StatefulRetryTaskletTests { }; private RetryTemplate retryTemplate = new RetryTemplate(); + + private Classifier rollbackClassifier = new Classifier() { + public Boolean classify(Throwable classifiable) { + return true; + } + }; private ItemSkipPolicy readSkipPolicy = new ItemSkipPolicy() { public boolean shouldSkip(Throwable t, int skipCount) throws SkipLimitExceededException { @@ -105,7 +112,7 @@ public class StatefulRetryTaskletTests { @Test public void testBasicHandle() throws Exception { handler = new StatefulRetryTasklet(itemReader, itemProcessor, itemWriter, chunkOperations, - retryTemplate, readSkipPolicy, writeSkipPolicy, writeSkipPolicy); + retryTemplate, rollbackClassifier, readSkipPolicy, writeSkipPolicy, writeSkipPolicy); StepContribution contribution = new StepExecution("foo", null).createStepContribution(); handler.execute(contribution, new BasicAttributeAccessor()); assertEquals(limit, contribution.getItemCount()); @@ -117,7 +124,7 @@ public class StatefulRetryTaskletTests { public Integer read() throws Exception, UnexpectedInputException, NoWorkFoundException, ParseException { throw new RuntimeException("Barf!"); } - }, itemProcessor, itemWriter, chunkOperations, retryTemplate, readSkipPolicy, writeSkipPolicy, writeSkipPolicy); + }, itemProcessor, itemWriter, chunkOperations, retryTemplate, rollbackClassifier, readSkipPolicy, writeSkipPolicy, writeSkipPolicy); chunkOperations.setCompletionPolicy(new SimpleCompletionPolicy(1)); StepContribution contribution = new StepExecution("foo", null).createStepContribution(); BasicAttributeAccessor attributes = new BasicAttributeAccessor(); @@ -139,7 +146,7 @@ public class StatefulRetryTaskletTests { written.addAll(items); throw new RuntimeException("Barf!"); } - }, chunkOperations, retryTemplate, readSkipPolicy, writeSkipPolicy, writeSkipPolicy); + }, chunkOperations, retryTemplate, rollbackClassifier, readSkipPolicy, writeSkipPolicy, writeSkipPolicy); chunkOperations.setCompletionPolicy(new SimpleCompletionPolicy(1)); StepContribution contribution = new StepExecution("foo", null).createStepContribution(); BasicAttributeAccessor attributes = new BasicAttributeAccessor(); @@ -165,7 +172,7 @@ public class StatefulRetryTaskletTests { written.addAll(items); throw new RuntimeException("Barf!"); } - }, chunkOperations, retryTemplate, readSkipPolicy, writeSkipPolicy, writeSkipPolicy); + }, chunkOperations, retryTemplate, rollbackClassifier, readSkipPolicy, writeSkipPolicy, writeSkipPolicy); chunkOperations.setCompletionPolicy(new SimpleCompletionPolicy(2)); StepContribution contribution = new StepExecution("foo", null).createStepContribution(); BasicAttributeAccessor attributes = new BasicAttributeAccessor(); @@ -217,7 +224,7 @@ public class StatefulRetryTaskletTests { throw new RuntimeException("Barf!"); } } - , itemWriter, chunkOperations, retryTemplate, readSkipPolicy, writeSkipPolicy, writeSkipPolicy); + , itemWriter, chunkOperations, retryTemplate, rollbackClassifier, readSkipPolicy, writeSkipPolicy, writeSkipPolicy); chunkOperations.setCompletionPolicy(new SimpleCompletionPolicy(2)); StepContribution contribution = new StepExecution("foo", null).createStepContribution(); BasicAttributeAccessor attributes = new BasicAttributeAccessor(); diff --git a/spring-batch-infrastructure/src/main/java/org/springframework/batch/retry/support/RetryTemplate.java b/spring-batch-infrastructure/src/main/java/org/springframework/batch/retry/support/RetryTemplate.java index 33600b8c9..b3c6f0c2a 100644 --- a/spring-batch-infrastructure/src/main/java/org/springframework/batch/retry/support/RetryTemplate.java +++ b/spring-batch-infrastructure/src/main/java/org/springframework/batch/retry/support/RetryTemplate.java @@ -448,6 +448,8 @@ public class RetryTemplate implements RetryOperations { protected boolean shouldRethrow(RetryPolicy retryPolicy, RetryContext context, RetryState state) { // Allow stateless behaviour to take over for certain exception types if (rollbackClassifier != null) { + // TODO: remove this. Make it part of the stateful execution parameters? + // Then we wouldn't have to make assertions about the stateless case. boolean rollback = rollbackClassifier.classify(context.getLastThrowable()); if (rollback && state == null && retryPolicy.canRetry(context)) { throw new RetryException("Inconsistent configuration. The retry policy says we can retry but " diff --git a/spring-batch-samples/src/main/java/org/springframework/batch/sample/support/ItemTrackingItemWriter.java b/spring-batch-samples/src/main/java/org/springframework/batch/sample/support/ItemTrackingItemWriter.java index 7d34c6daf..aa468f8d7 100644 --- a/spring-batch-samples/src/main/java/org/springframework/batch/sample/support/ItemTrackingItemWriter.java +++ b/spring-batch-samples/src/main/java/org/springframework/batch/sample/support/ItemTrackingItemWriter.java @@ -28,6 +28,7 @@ public class ItemTrackingItemWriter implements ItemWriter { counter += items.size(); if (current < failure && counter >= failure) { failed = items.get(failure-current-1); + this.items.remove(failed); throw new ValidationException("validation failed"); } } diff --git a/spring-batch-samples/src/test/java/org/springframework/batch/sample/support/ItemTrackingItemWriterTests.java b/spring-batch-samples/src/test/java/org/springframework/batch/sample/support/ItemTrackingItemWriterTests.java index 14968ff10..2e66cc197 100644 --- a/spring-batch-samples/src/test/java/org/springframework/batch/sample/support/ItemTrackingItemWriterTests.java +++ b/spring-batch-samples/src/test/java/org/springframework/batch/sample/support/ItemTrackingItemWriterTests.java @@ -51,9 +51,10 @@ public class ItemTrackingItemWriterTests { catch (ValidationException e) { // expected } - assertEquals(3, writer.getItems().size()); + // the failed item is removed + assertEquals(2, writer.getItems().size()); writer.write(Arrays.asList("a", "e", "c")); - assertEquals(6, writer.getItems().size()); + assertEquals(5, writer.getItems().size()); try { writer.write(Arrays.asList("f", "b", "g")); fail("Expected ValidationException"); @@ -61,7 +62,8 @@ public class ItemTrackingItemWriterTests { catch (ValidationException e) { // expected } - assertEquals(6, writer.getItems().size()); + // barf immediately if a failure is detected + assertEquals(5, writer.getItems().size()); } }