diff --git a/spring-batch-core/src/main/java/org/springframework/batch/core/domain/StepExecution.java b/spring-batch-core/src/main/java/org/springframework/batch/core/domain/StepExecution.java index e1b69293e..4f2032589 100644 --- a/spring-batch-core/src/main/java/org/springframework/batch/core/domain/StepExecution.java +++ b/spring-batch-core/src/main/java/org/springframework/batch/core/domain/StepExecution.java @@ -346,7 +346,8 @@ public class StepExecution extends Entity { */ public synchronized void apply(StepContribution contribution) { taskCount += contribution.getTaskCount(); - executionContext = contribution.getExecutionContext(); + // TODO: this should not be necessary - the step decides + // executionContext = contribution.getExecutionContext(); commitCount += contribution.getCommitCount(); } diff --git a/spring-batch-core/src/test/java/org/springframework/batch/core/domain/StepExecutionTests.java b/spring-batch-core/src/test/java/org/springframework/batch/core/domain/StepExecutionTests.java index b7119db0b..c6ec5038d 100644 --- a/spring-batch-core/src/test/java/org/springframework/batch/core/domain/StepExecutionTests.java +++ b/spring-batch-core/src/test/java/org/springframework/batch/core/domain/StepExecutionTests.java @@ -16,11 +16,14 @@ package org.springframework.batch.core.domain; import java.util.Date; +import java.util.HashSet; +import java.util.Set; import junit.framework.TestCase; import org.springframework.batch.item.ExecutionContext; import org.springframework.batch.repeat.ExitStatus; +import org.springframework.batch.support.PropertiesConverter; /** * @author Dave Syer @@ -260,6 +263,14 @@ public class StepExecutionTests extends TestCase { .hashCode()); } + public void testHashCodeViaHashSet() throws Exception { + Set set = new HashSet(); + set.add(execution); + assertTrue(set.contains(execution)); + execution.setExecutionContext(new ExecutionContext(PropertiesConverter.stringToProperties("foo=bar"))); + assertTrue(set.contains(execution)); + } + private StepExecution newStepExecution(Long long1, Long long2) { JobInstance job = new JobInstance(new Long(3), new JobParameters()); StepInstance step = new StepInstance(job, "foo", long1); diff --git a/spring-batch-execution/src/main/java/org/springframework/batch/execution/scope/SimpleStepContext.java b/spring-batch-execution/src/main/java/org/springframework/batch/execution/scope/SimpleStepContext.java index d74a47db8..2686b33f9 100644 --- a/spring-batch-execution/src/main/java/org/springframework/batch/execution/scope/SimpleStepContext.java +++ b/spring-batch-execution/src/main/java/org/springframework/batch/execution/scope/SimpleStepContext.java @@ -25,8 +25,6 @@ import java.util.Set; import org.springframework.batch.core.domain.StepExecution; import org.springframework.batch.io.exception.BatchCriticalException; -import org.springframework.batch.item.ItemStream; -import org.springframework.batch.item.ExecutionContext; import org.springframework.batch.item.stream.StreamManager; import org.springframework.batch.repeat.context.SynchronizedAttributeAccessor; @@ -46,8 +44,6 @@ public class SimpleStepContext extends SynchronizedAttributeAccessor implements private StreamManager streamManager; - private ExecutionContext executionContext; - /** * Default constructor. */ @@ -79,11 +75,6 @@ public class SimpleStepContext extends SynchronizedAttributeAccessor implements */ public void setAttribute(String name, Object value) { super.setAttribute(name, value); - if (streamManager != null && (value instanceof ItemStream)) { - ItemStream stream = (ItemStream) value; - stream.open(); - streamManager.register(this, stream, executionContext); - } } /* @@ -192,18 +183,4 @@ public class SimpleStepContext extends SynchronizedAttributeAccessor implements return stepExecution; } - /* (non-Javadoc) - * @see org.springframework.batch.item.ExecutionContextProvider#getExecutionContext() - */ - public ExecutionContext getExecutionContext() { - return streamManager.getExecutionContext(this); - } - - /* (non-Javadoc) - * @see org.springframework.batch.execution.scope.StepContext#restoreFrom(org.springframework.batch.item.ExecutionContext) - */ - public void restoreFrom(ExecutionContext executionContext) { - this.executionContext = executionContext; - } - } diff --git a/spring-batch-execution/src/main/java/org/springframework/batch/execution/scope/StepContext.java b/spring-batch-execution/src/main/java/org/springframework/batch/execution/scope/StepContext.java index 645b904aa..20f23ddbc 100644 --- a/spring-batch-execution/src/main/java/org/springframework/batch/execution/scope/StepContext.java +++ b/spring-batch-execution/src/main/java/org/springframework/batch/execution/scope/StepContext.java @@ -16,9 +16,6 @@ package org.springframework.batch.execution.scope; import org.springframework.batch.core.domain.StepExecution; -import org.springframework.batch.item.ItemStream; -import org.springframework.batch.item.ExecutionContext; -import org.springframework.batch.item.ExecutionContextProvider; import org.springframework.core.AttributeAccessor; /** @@ -27,7 +24,7 @@ import org.springframework.core.AttributeAccessor; * @author Dave Syer * */ -public interface StepContext extends AttributeAccessor, ExecutionContextProvider { +public interface StepContext extends AttributeAccessor { /** * Accessor for the {@link StepExecution} associated with the currently @@ -54,14 +51,4 @@ public interface StepContext extends AttributeAccessor, ExecutionContextProvider */ void close(); - /** - * Provide the stream context needed to restore {@link ItemStream} - * instances. Implementations are free to use this as necessary (e.g. lazily - * if all the streams are not available at once). If this is not set the - * streams will simply not be initialised and repositioned for restart - * (which is sometimes desirable). - * - * @param executionContext - */ - void restoreFrom(ExecutionContext executionContext); } \ No newline at end of file diff --git a/spring-batch-execution/src/main/java/org/springframework/batch/execution/step/simple/ChunkedStep.java b/spring-batch-execution/src/main/java/org/springframework/batch/execution/step/simple/ChunkedStep.java index 6d5962dc5..b6137beb5 100644 --- a/spring-batch-execution/src/main/java/org/springframework/batch/execution/step/simple/ChunkedStep.java +++ b/spring-batch-execution/src/main/java/org/springframework/batch/execution/step/simple/ChunkedStep.java @@ -25,9 +25,10 @@ import org.springframework.batch.core.domain.Chunker; import org.springframework.batch.core.domain.ChunkingResult; import org.springframework.batch.core.domain.Dechunker; import org.springframework.batch.core.domain.DechunkingResult; -import org.springframework.batch.core.domain.SkippedItemHandler; import org.springframework.batch.core.domain.ItemSkipPolicy; import org.springframework.batch.core.domain.JobInterruptedException; +import org.springframework.batch.core.domain.SkippedItemHandler; +import org.springframework.batch.core.domain.Step; import org.springframework.batch.core.domain.StepContribution; import org.springframework.batch.core.domain.StepExecution; import org.springframework.batch.core.domain.StepInstance; @@ -60,44 +61,64 @@ import org.springframework.transaction.TransactionStatus; import org.springframework.util.Assert; /** - *

Implementation of the {@link Step} interface that deals with input and output as 'chunks'. Reading is - * delegated to a {@link Chunker} that will read in a {@link Chunk} of items for processing. The number of - * items per chunks is configurable as the chunk size. Once the chunk has been read, any errors encountered - * while reading (usually skipped unless configured not to) will be logged out via the {@link SkippedItemHandler}. - * The chunk will then be 'dechunked', which in most scenarios will mean delegating to an {@link ItemWriter} - * by writing out one chunk at a time. The transaction boundary is around this process. If any errors are - * encountered, the dechunking process will error out, leaving the decision for retrying the chunk up to - * a {@link RepeatTemplate}. This template is configurable, allowing for the number of retries and how long - * to wait between retries (backoff) to be set. Once dechunking has been finished, any errors not fatal to - * the chunk (usually because the error didn't invalidate the transaction) will also be written out via - * the {@link SkippedItemHandler}

+ *

+ * Implementation of the {@link Step} interface that deals with input and output + * as 'chunks'. Reading is delegated to a {@link Chunker} that will read in a + * {@link Chunk} of items for processing. The number of items per chunks is + * configurable as the chunk size. Once the chunk has been read, any errors + * encountered while reading (usually skipped unless configured not to) will be + * logged out via the {@link SkippedItemHandler}. The chunk will then be + * 'dechunked', which in most scenarios will mean delegating to an + * {@link ItemWriter} by writing out one chunk at a time. The transaction + * boundary is around this process. If any errors are encountered, the + * dechunking process will error out, leaving the decision for retrying the + * chunk up to a {@link RepeatTemplate}. This template is configurable, + * allowing for the number of retries and how long to wait between retries + * (backoff) to be set. Once dechunking has been finished, any errors not fatal + * to the chunk (usually because the error didn't invalidate the transaction) + * will also be written out via the {@link SkippedItemHandler} + *

* - *

Clients can use {@link RepeatListener}s in the step operations to intercept or listen to the iteration - * on a step-wide basis, for instance to get a callback when the step is complete. The open and close methods of - * could easily be done with AOP, however, notifications in between complete chunks (before and after) can - * be quite useful

+ *

+ * Clients can use {@link RepeatListener}s in the step operations to intercept + * or listen to the iteration on a step-wide basis, for instance to get a + * callback when the step is complete. The open and close methods of could + * easily be done with AOP, however, notifications in between complete chunks + * (before and after) can be quite useful + *

* - *

Repository Usage: The {@link JobRepository} is used extensively to store metadata about the run such as - * when the {@link StepExecution} was started, or the commit count.

+ *

+ * Repository Usage: The {@link JobRepository} is used extensively to store + * metadata about the run such as when the {@link StepExecution} was started, or + * the commit count. + *

* - *

Interruption: At various times while processing, the step will check to see if it has been interrupted - * by calling the {@link StepInterruptionPolicy}. This policy could check if thread.isInterupted() is true, - * or RepeatContext.isTerminateOnly() is set. It could even be a check to see if a 'stop file' has been added - * to a particular directory. If the step should finish, a {@link JobInterruptedException} is thrown, and the - * step will clean up, set the status of the {@link StepExecution} to 'STOPPED' and rethrow. + * Interruption: At various times while processing, the step will check to see + * if it has been interrupted by calling the {@link StepInterruptionPolicy}. + * This policy could check if thread.isInterupted() is true, or + * RepeatContext.isTerminateOnly() is set. It could even be a check to see if a + * 'stop file' has been added to a particular directory. If the step should + * finish, a {@link JobInterruptedException} is thrown, and the step will clean + * up, set the status of the {@link StepExecution} to 'STOPPED' and rethrow.ExitStatusClassification: Any number of fatal errors could be thrown during processing. In general, the - * framework must remain fairly dumb as to what error code these exceptions should translate to. By default - * it's a fairly generic 'FATAL_EXECUTION'. However, this may be insufficient for many scenarios. If an - * enterprise scheduler is used to kick off a batch job, the exit code is the only means of communication as - * to what action must be taken. It may also be the only result that many batch operators see as well. Therefore, - * an {@link ExitStatusExceptionClassifier} may be used to classify an exception to a particular exit code.

+ *

+ * ExitStatusClassification: Any number of fatal errors could be thrown during + * processing. In general, the framework must remain fairly dumb as to what + * error code these exceptions should translate to. By default it's a fairly + * generic 'FATAL_EXECUTION'. However, this may be insufficient for many + * scenarios. If an enterprise scheduler is used to kick off a batch job, the + * exit code is the only means of communication as to what action must be taken. + * It may also be the only result that many batch operators see as well. + * Therefore, an {@link ExitStatusExceptionClassifier} may be used to classify + * an exception to a particular exit code. + *

* * @author Dave Syer * @author Lucas Ward * @author Ben Hale */ -public class ChunkedStep extends StepSupport implements InitializingBean{ +public class ChunkedStep extends StepSupport implements InitializingBean { private static final Log logger = LogFactory.getLog(ChunkedStep.class); @@ -105,60 +126,66 @@ public class ChunkedStep extends StepSupport implements InitializingBean{ private JobRepository jobRepository; - //default to simple exception classification. + // default to simple exception classification. private ExitStatusExceptionClassifier exceptionClassifier = new SimpleExitStatusExceptionClassifier(); // default to checking current thread for interruption. private StepInterruptionPolicy interruptionPolicy = new ThreadStepInterruptionPolicy(); - + private SkippedItemHandler failureLog = new DefaultItemFailureLog(); private StreamManager streamManager; private ItemReader itemReader; + private Chunker chunker; private ItemWriter itemWriter; + private Dechunker dechunker; - + private ItemSkipPolicy itemSkipPolicy; private RetryTemplate retryTemplate = new RetryTemplate(); - + private int chunkSize; - + public void setChunkSize(int chunkSize) { this.chunkSize = chunkSize; } - /** - * Public setter for the {@link StreamManager}. This will be used to create the {@link StepContext}, and hence any - * component that is a {@link ItemStream} and in step scope will be registered with the service. The - * {@link StepContext} is then a source of aggregate statistics for the step. + * Public setter for the {@link StreamManager}. This will be used to create + * the {@link StepContext}, and hence any component that is a + * {@link ItemStream} and in step scope will be registered with the service. + * The {@link StepContext} is then a source of aggregate statistics for the + * step. * - * @param streamManager the {@link StreamManager} to set. Default is a {@link SimpleStreamManager}. + * @param streamManager the {@link StreamManager} to set. Default is a + * {@link SimpleStreamManager}. */ public void setStreamManager(StreamManager streamManager) { this.streamManager = streamManager; } - + public void setFailureLog(SkippedItemHandler failureLog) { this.failureLog = failureLog; } /** - * Injected strategy for storage and retrieval of persistent step information. Mandatory property. + * Injected strategy for storage and retrieval of persistent step + * information. Mandatory property. * * @param jobRepository */ public void setJobRepository(JobRepository jobRepository) { this.jobRepository = jobRepository; } - + /** - * The {@link RepeatOperations} to use for the outer loop of the batch processing. Should be set up by the caller - * through a factory. Defaults to a plain {@link RepeatTemplate}. + * The {@link RepeatOperations} to use for the outer loop of the batch + * processing. Should be set up by the caller through a factory. Defaults to + * a plain {@link RepeatTemplate}. * * @param stepOperations a {@link RepeatOperations} instance. */ @@ -167,8 +194,9 @@ public class ChunkedStep extends StepSupport implements InitializingBean{ } /** - * Setter for the {@link StepInterruptionPolicy}. The policy is used to check whether an external request has been - * made to interrupt the job execution. + * Setter for the {@link StepInterruptionPolicy}. The policy is used to + * check whether an external request has been made to interrupt the job + * execution. * * @param interruptionPolicy a {@link StepInterruptionPolicy} */ @@ -177,8 +205,8 @@ public class ChunkedStep extends StepSupport implements InitializingBean{ } /** - * Setter for the {@link ExitStatusExceptionClassifier} that will be used to classify any exception that causes a job - * to fail. + * Setter for the {@link ExitStatusExceptionClassifier} that will be used to + * classify any exception that causes a job to fail. * * @param exceptionClassifier */ @@ -199,18 +227,18 @@ public class ChunkedStep extends StepSupport implements InitializingBean{ public void setItemWriter(ItemWriter itemWriter) { this.itemWriter = itemWriter; } - + public void setChunker(Chunker chunker) { this.chunker = chunker; } - + public void setDechunker(Dechunker dechunker) { this.dechunker = dechunker; } - + /** - * Set the skip policy. If set, it will be used for both reading - * and writing. + * Set the skip policy. If set, it will be used for both reading and + * writing. * * @param itemSkipPolicy */ @@ -224,35 +252,42 @@ public class ChunkedStep extends StepSupport implements InitializingBean{ * @see org.springframework.beans.factory.InitializingBean#afterPropertiesSet() */ public void afterPropertiesSet() throws Exception { - //This is currently a little bit funky, I don't want to require a chunker or - //dechunker to be wired in, since the developer should really only be wiring in a - //ItemReader and ItemWriter, a namespace should take care of the issue though. - if(chunker == null){ + // This is currently a little bit funky, I don't want to require a + // chunker or + // dechunker to be wired in, since the developer should really only be + // wiring in a + // ItemReader and ItemWriter, a namespace should take care of the issue + // though. + if (chunker == null) { chunker = new ItemChunker(itemReader); - if(itemSkipPolicy != null){ - ((ItemChunker)chunker).setItemSkipPolicy(itemSkipPolicy); + if (itemSkipPolicy != null) { + ((ItemChunker) chunker).setItemSkipPolicy(itemSkipPolicy); } } - - if(dechunker == null){ + + if (dechunker == null) { dechunker = new ItemDechunker(itemWriter); - if(itemSkipPolicy != null){ - ((ItemChunker)dechunker).setItemSkipPolicy(itemSkipPolicy); + if (itemSkipPolicy != null) { + ((ItemChunker) dechunker).setItemSkipPolicy(itemSkipPolicy); } } - + Assert.notNull(jobRepository, "JobRepository must not be null"); } /** - * Process the step and update its context so that progress can be monitored by the caller. The step is broken down - * into chunks, each one executing in a transaction. The step and its execution and execution context are all given - * an up to date {@link BatchStatus}, and the {@link JobRepository} is used to store the result. Various reporting - * information are also added to the current context (the {@link RepeatContext} governing the step execution, which - * would normally be available to the caller somehow through the step's {@link StepContext}.
+ * Process the step and update its context so that progress can be monitored + * by the caller. The step is broken down into chunks, each one executing in + * a transaction. The step and its execution and execution context are all + * given an up to date {@link BatchStatus}, and the {@link JobRepository} + * is used to store the result. Various reporting information are also added + * to the current context (the {@link RepeatContext} governing the step + * execution, which would normally be available to the caller somehow + * through the step's {@link StepContext}.
* * @throws JobInterruptedException if the step or a chunk is interrupted - * @throws RuntimeException if there is an exception during a chunk execution + * @throws RuntimeException if there is an exception during a chunk + * execution * @see StepExecutor#execute(StepExecution) */ public void execute(final StepExecution stepExecution) throws BatchCriticalException, JobInterruptedException { @@ -267,15 +302,18 @@ public class ChunkedStep extends StepSupport implements InitializingBean{ StepContext parentStepContext = StepSynchronizationManager.getContext(); final StepContext stepContext = new SimpleStepContext(stepExecution, parentStepContext, streamManager); StepSynchronizationManager.register(stepContext); + possiblyRegisterStreams(stepExecution); // Add the job identifier so that it can be used to identify // the conversation in StepScope stepContext.setAttribute(StepScope.ID_KEY, stepExecution.getJobExecution().getId()); final boolean saveExecutionContext = isSaveExecutionContext(); + streamManager.open(stepExecution); + if (saveExecutionContext && isRestart && stepInstance.getLastExecution() != null) { stepExecution.setExecutionContext(stepInstance.getLastExecution().getExecutionContext()); - stepContext.restoreFrom(stepExecution.getExecutionContext()); + streamManager.restoreFrom(stepExecution, stepExecution.getExecutionContext()); } try { @@ -288,27 +326,26 @@ public class ChunkedStep extends StepSupport implements InitializingBean{ public ExitStatus doInIteration(final RepeatContext context) throws Exception { - // Before starting a new transaction, check for // interruption. interruptionPolicy.checkInterrupted(context); - + ChunkingResult chunkingResult = chunker.chunk(chunkSize, stepExecution); - - if(chunkingResult == null){ + + if (chunkingResult == null) { return ExitStatus.FINISHED; } - + final Chunk chunk = chunkingResult.getChunk(); failureLog.handle(chunkingResult.getExceptions()); - retryTemplate.execute(new RetryCallback(){ + retryTemplate.execute(new RetryCallback() { - public Object doWithRetry(RetryContext context) - throws Throwable { + public Object doWithRetry(RetryContext context) throws Throwable { processChunk(chunk, stepExecution, stepContext); return null; - }}); + } + }); // Check for interruption after transaction as well, so that // the interrupted exception is correctly propagated up to @@ -321,57 +358,79 @@ public class ChunkedStep extends StepSupport implements InitializingBean{ }); updateStatus(stepExecution, BatchStatus.COMPLETED); - } catch (RuntimeException e) { + } + catch (RuntimeException e) { // classify exception so an exit code can be stored. status = exceptionClassifier.classifyForExitCode(e); if (e.getCause() instanceof JobInterruptedException) { updateStatus(stepExecution, BatchStatus.STOPPED); throw (JobInterruptedException) e.getCause(); - } else if (e instanceof ResetFailedException) { + } + else if (e instanceof ResetFailedException) { updateStatus(stepExecution, BatchStatus.UNKNOWN); throw (ResetFailedException) e; - } else { + } + else { updateStatus(stepExecution, BatchStatus.FAILED); throw e; } - } finally { + } + finally { stepExecution.setExitStatus(status); stepExecution.setEndTime(new Date(System.currentTimeMillis())); try { jobRepository.saveOrUpdate(stepExecution); - } finally { + } + finally { // clear any registered synchronizations StepSynchronizationManager.close(); + streamManager.close(stepExecution); } } } /** - * Execute a bunch of identical business logic operations all within a transaction. * - * @param stepExecution the current execution in which to process the chunk in. + */ + private void possiblyRegisterStreams(Object key) { + if (itemReader instanceof ItemStream) { + ItemStream stream = (ItemStream) itemReader; + streamManager.register(key, stream); + } + if (itemWriter instanceof ItemStream) { + ItemStream stream = (ItemStream) itemWriter; + streamManager.register(key, stream); + } + } + + /** + * Execute a bunch of identical business logic operations all within a + * transaction. + * + * @param stepExecution the current execution in which to process the chunk + * in. * @param chunk to be processed. * @param stepContext the current step context. * @return true if there is more data to process. */ void processChunk(Chunk chunk, final StepExecution stepExecution, StepContext stepContext) { - - TransactionStatus transaction = streamManager.getTransaction(stepContext); + + TransactionStatus transaction = streamManager.getTransaction(stepExecution); final StepContribution contribution = stepExecution.createStepContribution(); - + try { - + DechunkingResult chunkResult = dechunker.dechunk(chunk, stepExecution); failureLog.handle(chunkResult.getExceptions()); // TODO: check that stepExecution can // aggregate these contributions if they // come in asynchronously. - ExecutionContext statistics = stepContext.getExecutionContext(); + ExecutionContext statistics = streamManager.getExecutionContext(stepExecution); contribution.setExecutionContext(statistics); contribution.incrementCommitCount(); @@ -385,7 +444,7 @@ public class ChunkedStep extends StepSupport implements InitializingBean{ stepExecution.apply(contribution); if (isSaveExecutionContext()) { - stepExecution.setExecutionContext(stepContext.getExecutionContext()); + stepExecution.setExecutionContext(statistics); } jobRepository.saveOrUpdate(stepExecution); @@ -393,43 +452,47 @@ public class ChunkedStep extends StepSupport implements InitializingBean{ streamManager.commit(transaction); - } catch (Throwable t) { + } + catch (Throwable t) { /* - * Any exception thrown within the transaction template will automatically cause the transaction - * to rollback. We need to include exceptions during an attempted commit (e.g. Hibernate flush) - * so this catch block comes outside the transaction. + * Any exception thrown within the transaction template will + * automatically cause the transaction to rollback. We need to + * include exceptions during an attempted commit (e.g. Hibernate + * flush) so this catch block comes outside the transaction. */ synchronized (stepExecution) { stepExecution.rollback(); } try { streamManager.rollback(transaction); - } catch (ResetFailedException e) { + } + catch (ResetFailedException e) { // The original Throwable cause is in danger of // being lost here, so we log the reset // failure and re-throw with cause of the rollback. - logger.error("Encountered reset error on rollback: " - + "one of the streams may be in an inconsistent state, " - + "so this step should not proceed", e); + logger + .error("Encountered reset error on rollback: " + + "one of the streams may be in an inconsistent state, " + + "so this step should not proceed", e); throw new ResetFailedException("Encountered reset error on rollback. " - + "Consult logs for the cause of the reet failure. " - + "The cause of the original rollback is incuded here.", t); + + "Consult logs for the cause of the reet failure. " + + "The cause of the original rollback is incuded here.", t); } if (t instanceof RuntimeException) { throw (RuntimeException) t; - } else { + } + else { throw new RuntimeException(t); } } - + } /* * Convenience method to update the status in all relevant places. * - * @param stepInstance the current step - * @param stepExecution the current stepExecution - * @param status the status to set + * @param stepInstance the current step @param stepExecution the current + * stepExecution @param status the status to set */ private void updateStatus(StepExecution stepExecution, BatchStatus status) { stepExecution.setStatus(status); diff --git a/spring-batch-execution/src/main/java/org/springframework/batch/execution/step/simple/SimpleStepExecutor.java b/spring-batch-execution/src/main/java/org/springframework/batch/execution/step/simple/SimpleStepExecutor.java index 372a7d990..4d608cf0c 100644 --- a/spring-batch-execution/src/main/java/org/springframework/batch/execution/step/simple/SimpleStepExecutor.java +++ b/spring-batch-execution/src/main/java/org/springframework/batch/execution/step/simple/SimpleStepExecutor.java @@ -20,11 +20,11 @@ import java.util.Date; import org.apache.commons.logging.Log; import org.apache.commons.logging.LogFactory; import org.springframework.batch.core.domain.BatchStatus; +import org.springframework.batch.core.domain.JobInterruptedException; import org.springframework.batch.core.domain.Step; import org.springframework.batch.core.domain.StepContribution; import org.springframework.batch.core.domain.StepExecution; import org.springframework.batch.core.domain.StepInstance; -import org.springframework.batch.core.domain.JobInterruptedException; import org.springframework.batch.core.repository.JobRepository; import org.springframework.batch.core.runtime.ExitStatusExceptionClassifier; import org.springframework.batch.core.tasklet.Tasklet; @@ -294,27 +294,32 @@ public class SimpleStepExecutor implements InitializingBean { boolean isRestart = stepInstance.getStepExecutionCount() > 0 ? true : false; ExitStatus status = ExitStatus.FAILED; - - StepContext parentStepContext = StepSynchronizationManager.getContext(); - final StepContext stepContext = new SimpleStepContext(stepExecution, parentStepContext, streamManager); - StepSynchronizationManager.register(stepContext); - // Add the job identifier so that it can be used to identify - // the conversation in StepScope - stepContext.setAttribute(StepScope.ID_KEY, stepExecution.getJobExecution().getId()); - - final boolean saveExecutionContext = step.isSaveExecutionContext(); - - if (saveExecutionContext && isRestart && stepInstance.getLastExecution() != null) { - stepExecution.setExecutionContext(stepInstance.getLastExecution().getExecutionContext()); - stepContext.restoreFrom(stepExecution.getExecutionContext()); - } - + try { stepExecution.setStartTime(new Date(System.currentTimeMillis())); - stepInstance.setLastExecution(stepExecution); + // We need to save the step execution right away, before we start + // using its ID. It would be better to make the creation atomic in + // the caller. updateStatus(stepExecution, BatchStatus.STARTED); + StepContext parentStepContext = StepSynchronizationManager.getContext(); + final StepContext stepContext = new SimpleStepContext(stepExecution, parentStepContext, streamManager); + StepSynchronizationManager.register(stepContext); + possiblyRegisterStreams(stepExecution); + // Add the job identifier so that it can be used to identify + // the conversation in StepScope + stepContext.setAttribute(StepScope.ID_KEY, stepExecution.getJobExecution().getId()); + + final boolean saveExecutionContext = step.isSaveExecutionContext(); + + streamManager.open(stepExecution); + + if (saveExecutionContext && isRestart && stepInstance.getLastExecution() != null) { + stepExecution.setExecutionContext(stepInstance.getLastExecution().getExecutionContext()); + streamManager.restoreFrom(stepExecution, stepExecution.getExecutionContext()); + } + status = stepOperations.iterate(new RepeatCallback() { public ExitStatus doInIteration(final RepeatContext context) throws Exception { @@ -327,7 +332,7 @@ public class SimpleStepExecutor implements InitializingBean { ExitStatus result; - TransactionStatus transaction = streamManager.getTransaction(stepContext); + TransactionStatus transaction = streamManager.getTransaction(stepExecution); try { @@ -336,7 +341,7 @@ public class SimpleStepExecutor implements InitializingBean { // TODO: check that stepExecution can // aggregate these contributions if they // come in asynchronously. - ExecutionContext statistics = stepContext.getExecutionContext(); + ExecutionContext statistics = streamManager.getExecutionContext(stepExecution); contribution.setExecutionContext(statistics); contribution.incrementCommitCount(); @@ -350,7 +355,7 @@ public class SimpleStepExecutor implements InitializingBean { stepExecution.apply(contribution); if (saveExecutionContext) { - stepExecution.setExecutionContext(stepContext.getExecutionContext()); + stepExecution.setExecutionContext(statistics); } jobRepository.saveOrUpdate(stepExecution); @@ -427,15 +432,31 @@ public class SimpleStepExecutor implements InitializingBean { stepExecution.setEndTime(new Date(System.currentTimeMillis())); try { jobRepository.saveOrUpdate(stepExecution); + streamManager.close(stepExecution); } finally { // clear any registered synchronizations + StepSynchronizationManager.close(); } } } + /** + * + */ + private void possiblyRegisterStreams(Object key) { + if (itemReader instanceof ItemStream) { + ItemStream stream = (ItemStream) itemReader; + streamManager.register(key, stream); + } + if (itemWriter instanceof ItemStream) { + ItemStream stream = (ItemStream) itemWriter; + streamManager.register(key, stream); + } + } + /** * Execute a bunch of identical business logic operations all within a * transaction. The transaction is programmatically started and stopped diff --git a/spring-batch-execution/src/main/resources/schema-db2.sql b/spring-batch-execution/src/main/resources/schema-db2.sql index db7fe1745..a99cb490b 100644 --- a/spring-batch-execution/src/main/resources/schema-db2.sql +++ b/spring-batch-execution/src/main/resources/schema-db2.sql @@ -3,7 +3,7 @@ DROP TABLE BATCH_STEP_EXECUTION ; DROP TABLE BATCH_JOB_EXECUTION ; DROP TABLE BATCH_STEP_INSTANCE ; DROP TABLE BATCH_JOB_INSTANCE ; -DROP TABLE BATCH_JOB_INSTANCE_PARAMS ; +DROP TABLE BATCH_JOB_PARAMS ; DROP TABLE BATCH_STEP_EXECUTION_ATTRS ; DROP SEQUENCE BATCH_STEP_EXECUTION_SEQ ; diff --git a/spring-batch-execution/src/main/resources/schema-derby.sql b/spring-batch-execution/src/main/resources/schema-derby.sql index 1171d4358..6a1624e66 100644 --- a/spring-batch-execution/src/main/resources/schema-derby.sql +++ b/spring-batch-execution/src/main/resources/schema-derby.sql @@ -3,7 +3,7 @@ DROP TABLE BATCH_STEP_EXECUTION ; DROP TABLE BATCH_JOB_EXECUTION ; DROP TABLE BATCH_STEP_INSTANCE ; DROP TABLE BATCH_JOB_INSTANCE ; -DROP TABLE BATCH_JOB_INSTANCE_PARAMS ; +DROP TABLE BATCH_JOB_PARAMS ; DROP TABLE BATCH_STEP_EXECUTION_ATTRS ; DROP TABLE BATCH_STEP_EXECUTION_SEQ ; diff --git a/spring-batch-execution/src/main/resources/schema-hsqldb.sql b/spring-batch-execution/src/main/resources/schema-hsqldb.sql index 269abc738..04f595552 100644 --- a/spring-batch-execution/src/main/resources/schema-hsqldb.sql +++ b/spring-batch-execution/src/main/resources/schema-hsqldb.sql @@ -3,7 +3,7 @@ DROP TABLE BATCH_STEP_EXECUTION IF EXISTS; DROP TABLE BATCH_JOB_EXECUTION IF EXISTS; DROP TABLE BATCH_STEP_INSTANCE IF EXISTS; DROP TABLE BATCH_JOB_INSTANCE IF EXISTS; -DROP TABLE BATCH_JOB_INSTANCE_PARAMS IF EXISTS; +DROP TABLE BATCH_JOB_PARAMS IF EXISTS; DROP TABLE BATCH_STEP_EXECUTION_ATTRS IF EXISTS; DROP TABLE BATCH_STEP_EXECUTION_SEQ IF EXISTS; diff --git a/spring-batch-execution/src/main/resources/schema-mysql.sql b/spring-batch-execution/src/main/resources/schema-mysql.sql index 15fa7c74a..e494c6f6b 100644 --- a/spring-batch-execution/src/main/resources/schema-mysql.sql +++ b/spring-batch-execution/src/main/resources/schema-mysql.sql @@ -3,7 +3,7 @@ DROP TABLE IF EXISTS BATCH_STEP_EXECUTION ; DROP TABLE IF EXISTS BATCH_JOB_EXECUTION ; DROP TABLE IF EXISTS BATCH_STEP_INSTANCE ; DROP TABLE IF EXISTS BATCH_JOB_INSTANCE ; -DROP TABLE IF EXISTS BATCH_JOB_INSTANCE_PARAMS ; +DROP TABLE IF EXISTS BATCH_JOB_PARAMS ; DROP TABLE IF EXISTS BATCH_STEP_EXECUTION_ATTRS ; DROP TABLE IF EXISTS BATCH_STEP_EXECUTION_SEQ ; diff --git a/spring-batch-execution/src/main/resources/schema-oracle10g.sql b/spring-batch-execution/src/main/resources/schema-oracle10g.sql index 1a6e0a481..abf1c9a59 100644 --- a/spring-batch-execution/src/main/resources/schema-oracle10g.sql +++ b/spring-batch-execution/src/main/resources/schema-oracle10g.sql @@ -3,7 +3,7 @@ DROP TABLE BATCH_STEP_EXECUTION ; DROP TABLE BATCH_JOB_EXECUTION ; DROP TABLE BATCH_STEP_INSTANCE ; DROP TABLE BATCH_JOB_INSTANCE ; -DROP TABLE BATCH_JOB_INSTANCE_PARAMS ; +DROP TABLE BATCH_JOB_PARAMS ; DROP TABLE BATCH_STEP_EXECUTION_ATTRS ; DROP SEQUENCE BATCH_STEP_EXECUTION_SEQ ; diff --git a/spring-batch-execution/src/main/resources/schema-postgresql.sql b/spring-batch-execution/src/main/resources/schema-postgresql.sql index db7fe1745..a99cb490b 100644 --- a/spring-batch-execution/src/main/resources/schema-postgresql.sql +++ b/spring-batch-execution/src/main/resources/schema-postgresql.sql @@ -3,7 +3,7 @@ DROP TABLE BATCH_STEP_EXECUTION ; DROP TABLE BATCH_JOB_EXECUTION ; DROP TABLE BATCH_STEP_INSTANCE ; DROP TABLE BATCH_JOB_INSTANCE ; -DROP TABLE BATCH_JOB_INSTANCE_PARAMS ; +DROP TABLE BATCH_JOB_PARAMS ; DROP TABLE BATCH_STEP_EXECUTION_ATTRS ; DROP SEQUENCE BATCH_STEP_EXECUTION_SEQ ; diff --git a/spring-batch-execution/src/main/sql/destroy.sql.vpp b/spring-batch-execution/src/main/sql/destroy.sql.vpp index 2bb836b8e..ab9453646 100644 --- a/spring-batch-execution/src/main/sql/destroy.sql.vpp +++ b/spring-batch-execution/src/main/sql/destroy.sql.vpp @@ -3,7 +3,7 @@ DROP TABLE $!{IFEXISTSBEFORE} BATCH_STEP_EXECUTION $!{IFEXISTS}; DROP TABLE $!{IFEXISTSBEFORE} BATCH_JOB_EXECUTION $!{IFEXISTS}; DROP TABLE $!{IFEXISTSBEFORE} BATCH_STEP_INSTANCE $!{IFEXISTS}; DROP TABLE $!{IFEXISTSBEFORE} BATCH_JOB_INSTANCE $!{IFEXISTS}; -DROP TABLE $!{IFEXISTSBEFORE} BATCH_JOB_INSTANCE_PARAMS $!{IFEXISTS}; +DROP TABLE $!{IFEXISTSBEFORE} BATCH_JOB_PARAMS $!{IFEXISTS}; DROP TABLE $!{IFEXISTSBEFORE} BATCH_STEP_EXECUTION_ATTRS $!{IFEXISTS}; DROP ${SEQUENCE} $!{IFEXISTSBEFORE} BATCH_STEP_EXECUTION_SEQ $!{IFEXISTS}; diff --git a/spring-batch-execution/src/test/java/org/springframework/batch/execution/scope/SimpleStepContextTests.java b/spring-batch-execution/src/test/java/org/springframework/batch/execution/scope/SimpleStepContextTests.java index 389507c31..83a228092 100644 --- a/spring-batch-execution/src/test/java/org/springframework/batch/execution/scope/SimpleStepContextTests.java +++ b/spring-batch-execution/src/test/java/org/springframework/batch/execution/scope/SimpleStepContextTests.java @@ -16,18 +16,11 @@ package org.springframework.batch.execution.scope; import java.util.ArrayList; -import java.util.HashMap; import java.util.List; -import java.util.Map; import junit.framework.TestCase; import org.springframework.batch.core.domain.StepExecution; -import org.springframework.batch.item.ItemStream; -import org.springframework.batch.item.ExecutionContext; -import org.springframework.batch.item.stream.ItemStreamAdapter; -import org.springframework.batch.item.stream.SimpleStreamManager; -import org.springframework.batch.support.PropertiesConverter; /** * @author Dave Syer @@ -133,50 +126,4 @@ public class SimpleStepContextTests extends TestCase { assertTrue(list.contains("spam")); } - public void testExecutionContextWithNotNullService() throws Exception { - Map map = new HashMap(); - context = new SimpleStepContext(null, null, new StubStreamManager(map)); - assertEquals(1, context.getExecutionContext().getProperties().size()); - assertEquals("bar", context.getExecutionContext().getProperties().getProperty("foo")); - } - - public void testStreamManagerRegistration() throws Exception { - Map map = new HashMap(); - context = new SimpleStepContext(null, null, new StubStreamManager(map)); - ItemStreamAdapter provider = new ItemStreamAdapter(); - context.setAttribute("foo", provider); - assertEquals(1, map.size()); - assertEquals(context, map.keySet().iterator().next()); - assertEquals(provider, map.values().iterator().next()); - } - - /** - * @author Dave Syer - * - */ - private class StubStreamManager extends SimpleStreamManager { - private final Map map; - - private StubStreamManager(Map map) { - this.map = map; - } - - public void close(Object key) { - } - - public ExecutionContext getExecutionContext(Object key) { - return new ExecutionContext(PropertiesConverter.stringToProperties("foo=bar")); - } - - public void open(Object key) { - } - - public void register(Object key, ItemStream stream, ExecutionContext executionContext) { - map.put(key, stream); - } - - public void restoreFrom(Object key, ExecutionContext data) { - } - } - } diff --git a/spring-batch-execution/src/test/java/org/springframework/batch/execution/step/simple/SimpleStepExecutorTests.java b/spring-batch-execution/src/test/java/org/springframework/batch/execution/step/simple/SimpleStepExecutorTests.java index 20af14203..dac23b7cd 100644 --- a/spring-batch-execution/src/test/java/org/springframework/batch/execution/step/simple/SimpleStepExecutorTests.java +++ b/spring-batch-execution/src/test/java/org/springframework/batch/execution/step/simple/SimpleStepExecutorTests.java @@ -281,7 +281,7 @@ public class SimpleStepExecutorTests extends TestCase { */ public void testNonRestartedJob() throws Exception { StepInstance step = new StepInstance(new Long(1)); - MockRestartableTasklet tasklet = new MockRestartableTasklet(); + MockRestartableItemReader tasklet = new MockRestartableItemReader(); stepExecutor.setItemReader(tasklet); stepConfiguration.setSaveExecutionContext(true); JobExecution jobExecutionContext = new JobExecution(jobInstance); @@ -300,7 +300,7 @@ public class SimpleStepExecutorTests extends TestCase { public void testRestartedJob() throws Exception { StepInstance step = new StepInstance(new Long(1)); step.setStepExecutionCount(1); - MockRestartableTasklet tasklet = new MockRestartableTasklet(); + MockRestartableItemReader tasklet = new MockRestartableItemReader(); stepExecutor.setItemReader(tasklet); stepConfiguration.setSaveExecutionContext(true); JobExecution jobExecutionContext = new JobExecution(jobInstance); @@ -324,7 +324,7 @@ public class SimpleStepExecutorTests extends TestCase { public void testNoSaveExecutionAttributesRestartableJob() { StepInstance step = new StepInstance(new Long(1)); step.setStepExecutionCount(1); - MockRestartableTasklet tasklet = new MockRestartableTasklet(); + MockRestartableItemReader tasklet = new MockRestartableItemReader(); stepConfiguration.setItemReader(tasklet); stepConfiguration.setSaveExecutionContext(false); JobExecution jobExecutionContext = new JobExecution(jobInstance); @@ -431,7 +431,7 @@ public class SimpleStepExecutorTests extends TestCase { assertEquals(0, map.size()); } - private class MockRestartableTasklet extends ItemStreamAdapter implements ItemReader { + private class MockRestartableItemReader extends ItemStreamAdapter implements ItemReader { private boolean getExecutionAttributesCalled = false; diff --git a/spring-batch-execution/src/test/resources/org/springframework/batch/execution/repository/dao/destroy.sql b/spring-batch-execution/src/test/resources/org/springframework/batch/execution/repository/dao/destroy.sql index 68c2ff272..0974a9770 100644 --- a/spring-batch-execution/src/test/resources/org/springframework/batch/execution/repository/dao/destroy.sql +++ b/spring-batch-execution/src/test/resources/org/springframework/batch/execution/repository/dao/destroy.sql @@ -3,7 +3,7 @@ DROP TABLE BATCH_STEP_EXECUTION IF EXISTS; DROP TABLE BATCH_JOB_EXECUTION IF EXISTS; DROP TABLE BATCH_STEP_INSTANCE IF EXISTS; DROP TABLE BATCH_JOB_INSTANCE IF EXISTS; -DROP TABLE BATCH_JOB_INSTANCE_PARAMS IF EXISTS; +DROP TABLE BATCH_JOB_PARAMS IF EXISTS; DROP TABLE BATCH_STEP_EXECUTION_ATTRS IF EXISTS; DROP TABLE BATCH_STEP_EXECUTION_SEQ IF EXISTS; diff --git a/spring-batch-infrastructure/src/main/java/org/springframework/batch/item/reader/AggregateItemReader.java b/spring-batch-infrastructure/src/main/java/org/springframework/batch/item/reader/AggregateItemReader.java index 3869c871d..edeaa6d25 100644 --- a/spring-batch-infrastructure/src/main/java/org/springframework/batch/item/reader/AggregateItemReader.java +++ b/spring-batch-infrastructure/src/main/java/org/springframework/batch/item/reader/AggregateItemReader.java @@ -37,13 +37,11 @@ import org.springframework.batch.item.ItemReader; * @author Dave Syer * */ -public class AggregateItemReader extends AbstractItemReader { +public class AggregateItemReader extends DelegatingItemReader { private static final Log log = LogFactory .getLog(AggregateItemReader.class); - private ItemReader inputSource; - /** * Marker for the end of a multi-object record. */ @@ -63,7 +61,7 @@ public class AggregateItemReader extends AbstractItemReader { public Object read() throws Exception { ResultHolder holder = new ResultHolder(); - while (process(inputSource.read(), holder)) { + while (process(getItemReader().read(), holder)) { continue; } @@ -100,16 +98,6 @@ public class AggregateItemReader extends AbstractItemReader { return true; } - /** - * Injection setter for {@link ItemReader}. - * - * @param inputSource - * an {@link ItemReader}. - */ - public void setItemReader(ItemReader inputSource) { - this.inputSource = inputSource; - } - /** * Private class for temporary state management while item is being * collected. diff --git a/spring-batch-infrastructure/src/main/java/org/springframework/batch/item/reader/DelegatingItemReader.java b/spring-batch-infrastructure/src/main/java/org/springframework/batch/item/reader/DelegatingItemReader.java index 4d13c4818..3d59471ad 100644 --- a/spring-batch-infrastructure/src/main/java/org/springframework/batch/item/reader/DelegatingItemReader.java +++ b/spring-batch-infrastructure/src/main/java/org/springframework/batch/item/reader/DelegatingItemReader.java @@ -28,7 +28,7 @@ import org.springframework.util.Assert; * Simple wrapper around {@link ItemReader}. The input source is expected to * take care of open and close operations. If necessary it should be registered * as a step scoped bean to ensure that the lifecycle methods are called. - * + * * @author Dave Syer */ public class DelegatingItemReader extends AbstractItemReader implements Skippable, InitializingBean, ItemStream { @@ -38,9 +38,10 @@ public class DelegatingItemReader extends AbstractItemReader implements Skippabl public void afterPropertiesSet() throws Exception { Assert.notNull(inputSource, "ItemReader must not be null."); } + /** * Get the next object from the input source. - * @throws Exception + * @throws Exception * @see org.springframework.batch.item.ItemReader#read() */ public Object read() throws Exception { @@ -53,9 +54,10 @@ public class DelegatingItemReader extends AbstractItemReader implements Skippabl * {@link ItemStream}. */ public ExecutionContext getExecutionContext() { - // TODO: this is not necessary... - Assert.state(inputSource instanceof ItemStream, "Input source is not ItemStream"); - return ((ItemStream) inputSource).getExecutionContext(); + if (inputSource instanceof ItemStream) { + return ((ItemStream) inputSource).getExecutionContext(); + } + return new ExecutionContext(); } /** @@ -64,8 +66,9 @@ public class DelegatingItemReader extends AbstractItemReader implements Skippabl * {@link ItemStream}. */ public void restoreFrom(ExecutionContext data) { - Assert.state(inputSource instanceof ItemStream, "Input source is not ItemStream"); - ((ItemStream) inputSource).restoreFrom(data); + if (inputSource instanceof ItemStream) { + ((ItemStream) inputSource).restoreFrom(data); + } } /** @@ -82,11 +85,12 @@ public class DelegatingItemReader extends AbstractItemReader implements Skippabl public void skip() { if (inputSource instanceof Skippable) { - ((Skippable)inputSource).skip(); + ((Skippable) inputSource).skip(); } } - - /* (non-Javadoc) + + /* + * (non-Javadoc) * @see org.springframework.batch.item.ItemStream#open() */ public void open() throws StreamException { @@ -95,7 +99,8 @@ public class DelegatingItemReader extends AbstractItemReader implements Skippabl } } - /* (non-Javadoc) + /* + * (non-Javadoc) * @see org.springframework.batch.item.ItemStream#open() */ public void close() throws StreamException { @@ -116,7 +121,8 @@ public class DelegatingItemReader extends AbstractItemReader implements Skippabl return false; } - /* (non-Javadoc) + /* + * (non-Javadoc) * @see org.springframework.batch.item.ItemStream#mark(org.springframework.batch.item.ExecutionContext) */ public void mark() { @@ -125,7 +131,8 @@ public class DelegatingItemReader extends AbstractItemReader implements Skippabl } } - /* (non-Javadoc) + /* + * (non-Javadoc) * @see org.springframework.batch.item.ItemStream#reset(org.springframework.batch.item.ExecutionContext) */ public void reset() { diff --git a/spring-batch-infrastructure/src/main/java/org/springframework/batch/item/stream/SimpleStreamManager.java b/spring-batch-infrastructure/src/main/java/org/springframework/batch/item/stream/SimpleStreamManager.java index 9937cdb9b..d115714dc 100644 --- a/spring-batch-infrastructure/src/main/java/org/springframework/batch/item/stream/SimpleStreamManager.java +++ b/spring-batch-infrastructure/src/main/java/org/springframework/batch/item/stream/SimpleStreamManager.java @@ -116,7 +116,7 @@ public class SimpleStreamManager implements StreamManager { * @see org.springframework.batch.item.stream.StreamManager#register(java.lang.Object, * org.springframework.batch.item.ItemStream, ExecutionContext) */ - public void register(Object key, ItemStream stream, ExecutionContext executionContext) { + public void register(Object key, ItemStream stream) { synchronized (registry) { Set set = (Set) registry.get(key); if (set == null) { @@ -125,9 +125,19 @@ public class SimpleStreamManager implements StreamManager { } set.add(stream); } + } + + /* (non-Javadoc) + * @see org.springframework.batch.item.stream.StreamManager#restoreFrom(java.lang.Object, org.springframework.batch.item.ExecutionAttributes) + */ + public void restoreFrom(Object key, final ExecutionContext executionContext) { if (executionContext != null) { - stream.restoreFrom(extract(stream, executionContext)); - } + iterate(key, new Callback() { + public void execute(ItemStream stream) { + stream.restoreFrom(extract(stream, executionContext)); + } + }); + } } /** @@ -153,7 +163,7 @@ public class SimpleStreamManager implements StreamManager { /** * Broadcast the call to close from this {@link StreamManager}. - * @throws Exception + * @throws StreamException */ public void close(Object key) throws StreamException { iterate(key, new Callback() { @@ -163,6 +173,18 @@ public class SimpleStreamManager implements StreamManager { }); } + /** + * Broadcast the call to open from this {@link StreamManager}. + * @throws StreamException + */ + public void open(Object key) throws StreamException { + iterate(key, new Callback() { + public void execute(ItemStream stream) { + stream.open(); + } + }); + } + /** * Delegate to the {@link PlatformTransactionManager} to create a new * transaction. diff --git a/spring-batch-infrastructure/src/main/java/org/springframework/batch/item/stream/StreamManager.java b/spring-batch-infrastructure/src/main/java/org/springframework/batch/item/stream/StreamManager.java index 30b787c09..6665ebd1b 100644 --- a/spring-batch-infrastructure/src/main/java/org/springframework/batch/item/stream/StreamManager.java +++ b/spring-batch-infrastructure/src/main/java/org/springframework/batch/item/stream/StreamManager.java @@ -37,9 +37,8 @@ public interface StreamManager { * * @param key the key under which to add the provider * @param stream an {@link ItemStream} - * @param executionContext the context (may be null) to restore from on registration */ - void register(Object key, ItemStream stream, ExecutionContext executionContext); + void register(Object key, ItemStream stream); /** * Extract and aggregate the {@link ExecutionContext} from all streams under @@ -60,6 +59,15 @@ public interface StreamManager { * been registered. */ void close(Object key) throws StreamException; + + /** + * If any resources are needed for the stream to operate they need to be + * opened here. + * + * @param key the key under which {@link ItemStream} instances might have + * been registered. + */ + void open(Object key) throws StreamException; TransactionStatus getTransaction(Object key); @@ -67,4 +75,12 @@ public interface StreamManager { void rollback(TransactionStatus transaction); + /** + * Restore the streams registered under a goven key. + * + * @param key the key under which the streams are stored + * @param executionContext the context (may be null) to restore from + */ + void restoreFrom(Object key, ExecutionContext executionContext); + } diff --git a/spring-batch-infrastructure/src/test/java/org/springframework/batch/item/stream/SimpleStreamManagerTests.java b/spring-batch-infrastructure/src/test/java/org/springframework/batch/item/stream/SimpleStreamManagerTests.java index 1b124edbf..655ab46e8 100644 --- a/spring-batch-infrastructure/src/test/java/org/springframework/batch/item/stream/SimpleStreamManagerTests.java +++ b/spring-batch-infrastructure/src/test/java/org/springframework/batch/item/stream/SimpleStreamManagerTests.java @@ -84,7 +84,7 @@ public class SimpleStreamManagerTests extends TestCase { * {@link org.springframework.batch.item.stream.SimpleStreamManager#getExecutionContext(java.lang.Object)}. */ public void testGetStreamContextNotEmpty() { - manager.register("foo", stream, null); + manager.register("foo", stream); ExecutionContext streamContext = manager.getExecutionContext("foo"); assertEquals(1, streamContext.entrySet().size()); assertEquals("bar", streamContext.getString(ClassUtils.getQualifiedName(stream.getClass()) + ".foo")); @@ -99,7 +99,7 @@ public class SimpleStreamManagerTests extends TestCase { ExecutionContext context = manager.getExecutionContext("foo"); // Register again, now with the context that was created from the same // stream... - manager.register("foo", stream, context); + manager.restoreFrom("foo", context); assertEquals(1, list.size()); // The list should have the foo= map value from the sub-context assertEquals("bar", list.get(0)); @@ -112,7 +112,8 @@ public class SimpleStreamManagerTests extends TestCase { public void testGetStreamContextNotEmptyAndRestoreWithNoPrefix() { ExecutionContext context = new ExecutionContext(PropertiesConverter.stringToProperties("foo=bar")); manager.setUseClassNameAsPrefix(false); - manager.register("foo", stream, context); + manager.register("foo", stream); + manager.restoreFrom("foo", context); assertEquals(1, list.size()); // The list should have the foo= map value from the sub-context assertEquals("bar", list.get(0)); @@ -124,7 +125,7 @@ public class SimpleStreamManagerTests extends TestCase { */ public void testGetStreamContextWithNoPrefix() { manager.setUseClassNameAsPrefix(false); - manager.register("foo", stream, null); + manager.register("foo", stream); ExecutionContext context = manager.getExecutionContext("foo"); assertEquals(1, context.entrySet().size()); // The list should have the foo= map value from the sub-context @@ -140,12 +141,12 @@ public class SimpleStreamManagerTests extends TestCase { public ExecutionContext getExecutionContext() { return new ExecutionContext(PropertiesConverter.stringToProperties("foo=bar")); } - }, null); + }); manager.register("foo", new ItemStreamAdapter() { public ExecutionContext getExecutionContext() { return new ExecutionContext(PropertiesConverter.stringToProperties("foo=spam")); } - }, null); + }); ExecutionContext streamContext = manager.getExecutionContext("foo"); assertEquals(2, streamContext.entrySet().size()); } @@ -160,7 +161,7 @@ public class SimpleStreamManagerTests extends TestCase { list.add("bar"); super.close(); } - }, null); + }); manager.close("foo"); assertEquals(1, list.size()); } @@ -178,7 +179,7 @@ public class SimpleStreamManagerTests extends TestCase { public void mark() { list.add("bar"); } - }, null); + }); TransactionStatus status = manager.getTransaction("foo"); manager.commit(status); assertEquals(1, list.size()); @@ -193,7 +194,7 @@ public class SimpleStreamManagerTests extends TestCase { public void mark() { list.add("bar"); } - }, null); + }); TransactionStatus status = manager.getTransaction("foo"); manager.commit(status); assertEquals(0, list.size()); @@ -212,7 +213,7 @@ public class SimpleStreamManagerTests extends TestCase { public void reset() { list.add("bar"); } - }, null); + }); TransactionStatus status = manager.getTransaction("foo"); manager.rollback(status); assertEquals(1, list.size()); @@ -227,7 +228,7 @@ public class SimpleStreamManagerTests extends TestCase { public void reset() { list.add("bar"); } - }, null); + }); TransactionStatus status = manager.getTransaction("foo"); manager.rollback(status); assertEquals(0, list.size()); diff --git a/spring-batch-infrastructure/src/test/java/org/springframework/batch/repeat/context/SynchronizedAttributeAccessorTests.java b/spring-batch-infrastructure/src/test/java/org/springframework/batch/repeat/context/SynchronizedAttributeAccessorTests.java index 443fb2411..fc4580dd4 100644 --- a/spring-batch-infrastructure/src/test/java/org/springframework/batch/repeat/context/SynchronizedAttributeAccessorTests.java +++ b/spring-batch-infrastructure/src/test/java/org/springframework/batch/repeat/context/SynchronizedAttributeAccessorTests.java @@ -34,7 +34,8 @@ public class SynchronizedAttributeAccessorTests extends TestCase { SynchronizedAttributeAccessor another = new SynchronizedAttributeAccessor(); accessor.setAttribute("foo", "bar"); another.setAttribute("foo", "bar"); - assertEquals(accessor.hashCode(), another.hashCode()); + assertEquals(accessor, another); + assertEquals("Object.hashCode() contract broken", accessor.hashCode(), another.hashCode()); } public void testToStringWithNoAttributes() throws Exception { diff --git a/spring-batch-samples/src/main/java/org/springframework/batch/sample/item/reader/OrderItemReader.java b/spring-batch-samples/src/main/java/org/springframework/batch/sample/item/reader/OrderItemReader.java index 46da0c23d..facab9ef6 100644 --- a/spring-batch-samples/src/main/java/org/springframework/batch/sample/item/reader/OrderItemReader.java +++ b/spring-batch-samples/src/main/java/org/springframework/batch/sample/item/reader/OrderItemReader.java @@ -23,8 +23,7 @@ import org.apache.commons.logging.LogFactory; import org.springframework.batch.core.domain.StepExecution; import org.springframework.batch.io.file.mapping.FieldSet; import org.springframework.batch.io.file.mapping.FieldSetMapper; -import org.springframework.batch.item.ItemReader; -import org.springframework.batch.item.reader.AbstractItemReader; +import org.springframework.batch.item.reader.DelegatingItemReader; import org.springframework.batch.item.validator.Validator; import org.springframework.batch.sample.domain.Address; import org.springframework.batch.sample.domain.BillingInfo; @@ -33,177 +32,180 @@ import org.springframework.batch.sample.domain.LineItem; import org.springframework.batch.sample.domain.Order; import org.springframework.batch.sample.domain.ShippingInfo; - - /** * @author peter.zozom - * + * */ -public class OrderItemReader extends AbstractItemReader { - private static Log log = LogFactory.getLog(OrderItemReader.class); - private ItemReader inputSource; - private Order order; - private boolean recordFinished; - private FieldSetMapper headerMapper; - private FieldSetMapper customerMapper; - private FieldSetMapper addressMapper; - private FieldSetMapper billingMapper; - private FieldSetMapper itemMapper; - private FieldSetMapper shippingMapper; - private Validator validator; +public class OrderItemReader extends DelegatingItemReader { + private static Log log = LogFactory.getLog(OrderItemReader.class); - /** - * @throws Exception - * @see org.springframework.batch.item.ItemReader#read() - */ - public Object read() throws Exception { - recordFinished = false; + private Order order; - while (!recordFinished) { - process((FieldSet)inputSource.read()); - } + private boolean recordFinished; - if (order!=null) { - log.info("Mapped: "+order); - validator.validate(order); - } + private FieldSetMapper headerMapper; - Object result = order; - order = null; + private FieldSetMapper customerMapper; - return result; - } + private FieldSetMapper addressMapper; - /** - * @see org.springframework.batch.execution.io.FieldSetCallback#execute(StepExecution) - */ - private void process(FieldSet fieldSet) { - //finish processing if we hit the end of file - if (fieldSet == null) { - log.debug("FINISHED"); - recordFinished = true; - order = null; + private FieldSetMapper billingMapper; - return; - } + private FieldSetMapper itemMapper; - String lineId = fieldSet.readString(0); + private FieldSetMapper shippingMapper; - //start a new Order - if (Order.LINE_ID_HEADER.equals(lineId)) { - log.debug("STARTING NEW RECORD"); - order = (Order) headerMapper.mapLine(fieldSet); + private Validator validator; - return; - } + /** + * @throws Exception + * @see org.springframework.batch.item.ItemReader#read() + */ + public Object read() throws Exception { + recordFinished = false; - //mark we are finished with current Order - if (Order.LINE_ID_FOOTER.equals(lineId)) { - log.debug("END OF RECORD"); + while (!recordFinished) { + process((FieldSet) getItemReader().read()); + } - //Do mapping for footer here, because mapper does not allow to pass an Order object as input. - //Mapper always creates new object - order.setTotalPrice(fieldSet.readBigDecimal("TOTAL_PRICE")); - order.setTotalLines(fieldSet.readInt("TOTAL_LINE_ITEMS")); - order.setTotalItems(fieldSet.readInt("TOTAL_ITEMS")); + if (order != null) { + log.info("Mapped: " + order); + validator.validate(order); + } - recordFinished = true; + Object result = order; + order = null; - return; - } + return result; + } - if (Customer.LINE_ID_BUSINESS_CUST.equals(lineId)) { - log.debug("MAPPING CUSTOMER"); + /** + * @see org.springframework.batch.execution.io.FieldSetCallback#execute(StepExecution) + */ + private void process(FieldSet fieldSet) { + // finish processing if we hit the end of file + if (fieldSet == null) { + log.debug("FINISHED"); + recordFinished = true; + order = null; - if (order.getCustomer() == null) { - order.setCustomer((Customer) customerMapper.mapLine(fieldSet)); - order.getCustomer().setBusinessCustomer(true); - } + return; + } - return; - } + String lineId = fieldSet.readString(0); - if (Customer.LINE_ID_NON_BUSINESS_CUST.equals(lineId)) { - log.debug("MAPPING CUSTOMER"); + // start a new Order + if (Order.LINE_ID_HEADER.equals(lineId)) { + log.debug("STARTING NEW RECORD"); + order = (Order) headerMapper.mapLine(fieldSet); - if (order.getCustomer() == null) { - order.setCustomer((Customer) customerMapper.mapLine(fieldSet)); - order.getCustomer().setBusinessCustomer(false); - } + return; + } - return; - } + // mark we are finished with current Order + if (Order.LINE_ID_FOOTER.equals(lineId)) { + log.debug("END OF RECORD"); - if (Address.LINE_ID_BILLING_ADDR.equals(lineId)) { - log.debug("MAPPING BILLING ADDRESS"); - order.setBillingAddress((Address) addressMapper.mapLine(fieldSet)); - return; - } + // Do mapping for footer here, because mapper does not allow to pass + // an Order object as input. + // Mapper always creates new object + order.setTotalPrice(fieldSet.readBigDecimal("TOTAL_PRICE")); + order.setTotalLines(fieldSet.readInt("TOTAL_LINE_ITEMS")); + order.setTotalItems(fieldSet.readInt("TOTAL_ITEMS")); - if (Address.LINE_ID_SHIPPING_ADDR.equals(lineId)) { - log.debug("MAPPING SHIPPING ADDRESS"); - order.setShippingAddress((Address) addressMapper.mapLine(fieldSet)); - return; - } + recordFinished = true; - if (BillingInfo.LINE_ID_BILLING_INFO.equals(lineId)) { - log.debug("MAPPING BILLING INFO"); - order.setBilling((BillingInfo) billingMapper.mapLine(fieldSet)); - return; - } + return; + } - if (ShippingInfo.LINE_ID_SHIPPING_INFO.equals(lineId)) { - log.debug("MAPPING SHIPPING INFO"); - order.setShipping((ShippingInfo) shippingMapper.mapLine(fieldSet)); - return; - } + if (Customer.LINE_ID_BUSINESS_CUST.equals(lineId)) { + log.debug("MAPPING CUSTOMER"); - if (LineItem.LINE_ID_ITEM.equals(lineId)) { - log.debug("MAPPING LINE ITEM"); + if (order.getCustomer() == null) { + order.setCustomer((Customer) customerMapper.mapLine(fieldSet)); + order.getCustomer().setBusinessCustomer(true); + } - if (order.getLineItems() == null) { - order.setLineItems(new ArrayList()); - } + return; + } - order.getLineItems().add(itemMapper.mapLine(fieldSet)); + if (Customer.LINE_ID_NON_BUSINESS_CUST.equals(lineId)) { + log.debug("MAPPING CUSTOMER"); - return; - } + if (order.getCustomer() == null) { + order.setCustomer((Customer) customerMapper.mapLine(fieldSet)); + order.getCustomer().setBusinessCustomer(false); + } - log.debug("Could not map LINE_ID="+lineId); + return; + } - } + if (Address.LINE_ID_BILLING_ADDR.equals(lineId)) { + log.debug("MAPPING BILLING ADDRESS"); + order.setBillingAddress((Address) addressMapper.mapLine(fieldSet)); + return; + } - public void setAddressMapper(FieldSetMapper addressMapper) { - this.addressMapper = addressMapper; - } + if (Address.LINE_ID_SHIPPING_ADDR.equals(lineId)) { + log.debug("MAPPING SHIPPING ADDRESS"); + order.setShippingAddress((Address) addressMapper.mapLine(fieldSet)); + return; + } - public void setBillingMapper(FieldSetMapper billingMapper) { - this.billingMapper = billingMapper; - } + if (BillingInfo.LINE_ID_BILLING_INFO.equals(lineId)) { + log.debug("MAPPING BILLING INFO"); + order.setBilling((BillingInfo) billingMapper.mapLine(fieldSet)); + return; + } - public void setCustomerMapper(FieldSetMapper customerMapper) { - this.customerMapper = customerMapper; - } + if (ShippingInfo.LINE_ID_SHIPPING_INFO.equals(lineId)) { + log.debug("MAPPING SHIPPING INFO"); + order.setShipping((ShippingInfo) shippingMapper.mapLine(fieldSet)); + return; + } - public void setHeaderMapper(FieldSetMapper headerMapper) { - this.headerMapper = headerMapper; - } + if (LineItem.LINE_ID_ITEM.equals(lineId)) { + log.debug("MAPPING LINE ITEM"); - public void setItemReader(ItemReader inputSource) { - this.inputSource = inputSource; - } + if (order.getLineItems() == null) { + order.setLineItems(new ArrayList()); + } - public void setItemMapper(FieldSetMapper itemMapper) { - this.itemMapper = itemMapper; - } + order.getLineItems().add(itemMapper.mapLine(fieldSet)); - public void setShippingMapper(FieldSetMapper shippingMapper) { - this.shippingMapper = shippingMapper; - } + return; + } - public void setValidator(Validator validator) { - this.validator = validator; - } + log.debug("Could not map LINE_ID=" + lineId); + + } + + public void setAddressMapper(FieldSetMapper addressMapper) { + this.addressMapper = addressMapper; + } + + public void setBillingMapper(FieldSetMapper billingMapper) { + this.billingMapper = billingMapper; + } + + public void setCustomerMapper(FieldSetMapper customerMapper) { + this.customerMapper = customerMapper; + } + + public void setHeaderMapper(FieldSetMapper headerMapper) { + this.headerMapper = headerMapper; + } + + public void setItemMapper(FieldSetMapper itemMapper) { + this.itemMapper = itemMapper; + } + + public void setShippingMapper(FieldSetMapper shippingMapper) { + this.shippingMapper = shippingMapper; + } + + public void setValidator(Validator validator) { + this.validator = validator; + } } diff --git a/spring-batch-samples/src/main/java/org/springframework/batch/sample/tasklet/ExceptionThrowingItemReaderProxy.java b/spring-batch-samples/src/main/java/org/springframework/batch/sample/tasklet/ExceptionThrowingItemReaderProxy.java index 6bb7e7e48..d0e35ff77 100644 --- a/spring-batch-samples/src/main/java/org/springframework/batch/sample/tasklet/ExceptionThrowingItemReaderProxy.java +++ b/spring-batch-samples/src/main/java/org/springframework/batch/sample/tasklet/ExceptionThrowingItemReaderProxy.java @@ -19,6 +19,7 @@ package org.springframework.batch.sample.tasklet; import org.springframework.batch.io.exception.BatchCriticalException; import org.springframework.batch.item.ItemReader; +import org.springframework.batch.item.reader.DelegatingItemReader; /** * Hacked {@link ItemReader} that throws exception on a given record number @@ -28,17 +29,11 @@ import org.springframework.batch.item.ItemReader; * @author Lucas Ward * */ -public class ExceptionThrowingItemReaderProxy implements ItemReader { +public class ExceptionThrowingItemReaderProxy extends DelegatingItemReader { private int counter = 0; private int throwExceptionOnRecordNumber = 4; - private final ItemReader itemReader; - - public ExceptionThrowingItemReaderProxy(ItemReader itemReader) { - this.itemReader = itemReader; - } - /** * @param throwExceptionOnRecordNumber The number of record on which exception should be thrown */ @@ -57,7 +52,7 @@ public class ExceptionThrowingItemReaderProxy implements ItemReader { throw new BatchCriticalException("Planned failure on count="+counter); } - return itemReader.read(); + return getItemReader().read(); } } diff --git a/spring-batch-samples/src/main/resources/jobs/restartSample.xml b/spring-batch-samples/src/main/resources/jobs/restartSample.xml index 3112144ba..b86f5b166 100644 --- a/spring-batch-samples/src/main/resources/jobs/restartSample.xml +++ b/spring-batch-samples/src/main/resources/jobs/restartSample.xml @@ -39,7 +39,7 @@ - +