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 ee39e9977..0abf493c6 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 @@ -529,6 +529,7 @@ public class FaultTolerantChunkProcessor extends SimpleChunkProcessor implements ChunkProcessor, Initializi doAfterWrite(items); } catch (Exception e) { - listener.onWriteError(e, items); + doOnWriteError(e, items); throw e; } @@ -165,6 +165,9 @@ public class SimpleChunkProcessor implements ChunkProcessor, Initializi protected final void doAfterWrite(List items) { listener.afterWrite(items); } + protected final void doOnWriteError(Exception e, List items) { + listener.onWriteError(e, items); + } protected void writeItems(List items) throws Exception { if (itemWriter != null) { diff --git a/spring-batch-core/src/test/java/org/springframework/batch/core/step/item/FaultTolerantChunkProcessorTests.java b/spring-batch-core/src/test/java/org/springframework/batch/core/step/item/FaultTolerantChunkProcessorTests.java index d0e2de168..e05010d4d 100644 --- a/spring-batch-core/src/test/java/org/springframework/batch/core/step/item/FaultTolerantChunkProcessorTests.java +++ b/spring-batch-core/src/test/java/org/springframework/batch/core/step/item/FaultTolerantChunkProcessorTests.java @@ -30,6 +30,8 @@ public class FaultTolerantChunkProcessorTests { private List list = new ArrayList(); private List after = new ArrayList(); + + private List writeError = new ArrayList(); private FaultTolerantChunkProcessor processor; @@ -213,22 +215,10 @@ public class FaultTolerantChunkProcessorTests { } })); processor.setWriteSkipPolicy(new AlwaysSkipItemSkipPolicy()); - try { - processor.process(contribution, chunk); - fail(); - } - catch (RuntimeException e) { - assertEquals("Planned failure!", e.getMessage()); - } + processAndExpectPlannedRuntimeException(chunk); processor.process(contribution, chunk); assertEquals(2, chunk.getItems().size()); - try { - processor.process(contribution, chunk); - fail(); - } - catch (RuntimeException e) { - assertEquals("Planned failure!", e.getMessage()); - } + processAndExpectPlannedRuntimeException(chunk); assertEquals(1, chunk.getItems().size()); processor.process(contribution, chunk); assertEquals(0, chunk.getItems().size()); @@ -260,6 +250,59 @@ public class FaultTolerantChunkProcessorTests { })); processor.setWriteSkipPolicy(new AlwaysSkipItemSkipPolicy()); + processAndExpectPlannedRuntimeException(chunk); + processor.process(contribution, chunk); + processor.process(contribution, chunk); + + assertEquals("[foo, bar]", list.toString()); + assertEquals("[foo, bar]", after.toString()); + } + + @Test + public void testOnErrorInWrite() throws Exception{ + Chunk chunk = new Chunk(Arrays.asList("foo", "fail")); + processor.setListeners(Arrays.asList(new ItemListenerSupport() { + @Override + public void onWriteError(Exception e, List item) { + writeError.addAll(item); + } + })); + processor.setWriteSkipPolicy(new AlwaysSkipItemSkipPolicy()); + + processAndExpectPlannedRuntimeException(chunk);//Process foo, fail + processor.process(contribution, chunk);;//Process foo + processAndExpectPlannedRuntimeException(chunk);//Process fail + + assertEquals("[foo, fail, fail]", writeError.toString()); + } + + @Test + public void testOnErrorInWriteAllItemsFail() throws Exception{ + Chunk chunk = new Chunk(Arrays.asList("foo", "bar")); + processor = new FaultTolerantChunkProcessor(new PassThroughItemProcessor(), + new ItemWriter() { + public void write(List items) throws Exception { + //Always fail in writer + throw new RuntimeException("Planned failure!"); + } + }, batchRetryTemplate); + processor.setListeners(Arrays.asList(new ItemListenerSupport() { + @Override + public void onWriteError(Exception e, List item) { + writeError.addAll(item); + } + })); + processor.setWriteSkipPolicy(new AlwaysSkipItemSkipPolicy()); + + processAndExpectPlannedRuntimeException(chunk);//Process foo, bar + processAndExpectPlannedRuntimeException(chunk);//Process foo + processAndExpectPlannedRuntimeException(chunk);//Process bar + + assertEquals("[foo, bar, foo, bar]", writeError.toString()); + } + + protected void processAndExpectPlannedRuntimeException(Chunk chunk) + throws Exception { try { processor.process(contribution, chunk); fail(); @@ -267,10 +310,5 @@ public class FaultTolerantChunkProcessorTests { catch (RuntimeException e) { assertEquals("Planned failure!", e.getMessage()); } - processor.process(contribution, chunk); - processor.process(contribution, chunk); - - assertEquals("[foo, bar]", list.toString()); - assertEquals("[foo, bar]", after.toString()); } }