From 2f17add6b7ab13184fd7df7f9389148fc0e706c5 Mon Sep 17 00:00:00 2001 From: robokaso Date: Wed, 7 May 2008 08:03:00 +0000 Subject: [PATCH] RESOLVED - BATCH-613: StaxEventItemReader can run out of memory the reader now buffers items instead of xml events --- .../batch/item/xml/StaxEventItemReader.java | 60 +++++++++++++++---- .../batch/item/CommonItemReaderTests.java | 4 ++ .../item/xml/StaxEventItemReaderTests.java | 23 ------- 3 files changed, 52 insertions(+), 35 deletions(-) diff --git a/spring-batch-infrastructure/src/main/java/org/springframework/batch/item/xml/StaxEventItemReader.java b/spring-batch-infrastructure/src/main/java/org/springframework/batch/item/xml/StaxEventItemReader.java index 246dd901a..16e519d9f 100644 --- a/spring-batch-infrastructure/src/main/java/org/springframework/batch/item/xml/StaxEventItemReader.java +++ b/spring-batch-infrastructure/src/main/java/org/springframework/batch/item/xml/StaxEventItemReader.java @@ -2,6 +2,9 @@ package org.springframework.batch.item.xml; import java.io.IOException; import java.io.InputStream; +import java.util.ArrayList; +import java.util.Iterator; +import java.util.List; import javax.xml.namespace.QName; import javax.xml.stream.XMLEventReader; @@ -16,9 +19,7 @@ import org.springframework.batch.item.ItemStream; import org.springframework.batch.item.ItemStreamException; import org.springframework.batch.item.ReaderNotOpenException; import org.springframework.batch.item.xml.stax.DefaultFragmentEventReader; -import org.springframework.batch.item.xml.stax.DefaultTransactionalEventReader; import org.springframework.batch.item.xml.stax.FragmentEventReader; -import org.springframework.batch.item.xml.stax.TransactionalEventReader; import org.springframework.beans.factory.InitializingBean; import org.springframework.core.io.Resource; import org.springframework.dao.DataAccessResourceFailureException; @@ -26,7 +27,7 @@ import org.springframework.util.Assert; import org.springframework.util.ClassUtils; /** - * Input source for reading XML input based on StAX. + * Item reader for reading XML input based on StAX. * * It extracts fragments from the input XML document which correspond to records * for processing. The fragments are wrapped with StartDocument and EndDocument @@ -42,7 +43,7 @@ public class StaxEventItemReader extends ExecutionContextUserSupport implements private FragmentEventReader fragmentReader; - private TransactionalEventReader txReader; + private XMLEventReader eventReader; private EventReaderDeserializer eventReaderDeserializer; @@ -59,6 +60,15 @@ public class StaxEventItemReader extends ExecutionContextUserSupport implements private long currentRecordCount = 0; private boolean saveState = false; + + private List buffer = new ArrayList(); + + private Iterator bufferIterator = null; + + /** + * indicates the reader has been shouldReadBuffer and should read items from buffer + */ + private boolean shouldReadBuffer = false; public StaxEventItemReader() { setName(ClassUtils.getShortName(StaxEventItemReader.class)); @@ -74,15 +84,29 @@ public class StaxEventItemReader extends ExecutionContextUserSupport implements if (!initialized) { throw new ReaderNotOpenException("Reader must be open before it can be read."); } - Object item = null; currentRecordCount++; + + // read from buffer after rollback + if (shouldReadBuffer) { + if (bufferIterator.hasNext()) { + return bufferIterator.next(); + } else { + // buffer is exhausted, continue reading from file + shouldReadBuffer = false; + bufferIterator = null; + } + } + + Object item = null; + if (moveCursorToNextFragment(fragmentReader)) { fragmentReader.markStartFragment(); item = eventReaderDeserializer.deserializeFragment(fragmentReader); fragmentReader.markFragmentProcessed(); } + buffer.add(item); if (item == null) { currentRecordCount--; } @@ -92,6 +116,7 @@ public class StaxEventItemReader extends ExecutionContextUserSupport implements public void close(ExecutionContext executionContext) { initialized = false; currentRecordCount = 0; + clearBuffer(); try { if (fragmentReader != null) { fragmentReader.close(); @@ -117,9 +142,9 @@ public class StaxEventItemReader extends ExecutionContextUserSupport implements try { inputStream = resource.getInputStream(); - txReader = new DefaultTransactionalEventReader(XMLInputFactory.newInstance().createXMLEventReader( - inputStream)); - fragmentReader = new DefaultFragmentEventReader(txReader); + eventReader = XMLInputFactory.newInstance().createXMLEventReader( + inputStream); + fragmentReader = new DefaultFragmentEventReader(eventReader); } catch (XMLStreamException xse) { throw new DataAccessResourceFailureException("Unable to create XML reader", xse); @@ -134,8 +159,9 @@ public class StaxEventItemReader extends ExecutionContextUserSupport implements int REASONABLE_ADHOC_COMMIT_FREQUENCY = 100; while (currentRecordCount <= restoredRecordCount) { currentRecordCount++; + if (currentRecordCount % REASONABLE_ADHOC_COMMIT_FREQUENCY == 0) { - txReader.onCommit(); // reset the history buffer + mark(); // clear the history buffer } if (!fragmentReader.hasNext()) { throw new ItemStreamException("Restore point must be before end of input"); @@ -143,7 +169,7 @@ public class StaxEventItemReader extends ExecutionContextUserSupport implements fragmentReader.next(); moveCursorToNextFragment(fragmentReader); } - mark(); // reset the history buffer + mark(); // clear the history buffer } } @@ -236,7 +262,8 @@ public class StaxEventItemReader extends ExecutionContextUserSupport implements */ public void mark() { lastCommitPointRecordCount = currentRecordCount; - txReader.onCommit(); + clearBuffer(); + shouldReadBuffer = false; } /* @@ -246,7 +273,8 @@ public class StaxEventItemReader extends ExecutionContextUserSupport implements */ public void reset() { currentRecordCount = lastCommitPointRecordCount; - txReader.onRollback(); + shouldReadBuffer = true; + bufferIterator = buffer.listIterator(); fragmentReader.reset(); } @@ -260,4 +288,12 @@ public class StaxEventItemReader extends ExecutionContextUserSupport implements public void setSaveState(boolean saveState) { this.saveState = saveState; } + + /** + * Clear the buffer and release the iterator. + */ + private void clearBuffer() { + buffer.clear(); + bufferIterator = null; + } } diff --git a/spring-batch-infrastructure/src/test/java/org/springframework/batch/item/CommonItemReaderTests.java b/spring-batch-infrastructure/src/test/java/org/springframework/batch/item/CommonItemReaderTests.java index abea1b0ef..921c52735 100644 --- a/spring-batch-infrastructure/src/test/java/org/springframework/batch/item/CommonItemReaderTests.java +++ b/spring-batch-infrastructure/src/test/java/org/springframework/batch/item/CommonItemReaderTests.java @@ -74,6 +74,10 @@ public abstract class CommonItemReaderTests extends TestCase { Foo foo5 = (Foo) tested.read(); assertEquals(5, foo5.getValue()); + tested.reset(); + + assertEquals(foo5, tested.read()); + assertNull(tested.read()); } diff --git a/spring-batch-infrastructure/src/test/java/org/springframework/batch/item/xml/StaxEventItemReaderTests.java b/spring-batch-infrastructure/src/test/java/org/springframework/batch/item/xml/StaxEventItemReaderTests.java index 09344f0eb..f5857a3e8 100644 --- a/spring-batch-infrastructure/src/test/java/org/springframework/batch/item/xml/StaxEventItemReaderTests.java +++ b/spring-batch-infrastructure/src/test/java/org/springframework/batch/item/xml/StaxEventItemReaderTests.java @@ -162,17 +162,6 @@ public class StaxEventItemReaderTests extends TestCase { assertEquals(second, source.read()); - // rollback while deserializing record - source.reset(); - source.setFragmentDeserializer(new ExceptionFragmentDeserializer()); - try { - source.read(); - } catch (Exception expected) { - source.reset(); - } - source.setFragmentDeserializer(deserializer); - - assertEquals(second, source.read()); } /** @@ -355,18 +344,6 @@ public class StaxEventItemReaderTests extends TestCase { return str.indexOf(searchStr) != -1; } - /** - * Moves cursor inside the fragment body and causes rollback. - */ - private static class ExceptionFragmentDeserializer implements EventReaderDeserializer { - - public Object deserializeFragment(XMLEventReader eventReader) { - eventReader.next(); - throw new RuntimeException(); - } - - } - private static class MockStaxEventItemReader extends StaxEventItemReader { private boolean openCalled = false;