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 40083a880..bcfb8fa3d 100755 --- 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 @@ -26,8 +26,10 @@ import java.util.Map; import org.springframework.batch.classify.BinaryExceptionClassifier; import org.springframework.batch.classify.Classifier; import org.springframework.batch.classify.SubclassClassifier; +import org.springframework.batch.core.ChunkListener; import org.springframework.batch.core.JobInterruptedException; import org.springframework.batch.core.Step; +import org.springframework.batch.core.StepListener; import org.springframework.batch.core.step.FatalStepExecutionException; import org.springframework.batch.core.step.skip.CompositeSkipPolicy; import org.springframework.batch.core.step.skip.ExceptionClassifierSkipPolicy; @@ -78,6 +80,7 @@ import org.springframework.util.Assert; * * @author Dave Syer * @author Robert Kasanicky + * @author Morten Andersen-Gott * */ public class FaultTolerantStepFactoryBean extends SimpleStepFactoryBean { @@ -120,7 +123,8 @@ public class FaultTolerantStepFactoryBean extends SimpleStepFactoryBean extends SimpleStepFactoryBean extends SimpleStepFactoryBeanretryLimit == 1 by default. * - * @param retryLimit the retry limit to set, must be greater or equal to 1. + * @param retryLimit + * the retry limit to set, must be greater or equal to 1. */ public void setRetryLimit(int retryLimit) { this.retryLimit = retryLimit; @@ -162,8 +168,8 @@ public class FaultTolerantStepFactoryBean extends SimpleStepFactoryBean extends SimpleStepFactoryBean extends SimpleStepFactoryBean, Boolean> retryableExceptionClasses) { this.retryableExceptionClasses = retryableExceptionClasses; @@ -192,7 +200,8 @@ public class FaultTolerantStepFactoryBean extends SimpleStepFactoryBean extends SimpleStepFactoryBean extends SimpleStepFactoryBean extends SimpleStepFactoryBean extends SimpleStepFactoryBean * Defaults to all no exception. * - * @param exceptionClasses defaults to Exception + * @param exceptionClasses + * defaults to Exception */ public void setSkippableExceptionClasses(Map, Boolean> exceptionClasses) { this.skippableExceptionClasses = exceptionClasses; @@ -254,7 +267,8 @@ public class FaultTolerantStepFactoryBean extends SimpleStepFactoryBean * Defaults is empty. * - * @param noRollbackExceptionClasses the exception classes to set + * @param noRollbackExceptionClasses + * the exception classes to set */ public void setNoRollbackExceptionClasses(Collection> noRollbackExceptionClasses) { this.noRollbackExceptionClasses = noRollbackExceptionClasses; @@ -272,7 +286,7 @@ public class FaultTolerantStepFactoryBean extends SimpleStepFactoryBean getRollbackClassifier() { @@ -328,11 +342,11 @@ public class FaultTolerantStepFactoryBean extends SimpleStepFactoryBean extends SimpleStepFactoryBean) { streamIsReader = true; chunkMonitor.registerItemStream(stream); - } - else { + } else { step.registerStream(stream); } } @@ -361,11 +374,9 @@ public class FaultTolerantStepFactoryBean extends SimpleStepFactoryBean extends SimpleStepFactoryBean chunkProvider = new FaultTolerantChunkProvider(getItemReader(), - getChunkOperations()); - chunkProvider.setMaxSkipsOnRead(Math.max(getCommitInterval(), - FaultTolerantChunkProvider.DEFAULT_MAX_SKIPS_ON_READ)); + FaultTolerantChunkProvider chunkProvider = new FaultTolerantChunkProvider(getItemReader(), getChunkOperations()); + chunkProvider.setMaxSkipsOnRead(Math.max(getCommitInterval(), FaultTolerantChunkProvider.DEFAULT_MAX_SKIPS_ON_READ)); chunkProvider.setSkipPolicy(readSkipPolicy); chunkProvider.setRollbackClassifier(getRollbackClassifier()); @@ -394,16 +403,14 @@ public class FaultTolerantStepFactoryBean extends SimpleStepFactoryBean, Boolean> map = new HashMap, Boolean>( - skippableExceptionClasses); + Map, Boolean> map = new HashMap, Boolean>(skippableExceptionClasses); map.put(ForceRollbackForWriteSkipException.class, true); LimitCheckingItemSkipPolicy limitCheckingItemSkipPolicy = new LimitCheckingItemSkipPolicy(skipLimit, map); if (skipPolicy == null) { Assert.state(!(skippableExceptionClasses.isEmpty() && skipLimit > 0), "If a skip limit is provided then skippable exceptions must also be specified"); skipPolicy = limitCheckingItemSkipPolicy; - } - else if (limitCheckingItemSkipPolicy != null) { + } else if (limitCheckingItemSkipPolicy != null) { skipPolicy = new CompositeSkipPolicy(new SkipPolicy[] { skipPolicy, limitCheckingItemSkipPolicy }); } return skipPolicy; @@ -417,8 +424,8 @@ public class FaultTolerantStepFactoryBean extends SimpleStepFactoryBean chunkProcessor = new FaultTolerantChunkProcessor(getItemProcessor(), - getItemWriter(), batchRetryTemplate); + FaultTolerantChunkProcessor chunkProcessor = new FaultTolerantChunkProcessor(getItemProcessor(), getItemWriter(), + batchRetryTemplate); chunkProcessor.setBuffering(!isReaderTransactionalQueue()); chunkProcessor.setProcessorTransactional(processorTransactional); @@ -442,8 +449,7 @@ public class FaultTolerantStepFactoryBean extends SimpleStepFactoryBean, Boolean> map = new HashMap, Boolean>( - retryableExceptionClasses); + Map, Boolean> map = new HashMap, Boolean>(retryableExceptionClasses); map.put(ForceRollbackForWriteSkipException.class, true); simpleRetryPolicy = new SimpleRetryPolicy(retryLimit, map); @@ -451,8 +457,7 @@ public class FaultTolerantStepFactoryBean extends SimpleStepFactoryBean 0), "If a retry limit is provided then retryable exceptions must also be specified"); retryPolicy = simpleRetryPolicy; - } - else if ((!retryableExceptionClasses.isEmpty() && retryLimit > 0)) { + } else if ((!retryableExceptionClasses.isEmpty() && retryLimit > 0)) { CompositeRetryPolicy compositeRetryPolicy = new CompositeRetryPolicy(); compositeRetryPolicy.setPolicies(new RetryPolicy[] { retryPolicy, simpleRetryPolicy }); retryPolicy = compositeRetryPolicy; @@ -469,8 +474,8 @@ public class FaultTolerantStepFactoryBean extends SimpleStepFactoryBean extends SimpleStepFactoryBean 0) { batchRetryTemplate.setRetryContextCache(new MapRetryContextCache(cacheCapacity)); } - } - else { + } else { batchRetryTemplate.setRetryContextCache(retryContextCache); } @@ -501,8 +505,7 @@ public class FaultTolerantStepFactoryBean extends SimpleStepFactoryBean classifier = new SubclassClassifier( - retryPolicy); + SubclassClassifier classifier = new SubclassClassifier(retryPolicy); classifier.setTypeMap(map); ExceptionClassifierRetryPolicy retryPolicyWrapper = new ExceptionClassifierRetryPolicy(); @@ -515,7 +518,8 @@ public class FaultTolerantStepFactoryBean extends SimpleStepFactoryBean extends SimpleStepFactoryBean>) exceptions; } + @Override + protected void registerChunkListeners(TaskletStep step, StepListener listener) { + super.registerChunkListeners(step, new TerminateOnExceptionChunkListenerDelegate((ChunkListener) listener)); + } + + /** + * ChunkListener that wraps exceptions thrown from the ChunkListener in + * {@link FatalStepExecutionException} to force termination of StepExecution + * + * ChunkListeners shoulnd't throw exceptions and expect continued + * processing, they must be handled in the implementation or the step will + * terminate + * + */ + private class TerminateOnExceptionChunkListenerDelegate implements ChunkListener { + + private ChunkListener chunkListener; + + TerminateOnExceptionChunkListenerDelegate(ChunkListener chunkListener) { + this.chunkListener = chunkListener; + } + + public void beforeChunk() { + try { + chunkListener.beforeChunk(); + } catch (Throwable t) { + throw new FatalStepExecutionException("ChunkListener threw exception, rethrowing as fatal", t); + } + } + + public void afterChunk() { + try { + chunkListener.afterChunk(); + } catch (Throwable t) { + throw new FatalStepExecutionException("ChunkListener threw exception, rethrowing as fatal", t); + } + } + + } + } diff --git a/spring-batch-core/src/main/java/org/springframework/batch/core/step/item/SimpleStepFactoryBean.java b/spring-batch-core/src/main/java/org/springframework/batch/core/step/item/SimpleStepFactoryBean.java index 709edd04c..415e027ea 100755 --- a/spring-batch-core/src/main/java/org/springframework/batch/core/step/item/SimpleStepFactoryBean.java +++ b/spring-batch-core/src/main/java/org/springframework/batch/core/step/item/SimpleStepFactoryBean.java @@ -590,15 +590,21 @@ public class SimpleStepFactoryBean implements FactoryBean, BeanNameAware { step.registerStepExecutionListener((StepExecutionListener) listener); } if (listener instanceof ChunkListener) { - step.registerChunkListener((ChunkListener) listener); + registerChunkListeners(step, listener); } } } step.setStepExecutionListeners(BatchListenerFactoryHelper.getListeners(listeners, StepExecutionListener.class) .toArray(new StepExecutionListener[] {})); - step.setChunkListeners(BatchListenerFactoryHelper.getListeners(listeners, ChunkListener.class).toArray( - new ChunkListener[] {})); + + for(ChunkListener chunkListener: BatchListenerFactoryHelper.getListeners(listeners, ChunkListener.class)){ + registerChunkListeners(step,chunkListener); + } + } + + protected void registerChunkListeners(TaskletStep step, StepListener listener) { + step.registerChunkListener((ChunkListener) listener); } /** 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 5ef7a10ee..707114515 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 @@ -1,8 +1,11 @@ package org.springframework.batch.core.step.item; +import static org.hamcrest.CoreMatchers.instanceOf; import static org.junit.Assert.assertEquals; import static org.junit.Assert.assertFalse; +import static org.junit.Assert.assertThat; import static org.junit.Assert.assertTrue; +import static org.springframework.batch.core.BatchStatus.FAILED; import java.util.ArrayList; import java.util.Arrays; @@ -17,12 +20,15 @@ import org.apache.commons.logging.LogFactory; import org.junit.Before; import org.junit.Test; import org.springframework.batch.core.BatchStatus; +import org.springframework.batch.core.ChunkListener; import org.springframework.batch.core.JobExecution; import org.springframework.batch.core.JobParameters; import org.springframework.batch.core.Step; import org.springframework.batch.core.StepExecution; +import org.springframework.batch.core.StepListener; import org.springframework.batch.core.repository.JobRepository; import org.springframework.batch.core.repository.support.MapJobRepositoryFactoryBean; +import org.springframework.batch.core.step.FatalStepExecutionException; import org.springframework.batch.support.transaction.ResourcelessTransactionManager; import org.springframework.core.task.SimpleAsyncTaskExecutor; import org.springframework.transaction.interceptor.RollbackRuleAttribute; @@ -90,6 +96,32 @@ public class FaultTolerantStepFactoryBeanRollbackTests { stepExecution = jobExecution.createStepExecution(factory.getName()); repository.add(stepExecution); } + + @Test + public void testBeforeChunkListenerException() throws Exception{ + factory.setListeners(new StepListener []{new ExceptionThrowingChunkListener(true)}); + Step step = (Step) factory.getObject(); + step.execute(stepExecution); + assertEquals(FAILED, stepExecution.getStatus()); + assertEquals(FAILED.toString(), stepExecution.getExitStatus().getExitCode()); + assertTrue(stepExecution.getCommitCount() == 0);//Make sure exception was thrown in after, not before + Throwable e = stepExecution.getFailureExceptions().get(0); + assertThat(e, instanceOf(FatalStepExecutionException.class)); + assertThat(e.getCause(), instanceOf(IllegalArgumentException.class)); + } + + @Test + public void testAfterChunkListenerException() throws Exception{ + factory.setListeners(new StepListener []{new ExceptionThrowingChunkListener(false)}); + Step step = (Step) factory.getObject(); + step.execute(stepExecution); + assertEquals(FAILED, stepExecution.getStatus()); + assertEquals(FAILED.toString(), stepExecution.getExitStatus().getExitCode()); + assertTrue(stepExecution.getCommitCount() > 0);//Make sure exception was thrown in after, not before + Throwable e = stepExecution.getFailureExceptions().get(0); + assertThat(e, instanceOf(FatalStepExecutionException.class)); + assertThat(e.getCause(), instanceOf(IllegalArgumentException.class)); + } @Test public void testOverrideWithoutChangingRollbackRules() throws Exception { @@ -525,5 +557,25 @@ public class FaultTolerantStepFactoryBeanRollbackTests { } return map; } + + class ExceptionThrowingChunkListener implements ChunkListener{ + + private boolean throwBefore = true; + + public ExceptionThrowingChunkListener(boolean throwBefore) { + this.throwBefore = throwBefore; + } + + public void beforeChunk() { + if(throwBefore){ + throw new IllegalArgumentException("Planned exception"); + } + } + + public void afterChunk() { + throw new IllegalArgumentException("Planned exception"); + + } + } }