diff --git a/spring-batch-samples/src/main/java/org/springframework/batch/sample/item/provider/StagingItemReader.java b/spring-batch-samples/src/main/java/org/springframework/batch/sample/item/provider/StagingItemReader.java index 2d1beb144..23015df83 100644 --- a/spring-batch-samples/src/main/java/org/springframework/batch/sample/item/provider/StagingItemReader.java +++ b/spring-batch-samples/src/main/java/org/springframework/batch/sample/item/provider/StagingItemReader.java @@ -26,12 +26,11 @@ import org.springframework.transaction.support.TransactionSynchronizationAdapter import org.springframework.transaction.support.TransactionSynchronizationManager; import org.springframework.util.Assert; -public class StagingItemReader extends JdbcDaoSupport implements - ItemReader, ResourceLifecycle, DisposableBean, +public class StagingItemReader extends JdbcDaoSupport implements ItemReader, ResourceLifecycle, DisposableBean, StepContextAware { // Key for buffer in transaction synchronization manager - private static final String BUFFER_KEY = StagingItemReader.class.getName()+".BUFFER"; + private static final String BUFFER_KEY = StagingItemReader.class.getName() + ".BUFFER"; private static Log logger = LogFactory.getLog(StagingItemReader.class); @@ -41,7 +40,7 @@ public class StagingItemReader extends JdbcDaoSupport implements private Object lock = new Object(); - private boolean initialized = false; + private volatile boolean initialized = false; private volatile Iterator keys; @@ -72,13 +71,14 @@ public class StagingItemReader extends JdbcDaoSupport implements * @see org.springframework.batch.io.driving.DrivingQueryItemReader#open() */ public void open() { - Assert.state(keys == null || initialized, - "Cannot open an already open StagingItemProvider" - + ", call close() first."); + Assert.state(keys == null || initialized, "Cannot open an already open StagingItemProvider" + + ", call close() first."); synchronized (lock) { - keys = retrieveKeys().iterator(); + if (keys == null) { + keys = retrieveKeys().iterator(); + logger.info("Keys obtained for staging."); + } } - logger.info("keys: " + keys); registerSynchronization(); initialized = true; } @@ -86,8 +86,7 @@ public class StagingItemReader extends JdbcDaoSupport implements /** * Callback for injection of the step context. * - * @param stepContext - * the stepContext to set + * @param stepContext the stepContext to set */ public void setStepContext(StepContext stepContext) { this.stepContext = stepContext; @@ -97,24 +96,19 @@ public class StagingItemReader extends JdbcDaoSupport implements synchronized (lock) { - return getJdbcTemplate() - .query( + return getJdbcTemplate().query( - "SELECT ID FROM BATCH_STAGING WHERE JOB_ID=? AND PROCESSED=? ORDER BY ID", + "SELECT ID FROM BATCH_STAGING WHERE JOB_ID=? AND PROCESSED=? ORDER BY ID", - new Object[] { - stepContext.getStepExecution() - .getJobExecution().getJobId(), - StagingItemProcessor.NEW }, + new Object[] { stepContext.getStepExecution().getJobExecution().getJobId(), StagingItemProcessor.NEW }, - new RowMapper() { - public Object mapRow(ResultSet rs, int rowNum) - throws SQLException { - return new Long(rs.getLong(1)); - } - } + new RowMapper() { + public Object mapRow(ResultSet rs, int rowNum) throws SQLException { + return new Long(rs.getLong(1)); + } + } - ); + ); } @@ -130,26 +124,19 @@ public class StagingItemReader extends JdbcDaoSupport implements if (id == null) { return null; } - Object result = getJdbcTemplate().queryForObject( - "SELECT VALUE FROM BATCH_STAGING WHERE ID=?", + Object result = getJdbcTemplate().queryForObject("SELECT VALUE FROM BATCH_STAGING WHERE ID=?", new Object[] { id }, new RowMapper() { - public Object mapRow(ResultSet rs, int rowNum) - throws SQLException { + public Object mapRow(ResultSet rs, int rowNum) throws SQLException { byte[] blob = lobHandler.getBlobAsBytes(rs, 1); return SerializationUtils.deserialize(blob); } }); // Update now - changes will rollback if there is a problem later. - int count = getJdbcTemplate() - .update( - "UPDATE BATCH_STAGING SET PROCESSED=? WHERE ID=? AND PROCESSED=?", - new Object[] { StagingItemProcessor.DONE, id, - StagingItemProcessor.NEW }); + int count = getJdbcTemplate().update("UPDATE BATCH_STAGING SET PROCESSED=? WHERE ID=? AND PROCESSED=?", + new Object[] { StagingItemProcessor.DONE, id, StagingItemProcessor.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."); + throw new OptimisticLockingFailureException("The staging record with ID=" + id + + " was updated concurrently when trying to mark as complete (updated " + count + " records."); } return result; } @@ -169,7 +156,8 @@ public class StagingItemReader extends JdbcDaoSupport implements logger.debug("Retrieved key from list: " + key); } } - } else { + } + else { logger.debug("Retrieved key from buffer: " + key); } return key; @@ -178,11 +166,9 @@ public class StagingItemReader extends JdbcDaoSupport implements private StagingBuffer getBuffer() { if (!TransactionSynchronizationManager.hasResource(BUFFER_KEY)) { - TransactionSynchronizationManager.bindResource(BUFFER_KEY, - new StagingBuffer()); + TransactionSynchronizationManager.bindResource(BUFFER_KEY, new StagingBuffer()); } - return (StagingBuffer) TransactionSynchronizationManager - .getResource(BUFFER_KEY); + return (StagingBuffer) TransactionSynchronizationManager.getResource(BUFFER_KEY); } public boolean recover(Object data, Throwable cause) { @@ -196,8 +182,7 @@ public class StagingItemReader extends JdbcDaoSupport implements * initialized. */ protected void registerSynchronization() { - BatchTransactionSynchronizationManager - .registerSynchronization(synchronization); + BatchTransactionSynchronizationManager.registerSynchronization(synchronization); } /* @@ -221,12 +206,12 @@ public class StagingItemReader extends JdbcDaoSupport implements /** * Encapsulates transaction events handling. */ - private class StagingInputTransactionSynchronization extends - TransactionSynchronizationAdapter { + private class StagingInputTransactionSynchronization extends TransactionSynchronizationAdapter { public void afterCompletion(int status) { if (status == TransactionSynchronization.STATUS_ROLLED_BACK) { transactionRolledBack(); - } else if (status == TransactionSynchronization.STATUS_COMMITTED) { + } + else if (status == TransactionSynchronization.STATUS_COMMITTED) { transactionCommitted(); } } @@ -235,6 +220,7 @@ public class StagingItemReader extends JdbcDaoSupport implements private class StagingBuffer { private List list = new ArrayList(); + private Iterator iter = new ArrayList().iterator(); public Long next() { diff --git a/spring-batch-samples/src/main/resources/jobs/fixedLengthImportJob.xml b/spring-batch-samples/src/main/resources/jobs/fixedLengthImportJob.xml index e28698034..1c2f61f55 100644 --- a/spring-batch-samples/src/main/resources/jobs/fixedLengthImportJob.xml +++ b/spring-batch-samples/src/main/resources/jobs/fixedLengthImportJob.xml @@ -76,8 +76,6 @@ - - @@ -87,7 +85,4 @@ - - - \ No newline at end of file diff --git a/spring-batch-samples/src/main/resources/jobs/parallelJob.xml b/spring-batch-samples/src/main/resources/jobs/parallelJob.xml index aba1a9034..cb4349506 100644 --- a/spring-batch-samples/src/main/resources/jobs/parallelJob.xml +++ b/spring-batch-samples/src/main/resources/jobs/parallelJob.xml @@ -1,9 +1,9 @@ - - + - - - + - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - SELECT games.player_id, games.year_no, SUM(COMPLETES), - SUM(ATTEMPTS), SUM(PASSING_YARDS), SUM(PASSING_TD), - SUM(INTERCEPTIONS), SUM(RUSHES), SUM(RUSH_YARDS), - SUM(RECEPTIONS), SUM(RECEPTIONS_YARDS), SUM(TOTAL_TD) - from games, players where players.player_id = - games.player_id group by games.player_id, games.year_no - - - - - - - - - - - - + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + \ No newline at end of file diff --git a/spring-batch-samples/src/main/resources/log4j.properties b/spring-batch-samples/src/main/resources/log4j.properties index fc2243d7e..4cc711578 100644 --- a/spring-batch-samples/src/main/resources/log4j.properties +++ b/spring-batch-samples/src/main/resources/log4j.properties @@ -2,11 +2,20 @@ log4j.appender.stdout=org.apache.log4j.ConsoleAppender log4j.appender.stdout.Target=System.out log4j.appender.stdout.layout=org.apache.log4j.PatternLayout -log4j.appender.stdout.layout.ConversionPattern=%d{ABSOLUTE} %5p %c{1}:%L - %m%n +log4j.appender.stdout.layout.ConversionPattern=%d{ABSOLUTE} %5p %t %c{1}:%L - %m%n + +log4j.appender.chainsaw=org.apache.log4j.RollingFileAppender +log4j.appender.chainsaw.File=out.xml +log4j.appender.chainsaw.Append=false +log4j.appender.chainsaw.Threshold=debug +log4j.appender.chainsaw.MaxFileSize=10MB +log4j.appender.chainsaw.MaxBackupIndex=2 +log4j.appender.chainsaw.layout=org.apache.log4j.xml.XMLLayout ### set log levels - for more verbose logging change 'info' to 'debug' ### log4j.rootLogger=info, stdout +# log4j.rootLogger=info, stdout, chainsaw ### enable the following line if you want to track down connection ### ### leakages when using DriverManagerConnectionProvider ### diff --git a/spring-batch-samples/src/test/java/org/springframework/batch/sample/FootballJobFunctionalTests.java b/spring-batch-samples/src/test/java/org/springframework/batch/sample/FootballJobFunctionalTests.java index d944d35bc..48e1cafd3 100644 --- a/spring-batch-samples/src/test/java/org/springframework/batch/sample/FootballJobFunctionalTests.java +++ b/spring-batch-samples/src/test/java/org/springframework/batch/sample/FootballJobFunctionalTests.java @@ -2,7 +2,6 @@ package org.springframework.batch.sample; import javax.sql.DataSource; -import org.springframework.batch.sample.item.processor.StagingItemProcessor; import org.springframework.jdbc.core.JdbcOperations; import org.springframework.jdbc.core.JdbcTemplate; @@ -16,21 +15,12 @@ public class FootballJobFunctionalTests extends } protected String[] getConfigLocations() { - return new String[] { "jobs/parallelJob.xml##" }; + return new String[] { "jobs/footballJob.xml##" }; } protected void validatePostConditions() throws Exception { - int count; - count = jdbcTemplate.queryForInt( - "SELECT COUNT(*) from BATCH_STAGING where PROCESSED=?", - new Object[] {StagingItemProcessor.NEW}); - assertEquals(0, count); - int total = jdbcTemplate.queryForInt( - "SELECT COUNT(*) from BATCH_STAGING"); - count = jdbcTemplate.queryForInt( - "SELECT COUNT(*) from BATCH_STAGING where PROCESSED=?", - new Object[] {StagingItemProcessor.DONE}); - assertEquals(total, count); + int count = jdbcTemplate.queryForInt("SELECT COUNT(*) from PLAYER_SUMMARY"); + assertTrue(count>0); } } diff --git a/spring-batch-samples/src/test/java/org/springframework/batch/sample/ParallelJobFunctionalTests.java b/spring-batch-samples/src/test/java/org/springframework/batch/sample/ParallelJobFunctionalTests.java new file mode 100644 index 000000000..0b9178301 --- /dev/null +++ b/spring-batch-samples/src/test/java/org/springframework/batch/sample/ParallelJobFunctionalTests.java @@ -0,0 +1,36 @@ +package org.springframework.batch.sample; + +import javax.sql.DataSource; + +import org.springframework.batch.sample.item.processor.StagingItemProcessor; +import org.springframework.jdbc.core.JdbcOperations; +import org.springframework.jdbc.core.JdbcTemplate; + +public class ParallelJobFunctionalTests extends + AbstractValidatingBatchLauncherTests { + + private JdbcOperations jdbcTemplate; + + public void setDataSource(DataSource dataSource) { + this.jdbcTemplate = new JdbcTemplate(dataSource); + } + + protected String[] getConfigLocations() { + return new String[] { "jobs/parallelJob.xml" }; + } + + protected void validatePostConditions() throws Exception { + int count; + count = jdbcTemplate.queryForInt( + "SELECT COUNT(*) from BATCH_STAGING where PROCESSED=?", + new Object[] {StagingItemProcessor.NEW}); + assertEquals(0, count); + int total = jdbcTemplate.queryForInt( + "SELECT COUNT(*) from BATCH_STAGING"); + count = jdbcTemplate.queryForInt( + "SELECT COUNT(*) from BATCH_STAGING where PROCESSED=?", + new Object[] {StagingItemProcessor.DONE}); + assertEquals(total, count); + } + +} diff --git a/spring-batch-samples/src/test/resources/log4j.properties b/spring-batch-samples/src/test/resources/log4j.properties deleted file mode 100644 index 6657107d3..000000000 --- a/spring-batch-samples/src/test/resources/log4j.properties +++ /dev/null @@ -1,32 +0,0 @@ -### direct log messages to stdout ### -log4j.appender.stdout=org.apache.log4j.ConsoleAppender -log4j.appender.stdout.Target=System.out -log4j.appender.stdout.layout=org.apache.log4j.PatternLayout -log4j.appender.stdout.layout.ConversionPattern=%d{ABSOLUTE} %5p %c{1}:%L - %m%n - -log4j.appender.chainsaw=org.apache.log4j.RollingFileAppender -log4j.appender.chainsaw.File=out.xml -log4j.appender.chainsaw.Append=false -log4j.appender.chainsaw.Threshold=debug -log4j.appender.chainsaw.MaxFileSize=10MB -log4j.appender.chainsaw.MaxBackupIndex=2 -log4j.appender.chainsaw.layout=org.apache.log4j.xml.XMLLayout - -### set log levels - for more verbose logging change 'info' to 'debug' ### - -log4j.rootLogger=info, stdout -# log4j.rootLogger=info, stdout, chainsaw - -### enable the following line if you want to track down connection ### -### leakages when using DriverManagerConnectionProvider ### -#log4j.logger.org.hibernate.connection.DriverManagerConnectionProvider=trace - -### enable spring -log4j.logger.org.springframework=error -log4j.logger.org.springframework.batch.sample=info - -### debug your specific package or classes with the following example -log4j.logger.org.springframework.batch.sample.module.OrderDataProvider=debug -log4j.logger.org.springframework.batch.container.common.module.process.support.DefaultXmlDataProvider=debug -#log4j.logger.org.springframework.jdbc.core=debug -#log4j.logger.org.springframework.transaction=debug \ No newline at end of file