diff --git a/spring-batch-infrastructure/src/main/java/org/springframework/batch/item/database/AbstractPagingItemReader.java b/spring-batch-infrastructure/src/main/java/org/springframework/batch/item/database/AbstractPagingItemReader.java index b2106aeac..b57618bf7 100644 --- a/spring-batch-infrastructure/src/main/java/org/springframework/batch/item/database/AbstractPagingItemReader.java +++ b/spring-batch-infrastructure/src/main/java/org/springframework/batch/item/database/AbstractPagingItemReader.java @@ -94,30 +94,30 @@ public abstract class AbstractPagingItemReader extends AbstractItemCountingIt @Override protected T doRead() throws Exception { - if (results == null || current.get() >= pageSize) { + synchronized (lock) { - if (logger.isDebugEnabled()) { - logger.debug("Reading page " + getPage()); - } + if (results == null || current.get() >= pageSize) { - synchronized (lock) { - if (results == null || current.get() >= pageSize) { - doReadPage(); - page++; - if (current.get() >= pageSize) { - current.set(0); - } + if (logger.isDebugEnabled()) { + logger.debug("Reading page " + getPage()); } + + doReadPage(); + page++; + if (current.get() >= pageSize) { + current.set(0); + } + } - } + int next = current.getAndIncrement(); + if (next < results.size()) { + return results.get(next); + } + else { + return null; + } - int next = current.getAndIncrement(); - if (next < results.size()) { - return results.get(next); - } - else { - return null; } } diff --git a/spring-batch-infrastructure/src/test/java/org/springframework/batch/item/database/IbatisPagingItemReaderAsyncTests.java b/spring-batch-infrastructure/src/test/java/org/springframework/batch/item/database/IbatisPagingItemReaderAsyncTests.java index 7a486de68..639111e1e 100644 --- a/spring-batch-infrastructure/src/test/java/org/springframework/batch/item/database/IbatisPagingItemReaderAsyncTests.java +++ b/spring-batch-infrastructure/src/test/java/org/springframework/batch/item/database/IbatisPagingItemReaderAsyncTests.java @@ -4,7 +4,9 @@ import static org.junit.Assert.assertEquals; import static org.junit.Assert.assertNotNull; import java.util.ArrayList; +import java.util.HashSet; import java.util.List; +import java.util.Set; import java.util.concurrent.Callable; import java.util.concurrent.CompletionService; import java.util.concurrent.ExecutionException; @@ -101,7 +103,7 @@ public class IbatisPagingItemReaderAsyncTests { Foo next = null; do { next = reader.read(); - Thread.sleep(10L); + Thread.sleep(10L); // try to make it fairer logger.debug("Reading item: " + next); if (next != null) { list.add(next); @@ -112,14 +114,17 @@ public class IbatisPagingItemReaderAsyncTests { }); } int count = 0; + Set results = new HashSet(); for (int i = 0; i < THREAD_COUNT; i++) { List items = completionService.take().get(); count += items.size(); logger.debug("Finished items count: " + items.size()); logger.debug("Finished items: " + items); assertNotNull(items); + results.addAll(items); } assertEquals(ITEM_COUNT, count); + assertEquals(ITEM_COUNT, results.size()); reader.close(); } @@ -149,5 +154,4 @@ public class IbatisPagingItemReaderAsyncTests { return (SqlMapClient) factory.getObject(); } - } diff --git a/spring-batch-infrastructure/src/test/java/org/springframework/batch/item/database/JdbcPagingItemReaderAsyncTests.java b/spring-batch-infrastructure/src/test/java/org/springframework/batch/item/database/JdbcPagingItemReaderAsyncTests.java index 125336be5..312be6f0a 100644 --- a/spring-batch-infrastructure/src/test/java/org/springframework/batch/item/database/JdbcPagingItemReaderAsyncTests.java +++ b/spring-batch-infrastructure/src/test/java/org/springframework/batch/item/database/JdbcPagingItemReaderAsyncTests.java @@ -6,7 +6,9 @@ import static org.junit.Assert.assertNotNull; import java.sql.ResultSet; import java.sql.SQLException; import java.util.ArrayList; +import java.util.HashSet; import java.util.List; +import java.util.Set; import java.util.concurrent.Callable; import java.util.concurrent.CompletionService; import java.util.concurrent.ExecutionException; @@ -118,14 +120,17 @@ public class JdbcPagingItemReaderAsyncTests { }); } int count = 0; + Set results = new HashSet(); for (int i = 0; i < THREAD_COUNT; i++) { List items = completionService.take().get(); count += items.size(); logger.debug("Finished items count: " + items.size()); logger.debug("Finished items: " + items); assertNotNull(items); + results.addAll(items); } assertEquals(ITEM_COUNT, count); + assertEquals(ITEM_COUNT, results.size()); } protected ItemReader getItemReader() throws Exception { diff --git a/spring-batch-infrastructure/src/test/java/org/springframework/batch/item/database/JpaPagingItemReaderAsyncTests.java b/spring-batch-infrastructure/src/test/java/org/springframework/batch/item/database/JpaPagingItemReaderAsyncTests.java index 006c05018..34674a990 100644 --- a/spring-batch-infrastructure/src/test/java/org/springframework/batch/item/database/JpaPagingItemReaderAsyncTests.java +++ b/spring-batch-infrastructure/src/test/java/org/springframework/batch/item/database/JpaPagingItemReaderAsyncTests.java @@ -4,7 +4,9 @@ import static org.junit.Assert.assertEquals; import static org.junit.Assert.assertNotNull; import java.util.ArrayList; +import java.util.HashSet; import java.util.List; +import java.util.Set; import java.util.concurrent.Callable; import java.util.concurrent.CompletionService; import java.util.concurrent.ExecutionException; @@ -113,14 +115,17 @@ public class JpaPagingItemReaderAsyncTests { }); } int count = 0; + Set results = new HashSet(); for (int i = 0; i < THREAD_COUNT; i++) { List items = completionService.take().get(); count += items.size(); logger.debug("Finished items count: " + items.size()); logger.debug("Finished items: " + items); assertNotNull(items); + results.addAll(items); } assertEquals(ITEM_COUNT, count); + assertEquals(ITEM_COUNT, results.size()); reader.close(); }