From e3a9e29b67ea50ddcfb816220da6442f0a3dd2aa Mon Sep 17 00:00:00 2001 From: robokaso Date: Fri, 21 Nov 2008 10:20:03 +0000 Subject: [PATCH] OPEN - BATCH-931: Write failures don't fail immediately. write RecoveryCallback now rethrows non-skippable exception immediately --- ...ractFaultTolerantChunkOrientedTasklet.java | 8 ++- ...aultTolerantChunkOrientedTaskletTests.java | 55 ++++++++++++++++--- 2 files changed, 52 insertions(+), 11 deletions(-) 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 4cf6c1883..f39eef2ed 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 @@ -232,13 +232,15 @@ public abstract class AbstractFaultTolerantChunkOrientedTasklet extends Ab RecoveryCallback recoveryCallback = new RecoveryCallback() { public Object recover(RetryContext context) throws Exception { + Exception le = (Exception) context.getLastThrowable(); + if (!writeSkipPolicy.shouldSkip(le, contribution.getSkipCount())) { + throw le; + } if (chunk.size() == 1) { - Exception e = (Exception) context.getLastThrowable(); O item = chunk.get(0); - checkSkipPolicy(item, e, contribution); + checkSkipPolicy(item, le, contribution); return null; } - Exception le = (Exception) context.getLastThrowable(); if (!rollbackClassifier.classify(le)) { throw new RetryException( "Invalid retry state during write caused by exception that does not classify for rollback: ", 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 dc597dcb4..4f6a2d135 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 @@ -15,13 +15,12 @@ */ package org.springframework.batch.core.step.item; -import static org.junit.Assert.assertEquals; -import static org.junit.Assert.assertTrue; -import static org.junit.Assert.fail; - +import static org.junit.Assert.*; import static org.easymock.EasyMock.*; import java.util.ArrayList; +import java.util.Arrays; +import java.util.HashMap; import java.util.List; import java.util.Map; @@ -33,6 +32,7 @@ import org.springframework.batch.core.SkipListener; import org.springframework.batch.core.StepContribution; import org.springframework.batch.core.StepExecution; import org.springframework.batch.core.scope.ChunkContext; +import org.springframework.batch.core.step.skip.NeverSkipItemSkipPolicy; import org.springframework.batch.core.step.skip.SkipPolicy; import org.springframework.batch.core.step.skip.SkipLimitExceededException; import org.springframework.batch.item.ItemProcessor; @@ -285,7 +285,7 @@ public class FaultTolerantChunkOrientedTaskletTests { assertTrue(attributes.hasAttribute("INPUT_BUFFER_KEY")); } @SuppressWarnings("unchecked") - Map skips = (Map) attributes.getAttribute("SKIPPED_INPUTS_KEY"); + Map skips = (Map) attributes.getAttribute("SKIPPED_INPUTS_KEY"); assertEquals(1, skips.size()); // The last recovery for this chunk... @@ -365,14 +365,14 @@ public class FaultTolerantChunkOrientedTaskletTests { new ItemProcessor() { public String process(Integer item) throws Exception { if (item == 1) { - throw processorException; + throw processorException; } processed.add(item); return String.valueOf(item); } }, new ItemWriter() { public void write(List items) throws Exception { - throw writerException; + throw writerException; } }, chunkOperations, retryTemplate, rollbackClassifier, readSkipPolicy, writeSkipPolicy, writeSkipPolicy); @@ -389,7 +389,7 @@ public class FaultTolerantChunkOrientedTaskletTests { skipListener.onSkipInWrite("2", writerException); expectLastCall().once(); replay(skipListener); - + // processor fails first try { tasklet.execute(contribution, attributes); @@ -424,4 +424,43 @@ public class FaultTolerantChunkOrientedTaskletTests { verify(skipListener); } + @Test + public void testRethrowNonSkippableExceptionOnWriteAsap() throws Exception { + final List chunk = Arrays.asList(new String[] { "1", "2" }); + final Exception ex = new RuntimeException(); + final StepContribution contribution = new StepExecution("foo", null).createStepContribution(); + final Map skipped = new HashMap(); + writeSkipPolicy = new NeverSkipItemSkipPolicy(); + + @SuppressWarnings("unchecked") + ItemWriter itemWriter = createMock(ItemWriter.class); + itemWriter.write(chunk); + expectLastCall().andThrow(ex); + replay(itemWriter); + tasklet = new FaultTolerantChunkOrientedTasklet(itemReader, itemProcessor, itemWriter, + chunkOperations, retryTemplate, rollbackClassifier, readSkipPolicy, writeSkipPolicy, writeSkipPolicy); + + try { + tasklet.write(chunk, contribution, skipped); + fail(); + } + catch (Exception e) { + assertSame(ex, e); + } + + try { + tasklet.write(chunk, contribution, skipped); + fail(); + } + catch (Exception e) { + assertSame(ex, e); + } + + /* + * writer was called only on first failed attempt, exception is rethrown + * immediately when chunk is reprocessed because it is not skippable + */ + verify(itemWriter); + } + }