From f55bc5e44f84efab08ba37d9580400ab033b3450 Mon Sep 17 00:00:00 2001 From: Dave Syer Date: Mon, 22 Nov 2010 11:03:39 +0000 Subject: [PATCH] BATCH-1656: fix logic error for skip limit exceeded on no-rollback exception --- ...epFactoryBeanRollbackIntegrationTests.java | 30 +++++++------ .../item/FaultTolerantChunkProcessor.java | 18 ++++++-- .../item/FaultTolerantStepFactoryBean.java | 21 +++++---- ...tTolerantStepFactoryBeanRollbackTests.java | 43 +++++++++++++++++++ 4 files changed, 83 insertions(+), 29 deletions(-) diff --git a/spring-batch-core-tests/src/test/java/org/springframework/batch/core/test/step/FaultTolerantStepFactoryBeanRollbackIntegrationTests.java b/spring-batch-core-tests/src/test/java/org/springframework/batch/core/test/step/FaultTolerantStepFactoryBeanRollbackIntegrationTests.java index 93de13148..92a5dafe1 100644 --- a/spring-batch-core-tests/src/test/java/org/springframework/batch/core/test/step/FaultTolerantStepFactoryBeanRollbackIntegrationTests.java +++ b/spring-batch-core-tests/src/test/java/org/springframework/batch/core/test/step/FaultTolerantStepFactoryBeanRollbackIntegrationTests.java @@ -74,7 +74,6 @@ public class FaultTolerantStepFactoryBeanRollbackIntegrationTests { @Autowired private PlatformTransactionManager transactionManager; - @SuppressWarnings("unchecked") @Before public void setUp() throws Exception { @@ -89,19 +88,11 @@ public class FaultTolerantStepFactoryBeanRollbackIntegrationTests { factory.setJobRepository(repository); factory.setCommitInterval(3); factory.setSkipLimit(10); - ThreadPoolTaskExecutor taskExecutor = new ThreadPoolTaskExecutor(); - taskExecutor.setCorePoolSize(3); - taskExecutor.setMaxPoolSize(6); - taskExecutor.setQueueCapacity(0); - taskExecutor.afterPropertiesSet(); - factory.setTaskExecutor(taskExecutor); - - factory.setSkippableExceptionClasses(getExceptionMap(Exception.class)); SimpleJdbcTestUtils.deleteFromTables(new SimpleJdbcTemplate(dataSource), "ERROR_LOG"); } - + @Test public void testUpdatesNoRollback() throws Exception { @@ -120,12 +111,23 @@ public class FaultTolerantStepFactoryBeanRollbackIntegrationTests { @Test public void testMultithreadedSkipInWriter() throws Throwable { + ThreadPoolTaskExecutor taskExecutor = new ThreadPoolTaskExecutor(); + taskExecutor.setCorePoolSize(3); + taskExecutor.setMaxPoolSize(6); + taskExecutor.setQueueCapacity(0); + taskExecutor.afterPropertiesSet(); + factory.setTaskExecutor(taskExecutor); + + @SuppressWarnings("unchecked") + Map, Boolean> skippable = getExceptionMap(Exception.class); + factory.setSkippableExceptionClasses(skippable); + jobExecution = repository.createJobExecution("skipJob", new JobParameters()); for (int i = 0; i < MAX_COUNT; i++) { - if (i%100==0) { - logger.info("Starting step: "+i); + if (i % 100 == 0) { + logger.info("Starting step: " + i); } SimpleJdbcTemplate jdbcTemplate = new SimpleJdbcTemplate(dataSource); @@ -263,7 +265,7 @@ public class FaultTolerantStepFactoryBeanRollbackIntegrationTests { public SkipProcessorStub(DataSource dataSource) { jdbcTemplate = new SimpleJdbcTemplate(dataSource); } - + /** * @return the processed */ @@ -287,7 +289,7 @@ public class FaultTolerantStepFactoryBeanRollbackIntegrationTests { public String process(String item) throws Exception { processed.add(item); - logger.debug("Processed item: "+item); + logger.debug("Processed item: " + item); jdbcTemplate.update("INSERT INTO ERROR_LOG (MESSAGE, STEP_NAME) VALUES (?, ?)", item, "processed"); return item; } 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 74ac3aad1..ee39e9977 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 @@ -29,6 +29,7 @@ import org.springframework.batch.classify.Classifier; import org.springframework.batch.core.StepContribution; import org.springframework.batch.core.step.skip.LimitCheckingItemSkipPolicy; import org.springframework.batch.core.step.skip.NonSkippableProcessException; +import org.springframework.batch.core.step.skip.SkipLimitExceededException; import org.springframework.batch.core.step.skip.SkipListenerFailedException; import org.springframework.batch.core.step.skip.SkipPolicy; import org.springframework.batch.item.ItemProcessor; @@ -269,13 +270,19 @@ public class FaultTolerantChunkProcessor extends SimpleChunkProcessor extends SimpleChunkProcessor extends SimpleChunkProcessor extends SimpleStepFactoryBean extends SimpleStepFactoryBean>(); - for (Class exceptionClass : nonSkippableExceptionClasses) { + private void addNonSkippableExceptionIfMissing(Class... cls) { + List> exceptions = new ArrayList>(); + for (Class exceptionClass : nonSkippableExceptionClasses) { exceptions.add(exceptionClass); } - for (Class fatal : cls) { + for (Class fatal : cls) { if (!exceptions.contains(fatal)) { exceptions.add(fatal); } @@ -534,18 +534,17 @@ public class FaultTolerantStepFactoryBean extends SimpleStepFactoryBean>(); - for (Class exceptionClass : nonRetryableExceptionClasses) { + private void addNonRetryableExceptionIfMissing(Class... cls) { + List> exceptions = new ArrayList>(); + for (Class exceptionClass : nonRetryableExceptionClasses) { exceptions.add(exceptionClass); } - for (Class fatal : cls) { + for (Class fatal : cls) { if (!exceptions.contains(fatal)) { exceptions.add(fatal); } } - nonRetryableExceptionClasses = exceptions; + nonRetryableExceptionClasses = (List>)exceptions; } } 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 511200a10..dadf353ee 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 @@ -4,8 +4,10 @@ import static org.junit.Assert.assertEquals; import static org.junit.Assert.assertFalse; import static org.junit.Assert.assertTrue; +import java.util.ArrayList; import java.util.Arrays; import java.util.Collection; +import java.util.Collections; import java.util.HashMap; import java.util.List; import java.util.Map; @@ -196,6 +198,47 @@ public class FaultTolerantStepFactoryBeanRollbackTests { assertEquals(2, stepExecution.getRollbackCount()); } + @Test + public void testNoRollbackInProcessorWhenSkipExceeded() throws Throwable { + + jobExecution = repository.createJobExecution("noRollbackJob", new JobParameters()); + + factory.setSkipLimit(0); + + reader.clear(); + reader.setItems("1", "2", "3", "4", "5"); + factory.setItemReader(reader); + writer.clear(); + factory.setItemWriter(writer); + processor.clear(); + factory.setItemProcessor(processor); + + @SuppressWarnings("unchecked") + List> exceptions = Arrays.>asList(Exception.class); + factory.setNoRollbackExceptionClasses(exceptions); + @SuppressWarnings("unchecked") + Map, Boolean> skippable = getExceptionMap(Exception.class); + factory.setSkippableExceptionClasses(skippable); + + processor.setFailures("2"); + + Step step = (Step) factory.getObject(); + + stepExecution = jobExecution.createStepExecution(factory.getName()); + repository.add(stepExecution); + step.execute(stepExecution); + assertEquals(BatchStatus.COMPLETED, stepExecution.getStatus()); + + assertEquals("[1, 3, 4, 5]", writer.getCommitted().toString()); + // No rollback on 2 so processor has side effect + assertEquals("[1, 2, 3, 4, 5]", processor.getCommitted().toString()); + List processed = new ArrayList(processor.getProcessed()); + Collections.sort(processed); + assertEquals("[1, 2, 3, 4, 5]", processed.toString()); + assertEquals(0, stepExecution.getSkipCount()); + + } + @Test public void testProcessSkipWithNoRollbackForCheckedException() throws Exception { processor.setFailures("4");