diff --git a/spring-batch-infrastructure/src/main/java/org/springframework/batch/item/MultiResourceItemReader.java b/spring-batch-infrastructure/src/main/java/org/springframework/batch/item/MultiResourceItemReader.java index 299fc040b..347cbe712 100644 --- a/spring-batch-infrastructure/src/main/java/org/springframework/batch/item/MultiResourceItemReader.java +++ b/spring-batch-infrastructure/src/main/java/org/springframework/batch/item/MultiResourceItemReader.java @@ -24,16 +24,6 @@ import org.springframework.util.ClassUtils; public class MultiResourceItemReader extends ExecutionContextUserSupport implements ItemReader, ItemStream, InitializingBean { - /** - * Key for the index of the current resource - */ - private static final String RESOURCE_INDEX = "resourceIndex"; - - /** - * Key for item count within current resource - */ - private static final String ITEM_COUNT = "itemIndex"; - /** * Unique object instance that marks resource boundaries in the item buffer */ @@ -43,13 +33,7 @@ public class MultiResourceItemReader extends ExecutionContextUserSupport impleme private Resource[] resources; - private int currentResourceIndex = 0; - - private int lastMarkedResourceIndex = 0; - - private long currentItemIndex = 0; - - private long lastMarkedItemIndex = 0; + private MultiResourceIndex index = new MultiResourceIndex(); private List itemBuffer = new ArrayList(); @@ -70,20 +54,14 @@ public class MultiResourceItemReader extends ExecutionContextUserSupport impleme Object item; if (shouldReadBuffer) { - if (itemBufferIterator.hasNext()) { - item = readBufferedItem(); - } - else { - // buffer is exhausted, continue reading from file - shouldReadBuffer = false; - itemBufferIterator = null; - item = readNextItem(); - } + item = readBufferedItem(); } else { item = readNextItem(); } + index.incrementItemCount(); + return item; } @@ -98,15 +76,15 @@ public class MultiResourceItemReader extends ExecutionContextUserSupport impleme while (item == null) { - incrementResourceIndex(); + index.incrementResourceCount(); - if (currentResourceIndex >= resources.length) { + if (index.currentResource >= resources.length) { return null; } itemBuffer.add(END_OF_RESOURCE_MARKER); delegate.close(new ExecutionContext()); - delegate.setResource(resources[currentResourceIndex]); + delegate.setResource(resources[index.currentResource]); delegate.open(new ExecutionContext()); item = delegate.read(); @@ -114,8 +92,6 @@ public class MultiResourceItemReader extends ExecutionContextUserSupport impleme itemBuffer.add(item); - currentItemIndex++; - return item; } @@ -125,33 +101,25 @@ public class MultiResourceItemReader extends ExecutionContextUserSupport impleme * @return next item from buffer */ private Object readBufferedItem() { + Object buffered = itemBufferIterator.next(); while (buffered == END_OF_RESOURCE_MARKER) { - incrementResourceIndex(); + index.incrementResourceCount(); buffered = itemBufferIterator.next(); } - currentItemIndex++; + + if (!itemBufferIterator.hasNext()) { + // buffer is exhausted, continue reading from file + shouldReadBuffer = false; + itemBufferIterator = null; + } return buffered; } - /** - * Adjust indexes for next resource. - */ - private void incrementResourceIndex() { - currentResourceIndex++; - currentItemIndex = 0; - } - - /** - * Clears the item buffer and cancels reading from buffer if it applies. - * - * @see ItemReader#mark() - */ public void mark() throws MarkFailedException { emptyBuffer(); - lastMarkedResourceIndex = currentResourceIndex; - lastMarkedItemIndex = currentItemIndex; + index.mark(); delegate.mark(); } @@ -176,10 +144,11 @@ public class MultiResourceItemReader extends ExecutionContextUserSupport impleme * @see ItemReader#reset() */ public void reset() throws ResetFailedException { - shouldReadBuffer = true; - itemBufferIterator = itemBuffer.listIterator(); - currentResourceIndex = lastMarkedResourceIndex; - currentItemIndex = lastMarkedItemIndex; + if (!itemBuffer.isEmpty()) { + shouldReadBuffer = true; + itemBufferIterator = itemBuffer.listIterator(); + } + index.reset(); } /** @@ -188,13 +157,10 @@ public class MultiResourceItemReader extends ExecutionContextUserSupport impleme */ public void close(ExecutionContext executionContext) throws ItemStreamException { shouldReadBuffer = false; - lastMarkedResourceIndex = 0; - currentResourceIndex = 0; - currentItemIndex = 0; - lastMarkedItemIndex = 0; itemBufferIterator = null; + index = new MultiResourceIndex(); itemBuffer.clear(); - delegate.close(executionContext); + delegate.close(new ExecutionContext()); } /** @@ -203,22 +169,14 @@ public class MultiResourceItemReader extends ExecutionContextUserSupport impleme */ public void open(ExecutionContext executionContext) throws ItemStreamException { - if (executionContext.containsKey(getKey(RESOURCE_INDEX))) { - currentResourceIndex = Long.valueOf(executionContext.getLong(getKey(RESOURCE_INDEX))).intValue(); - lastMarkedResourceIndex = currentResourceIndex; - } + index.open(executionContext); - if (executionContext.containsKey(getKey(ITEM_COUNT))) { - currentItemIndex = executionContext.getLong(getKey(ITEM_COUNT)); - lastMarkedItemIndex = currentItemIndex; - } - - delegate.setResource(resources[currentResourceIndex]); + delegate.setResource(resources[index.currentResource]); delegate.open(new ExecutionContext()); try { - for (int i = 0; i < currentItemIndex; i++) { + for (int i = 0; i < index.currentItem; i++) { delegate.read(); delegate.mark(); } @@ -233,8 +191,7 @@ public class MultiResourceItemReader extends ExecutionContextUserSupport impleme */ public void update(ExecutionContext executionContext) throws ItemStreamException { if (saveState) { - executionContext.putLong(getKey(RESOURCE_INDEX), currentResourceIndex); - executionContext.putLong(getKey(ITEM_COUNT), currentItemIndex); + index.update(executionContext); } } @@ -267,4 +224,56 @@ public class MultiResourceItemReader extends ExecutionContextUserSupport impleme this.saveState = saveState; } + /** + * Facilitates keeping track of the position within multi-resource input. + */ + private class MultiResourceIndex { + + private static final String RESOURCE_KEY = "resourceIndex"; + + private static final String ITEM_KEY = "itemIndex"; + + private int currentResource = 0; + + private int markedResource = 0; + + private long currentItem = 0; + + private long markedItem = 0; + + public void incrementItemCount() { + currentItem++; + } + + public void incrementResourceCount() { + currentResource++; + currentItem = 0; + } + + public void mark() { + markedResource = currentResource; + markedItem = currentItem; + } + + public void reset() { + currentResource = markedResource; + currentItem = markedItem; + } + + public void open(ExecutionContext ctx) { + if (ctx.containsKey(getKey(RESOURCE_KEY))) { + currentResource = Long.valueOf(ctx.getLong(getKey(RESOURCE_KEY))).intValue(); + } + + if (ctx.containsKey(getKey(ITEM_KEY))) { + currentItem = ctx.getLong(getKey(ITEM_KEY)); + } + } + + public void update(ExecutionContext ctx) { + ctx.putLong(getKey(RESOURCE_KEY), index.currentResource); + ctx.putLong(getKey(ITEM_KEY), index.currentItem); + } + } + }