diff --git a/spring-batch-samples/src/main/java/org/springframework/batch/sample/common/ProcessIndicatorItemWrapper.java b/spring-batch-samples/src/main/java/org/springframework/batch/sample/common/ProcessIndicatorItemWrapper.java new file mode 100644 index 000000000..d1b0c4ea1 --- /dev/null +++ b/spring-batch-samples/src/main/java/org/springframework/batch/sample/common/ProcessIndicatorItemWrapper.java @@ -0,0 +1,39 @@ +package org.springframework.batch.sample.common; + +/** + * Item wrapper useful in "process indicator" usecase, where input is marked as + * processed by the processor/writer. This requires passing a technical + * identifier of the input data so that it can be modified in later stages. + * + * @param item type + * + * @see StagingItemReader + * @see StagingItemProcessor + * + * @author Robert Kasanicky + */ +public class ProcessIndicatorItemWrapper { + + private long id; + + private T item; + + public ProcessIndicatorItemWrapper(long id, T item) { + this.id = id; + this.item = item; + } + + /** + * @return id identifying the input data (typically row in database) + */ + public long getId() { + return id; + } + + /** + * @return item (domain object for business processing) + */ + public T getItem() { + return item; + } +} diff --git a/spring-batch-samples/src/main/java/org/springframework/batch/sample/common/StagingItemProcessor.java b/spring-batch-samples/src/main/java/org/springframework/batch/sample/common/StagingItemProcessor.java new file mode 100644 index 000000000..ccf23b233 --- /dev/null +++ b/spring-batch-samples/src/main/java/org/springframework/batch/sample/common/StagingItemProcessor.java @@ -0,0 +1,55 @@ +package org.springframework.batch.sample.common; + +import javax.sql.DataSource; + +import org.springframework.batch.item.ItemProcessor; +import org.springframework.beans.factory.InitializingBean; +import org.springframework.dao.OptimisticLockingFailureException; +import org.springframework.jdbc.core.simple.SimpleJdbcOperations; +import org.springframework.jdbc.core.simple.SimpleJdbcTemplate; +import org.springframework.util.Assert; + +/** + * Marks the input row as 'processed'. (This change will rollback if there is + * problem later) + * + * @param item type + * + * @see StagingItemReader + * @see StagingItemWriter + * @see ProcessIndicatorItemWrapper + * + * @author Robert Kasanicky + */ +public class StagingItemProcessor implements ItemProcessor, T>, InitializingBean { + + private SimpleJdbcOperations jdbcTemplate; + + public void setJdbcTemplate(SimpleJdbcOperations jdbcTemplate) { + this.jdbcTemplate = jdbcTemplate; + } + + public void setDataSource(DataSource dataSource) { + this.jdbcTemplate = new SimpleJdbcTemplate(dataSource); + } + + public void afterPropertiesSet() throws Exception { + Assert.notNull(jdbcTemplate, "Either jdbcTemplate or dataSource must be set"); + } + + /** + * Use the technical identifier to mark the input row as processed and + * return unwrapped item. + */ + public T process(ProcessIndicatorItemWrapper wrapper) throws Exception { + + int count = jdbcTemplate.update("UPDATE BATCH_STAGING SET PROCESSED=? WHERE ID=? AND PROCESSED=?", + StagingItemWriter.DONE, wrapper.getId(), StagingItemWriter.NEW); + if (count != 1) { + throw new OptimisticLockingFailureException("The staging record with ID=" + wrapper.getId() + + " was updated concurrently when trying to mark as complete (updated " + count + " records."); + } + return wrapper.getItem(); + } + +} diff --git a/spring-batch-samples/src/main/java/org/springframework/batch/sample/common/StagingItemReader.java b/spring-batch-samples/src/main/java/org/springframework/batch/sample/common/StagingItemReader.java index 27820252f..e7a0d1dfd 100644 --- a/spring-batch-samples/src/main/java/org/springframework/batch/sample/common/StagingItemReader.java +++ b/spring-batch-samples/src/main/java/org/springframework/batch/sample/common/StagingItemReader.java @@ -34,7 +34,6 @@ import org.springframework.batch.item.ReaderNotOpenException; import org.springframework.beans.factory.DisposableBean; import org.springframework.beans.factory.InitializingBean; import org.springframework.dao.DataAccessException; -import org.springframework.dao.OptimisticLockingFailureException; import org.springframework.jdbc.core.simple.ParameterizedRowMapper; import org.springframework.jdbc.core.simple.SimpleJdbcTemplate; import org.springframework.util.Assert; @@ -42,8 +41,11 @@ import org.springframework.util.Assert; /** * Thread-safe database {@link ItemReader} implementing the process indicator * pattern. + * + * To achieve restartability use together with {@link StagingItemProcessor}. */ -public class StagingItemReader implements ItemReader, StepExecutionListener, InitializingBean, DisposableBean { +public class StagingItemReader implements ItemReader>, StepExecutionListener, + InitializingBean, DisposableBean { private static Log logger = LogFactory.getLog(StagingItemReader.class); @@ -90,7 +92,7 @@ public class StagingItemReader implements ItemReader, StepExecutionListene } - public T read() throws DataAccessException { + public ProcessIndicatorItemWrapper read() throws DataAccessException { if (!initialized) { throw new ReaderNotOpenException("ItemStream must be open before it can be read."); @@ -116,15 +118,7 @@ public class StagingItemReader implements ItemReader, StepExecutionListene } }, id); - // Update now - changes will rollback if there is a problem later. - int count = jdbcTemplate.update("UPDATE BATCH_STAGING SET PROCESSED=? WHERE ID=? AND PROCESSED=?", - StagingItemWriter.DONE, id, StagingItemWriter.NEW); - if (count != 1) { - throw new OptimisticLockingFailureException("The staging record with ID=" + id - + " was updated concurrently when trying to mark as complete (updated " + count + " records."); - } - - return result; + return new ProcessIndicatorItemWrapper(id, result); } diff --git a/spring-batch-samples/src/main/resources/jobs/parallelJob.xml b/spring-batch-samples/src/main/resources/jobs/parallelJob.xml index 12fb9d506..c2adf1c2d 100644 --- a/spring-batch-samples/src/main/resources/jobs/parallelJob.xml +++ b/spring-batch-samples/src/main/resources/jobs/parallelJob.xml @@ -10,8 +10,8 @@ http://www.springframework.org/schema/tx http://www.springframework.org/schema/tx/spring-tx-2.0.xsd"> - - + + @@ -35,7 +35,7 @@ - + @@ -45,6 +45,11 @@ + + + + + diff --git a/spring-batch-samples/src/test/java/org/springframework/batch/sample/common/StagingItemReaderTests.java b/spring-batch-samples/src/test/java/org/springframework/batch/sample/common/StagingItemReaderTests.java index b1dabd5ff..28dc8871f 100644 --- a/spring-batch-samples/src/test/java/org/springframework/batch/sample/common/StagingItemReaderTests.java +++ b/spring-batch-samples/src/test/java/org/springframework/batch/sample/common/StagingItemReaderTests.java @@ -65,15 +65,20 @@ public class StagingItemReaderTests { @Transactional @Test - public void testReaderUpdatesProcessIndicator() throws Exception { + public void testReaderWithProcessorUpdatesProcessIndicator() throws Exception { long id = simpleJdbcTemplate.queryForLong("SELECT MIN(ID) from BATCH_STAGING where JOB_ID=?", jobId); String before = simpleJdbcTemplate.queryForObject("SELECT PROCESSED from BATCH_STAGING where ID=?", String.class, id); assertEquals(StagingItemWriter.NEW, before); - String item = reader.read(); + ProcessIndicatorItemWrapper wrapper = reader.read(); + String item = wrapper.getItem(); assertEquals("FOO", item); + + StagingItemProcessor updater = new StagingItemProcessor(); + updater.setJdbcTemplate(simpleJdbcTemplate); + updater.process(wrapper); String after = simpleJdbcTemplate.queryForObject("SELECT PROCESSED from BATCH_STAGING where ID=?", String.class, id); @@ -89,7 +94,7 @@ public class StagingItemReaderTests { txTemplate.execute(new TransactionCallback() { public Object doInTransaction(TransactionStatus transactionStatus) { try { - testReaderUpdatesProcessIndicator(); + testReaderWithProcessorUpdatesProcessIndicator(); } catch (Exception e) { fail("Unxpected Exception: " + e); @@ -118,8 +123,8 @@ public class StagingItemReaderTests { String.class, id); assertEquals(StagingItemWriter.NEW, before); - Object item = reader.read(); - assertEquals("FOO", item); + ProcessIndicatorItemWrapper wrapper = reader.read(); + assertEquals("FOO", wrapper.getItem()); transactionStatus.setRollbackOnly();