diff --git a/spring-batch-core/src/main/java/org/springframework/batch/core/StepContribution.java b/spring-batch-core/src/main/java/org/springframework/batch/core/StepContribution.java index fdc38f9c9..d3b180c7c 100644 --- a/spring-batch-core/src/main/java/org/springframework/batch/core/StepContribution.java +++ b/spring-batch-core/src/main/java/org/springframework/batch/core/StepContribution.java @@ -25,6 +25,10 @@ package org.springframework.batch.core; public class StepContribution { private volatile int itemCount = 0; + + private volatile int readCount = 0; + + private volatile int writeCount = 0; private volatile int filterCount = 0; @@ -61,7 +65,21 @@ public class StepContribution { * Increment the counter for the number of items processed. */ public void incrementItemCount(int count) { - itemCount+=count; + itemCount += count; + } + + /** + * Increment the counter for the number of items read. + */ + public void incrementReadCount() { + readCount++; + } + + /** + * Increment the counter for the number of items written. + */ + public void incrementWriteCount() { + writeCount++; } /** @@ -72,6 +90,24 @@ public class StepContribution { public int getItemCount() { return itemCount; } + + /** + * Public access to the read counter. + * + * @return the item counter. + */ + public int getReadCount() { + return readCount; + } + + /** + * Public access to the write counter. + * + * @return the item counter. + */ + public int getWriteCount() { + return writeCount; + } /** * Public getter for the filter counter. diff --git a/spring-batch-core/src/main/java/org/springframework/batch/core/StepExecution.java b/spring-batch-core/src/main/java/org/springframework/batch/core/StepExecution.java index 460429b01..5b6a818e9 100644 --- a/spring-batch-core/src/main/java/org/springframework/batch/core/StepExecution.java +++ b/spring-batch-core/src/main/java/org/springframework/batch/core/StepExecution.java @@ -42,6 +42,10 @@ public class StepExecution extends Entity { private volatile BatchStatus status = BatchStatus.STARTING; private volatile int itemCount = 0; + + private volatile int readCount = 0; + + private volatile int writeCount = 0; private volatile int commitCount = 0; @@ -63,7 +67,7 @@ public class StepExecution extends Entity { private volatile boolean terminateOnly; - private int filterCount; + private volatile int filterCount; /** * Constructor with mandatory properties. @@ -164,6 +168,42 @@ public class StepExecution extends Entity { public void setItemCount(int itemCount) { this.itemCount = itemCount; } + + /** + * Returns the current number of items read for this execution + * + * @return the current number of items read for this execution + */ + public int getReadCount() { + return readCount; + } + + /** + * Sets the current number of read items for this execution + * + * @param readCount the current number of read items for this execution + */ + public void setReadCount(int readCount) { + this.readCount = readCount; + } + + /** + * Returns the current number of items written for this execution + * + * @return the current number of items written for this execution + */ + public int getWriteCount() { + return writeCount; + } + + /** + * Sets the current number of written items for this execution + * + * @param writeCount the current number of written items for this execution + */ + public void setWriteCount(int writeCount) { + this.writeCount = writeCount; + } /** * Returns the current number of rollbacks for this execution @@ -298,6 +338,8 @@ public class StepExecution extends Entity { readSkipCount += contribution.getReadSkipCount(); writeSkipCount += contribution.getWriteSkipCount(); filterCount += contribution.getFilterCount(); + readCount += contribution.getReadCount(); + writeCount += contribution.getWriteCount(); } /** diff --git a/spring-batch-core/src/main/java/org/springframework/batch/core/step/item/ChunkOrientedTasklet.java b/spring-batch-core/src/main/java/org/springframework/batch/core/step/item/ChunkOrientedTasklet.java index 7a26e1ce1..ed38b443d 100644 --- a/spring-batch-core/src/main/java/org/springframework/batch/core/step/item/ChunkOrientedTasklet.java +++ b/spring-batch-core/src/main/java/org/springframework/batch/core/step/item/ChunkOrientedTasklet.java @@ -205,6 +205,7 @@ public class ChunkOrientedTasklet implements Tasklet { // TODO: segregate read / write / filter count // (this is read count) contribution.incrementItemCount(); + contribution.incrementReadCount(); S output = doProcess(item); if (output != null) { outputs.add(output); @@ -241,7 +242,7 @@ public class ChunkOrientedTasklet implements Tasklet { * @param contribution current context */ protected void write(Chunk chunk, StepContribution contribution) throws Exception { - doWrite(chunk.getItems()); + doWrite(chunk.getItems(), contribution); chunk.clear(); } @@ -249,11 +250,11 @@ public class ChunkOrientedTasklet implements Tasklet { * @param items * @throws Exception */ - protected final void doWrite(List items) throws Exception { + protected final void doWrite(List items, StepContribution contribution) throws Exception { try { listener.beforeWrite(items); itemWriter.write(items); - // TODO: increment write count + contribution.incrementWriteCount(); listener.afterWrite(items); } catch (Exception e) { 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 f66cd35a3..0b5db8a49 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 @@ -332,6 +332,7 @@ public class SkipLimitStepFactoryBean extends SimpleStepFactoryBean * count * @return next item for processing */ + @Override protected ItemWrapper read(StepContribution contribution) throws Exception { int skipCount = 0; @@ -453,7 +454,7 @@ public class SkipLimitStepFactoryBean extends SimpleStepFactoryBean RetryCallback retryCallback = new RetryCallback() { public Object doWithRetry(RetryContext context) throws Exception { - doWrite(chunk.getItems()); + doWrite(chunk.getItems(), contribution); return null; } }; @@ -473,7 +474,7 @@ public class SkipLimitStepFactoryBean extends SimpleStepFactoryBean for (Chunk.ChunkIterator iterator = chunk.iterator(); iterator.hasNext();) { S item = iterator.next(); try { - doWrite(Collections.singletonList(item)); + doWrite(Collections.singletonList(item), contribution); } catch (Exception e) { checkSkipPolicy(contribution, iterator, e); diff --git a/spring-batch-core/src/test/java/org/springframework/batch/core/StepExecutionTests.java b/spring-batch-core/src/test/java/org/springframework/batch/core/StepExecutionTests.java index c667e3567..ffa07b5c2 100644 --- a/spring-batch-core/src/test/java/org/springframework/batch/core/StepExecutionTests.java +++ b/spring-batch-core/src/test/java/org/springframework/batch/core/StepExecutionTests.java @@ -133,10 +133,16 @@ public class StepExecutionTests extends TestCase { contribution.incrementReadSkipCount(); contribution.incrementWriteSkipCount(); contribution.incrementItemCount(); + contribution.incrementReadCount(); + contribution.incrementWriteCount(); + contribution.incrementFilterCount(1); execution.apply(contribution); assertEquals(1, execution.getReadSkipCount()); assertEquals(1, execution.getWriteSkipCount()); assertEquals(1, execution.getItemCount()); + assertEquals(1, execution.getReadCount()); + assertEquals(1, execution.getWriteCount()); + assertEquals(1, execution.getFilterCount()); } public void testTerminateOnly() throws Exception { diff --git a/spring-batch-core/src/test/java/org/springframework/batch/core/step/item/ChunkOrientedTaskletTests.java b/spring-batch-core/src/test/java/org/springframework/batch/core/step/item/ChunkOrientedTaskletTests.java index a5736e083..5684b89ea 100644 --- a/spring-batch-core/src/test/java/org/springframework/batch/core/step/item/ChunkOrientedTaskletTests.java +++ b/spring-batch-core/src/test/java/org/springframework/batch/core/step/item/ChunkOrientedTaskletTests.java @@ -92,21 +92,24 @@ public class ChunkOrientedTaskletTests { @Test public void testHandleCompositeItem() throws Exception { ChunkOrientedTasklet handler = new ChunkOrientedTasklet(itemReader, - new AgrgegateItemProcessor(), itemWriter, repeatTemplate); + new AggregateItemProcessor(), itemWriter, repeatTemplate); StepContribution contribution = new StepContribution(new StepExecution("foo", new JobExecution(new JobInstance( 123L, new JobParameters(), "job")))); handler.execute(contribution, context); assertEquals(2, itemReader.count); assertEquals(2, contribution.getItemCount()); + assertEquals(2, contribution.getReadCount()); assertEquals(1, contribution.getFilterCount()); + assertEquals(1, contribution.getWriteCount()); assertEquals("12", itemWriter.values); } + /** * @author Dave Syer * */ - private final class AgrgegateItemProcessor implements ItemProcessor { + private final class AggregateItemProcessor implements ItemProcessor { private int count = 0; private String value = "";