[BATCH-430] Incremental commit of exception moving

This commit is contained in:
nebhale
2008-03-07 11:54:20 +00:00
parent a60fc1e621
commit 8c2901093e
106 changed files with 736 additions and 848 deletions

View File

@@ -16,12 +16,12 @@
package org.springframework.batch.execution.step;
import org.springframework.batch.core.StepContribution;
import org.springframework.batch.item.ClearFailedException;
import org.springframework.batch.item.FlushFailedException;
import org.springframework.batch.item.ItemReader;
import org.springframework.batch.item.ItemWriter;
import org.springframework.batch.item.exception.ClearFailedException;
import org.springframework.batch.item.exception.FlushFailedException;
import org.springframework.batch.item.exception.MarkFailedException;
import org.springframework.batch.item.exception.ResetFailedException;
import org.springframework.batch.item.MarkFailedException;
import org.springframework.batch.item.ResetFailedException;
import org.springframework.batch.repeat.ExitStatus;
/**

View File

@@ -39,7 +39,6 @@ import org.springframework.batch.item.ExecutionContext;
import org.springframework.batch.item.ItemReader;
import org.springframework.batch.item.ItemStream;
import org.springframework.batch.item.ItemWriter;
import org.springframework.batch.item.exception.CommitFailedException;
import org.springframework.batch.item.stream.CompositeItemStream;
import org.springframework.batch.repeat.ExitStatus;
import org.springframework.batch.repeat.RepeatCallback;
@@ -51,20 +50,16 @@ import org.springframework.transaction.TransactionStatus;
import org.springframework.transaction.support.DefaultTransactionDefinition;
/**
* Simple implementation of executing the step as a set of chunks, each chunk
* surrounded by a transaction. The structure is therefore that of two nested
* loops, with transaction boundary around the whole inner loop. The outer loop
* is controlled by the step operations ({@link #setStepOperations(RepeatOperations)}),
* and the inner loop by the chunk operations ({@link #setChunkOperations(RepeatOperations)}).
* The inner loop should always be executed in a single thread, so the chunk
* operations should not do any concurrent execution. N.B. usually that means
* that the chunk operations should be a {@link RepeatTemplate} (which is the
* default).<br/>
* Simple implementation of executing the step as a set of chunks, each chunk surrounded by a transaction. The structure
* is therefore that of two nested loops, with transaction boundary around the whole inner loop. The outer loop is
* controlled by the step operations ({@link #setStepOperations(RepeatOperations)}), and the inner loop by the chunk
* operations ({@link #setChunkOperations(RepeatOperations)}). The inner loop should always be executed in a single
* thread, so the chunk operations should not do any concurrent execution. N.B. usually that means that the chunk
* operations should be a {@link RepeatTemplate} (which is the default).<br/>
*
* Clients can use interceptors 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. Those that want callbacks at the level of an individual
* tasks, can specify interceptors for the chunk operations.
* Clients can use interceptors 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. Those that want callbacks at the level of an individual tasks,
* can specify interceptors for the chunk operations.
*
* @author Dave Syer
* @author Lucas Ward
@@ -122,6 +117,7 @@ public class ItemOrientedStep extends AbstractStep {
/**
* Public setter for the {@link ItemHandler}.
*
* @param itemHandler the {@link ItemHandler} to set
*/
public void setItemHandler(ItemHandler itemHandler) {
@@ -129,13 +125,10 @@ public class ItemOrientedStep extends AbstractStep {
}
/**
* Register each of the streams for callbacks at the appropriate time in the
* step. The {@link ItemReader} and {@link ItemWriter} are automatically
* registered, but it doesn't hurt to also register them here. Injected
* dependencies of the reader and writer are not automatically registered,
* so if you implement {@link ItemWriter} using delegation to another object
* which itself is a {@link ItemStream}, you need to register the delegate
* here.
* Register each of the streams for callbacks at the appropriate time in the step. The {@link ItemReader} and
* {@link ItemWriter} are automatically registered, but it doesn't hurt to also register them here. Injected
* dependencies of the reader and writer are not automatically registered, so if you implement {@link ItemWriter}
* using delegation to another object which itself is a {@link ItemStream}, you need to register the delegate here.
*
* @param streams an array of {@link ItemStream} objects.
*/
@@ -146,8 +139,8 @@ public class ItemOrientedStep extends AbstractStep {
}
/**
* Register a single {@link ItemStream} for callbacks to the stream
* interface.
* Register a single {@link ItemStream} for callbacks to the stream interface.
*
* @param stream
*/
public void registerStream(ItemStream stream) {
@@ -155,11 +148,9 @@ public class ItemOrientedStep extends AbstractStep {
}
/**
* Register each of the objects as listeners. If the {@link ItemReader} or
* {@link ItemWriter} themselves implements this interface they will be
* registered automatically, but their injected dependencies will not be.
* This is a good way to get access to job parameters and execution context
* if the tasklet is parameterised.
* Register each of the objects as listeners. If the {@link ItemReader} or {@link ItemWriter} themselves implements
* this interface they will be registered automatically, but their injected dependencies will not be. This is a good
* way to get access to job parameters and execution context if the tasklet is parameterised.
*
* @param listeners an array of listener objects of known types.
*/
@@ -170,8 +161,7 @@ public class ItemOrientedStep extends AbstractStep {
}
/**
* Register a step listener for callbacks at the appropriate stages in a
* step execution.
* Register a step listener for callbacks at the appropriate stages in a step execution.
*
* @param listener a {@link StepListener}
*/
@@ -180,9 +170,8 @@ public class ItemOrientedStep extends AbstractStep {
}
/**
* 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.
*/
@@ -191,9 +180,8 @@ public class ItemOrientedStep extends AbstractStep {
}
/**
* The {@link RepeatOperations} to use for the inner 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 inner loop of the batch processing. should be set up by the caller
* through a factory. defaults to a plain {@link RepeatTemplate}.
*
* @param chunkOperations a {@link RepeatOperations} instance.
*/
@@ -202,9 +190,8 @@ public class ItemOrientedStep extends AbstractStep {
}
/**
* 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}
*/
@@ -213,8 +200,8 @@ public class ItemOrientedStep extends AbstractStep {
}
/**
* 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
*/
@@ -223,9 +210,9 @@ public class ItemOrientedStep extends AbstractStep {
}
/**
* Mostly useful for testing, but could be used to remove dependence on
* backport concurrency utilities. Public setter for the
* {@link StepExecutionSynchronizer}.
* Mostly useful for testing, but could be used to remove dependence on backport concurrency utilities. Public
* setter for the {@link StepExecutionSynchronizer}.
*
* @param synchronizer the {@link StepExecutionSynchronizer} to set
*/
public void setSynchronizer(StepExecutionSynchronizer synchronizer) {
@@ -233,18 +220,14 @@ public class ItemOrientedStep extends AbstractStep {
}
/**
* 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 JobExecutionContext}.<br/>
* 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 JobExecutionContext}.<br/>
*
* @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 InfrastructureException, JobInterruptedException {
@@ -267,8 +250,7 @@ public class ItemOrientedStep extends AbstractStep {
if (isRestart && lastStepExecution != null) {
stepExecution.setExecutionContext(lastStepExecution.getExecutionContext());
}
else {
} else {
stepExecution.setExecutionContext(new ExecutionContext());
}
@@ -291,7 +273,7 @@ public class ItemOrientedStep extends AbstractStep {
ExitStatus result = ExitStatus.CONTINUABLE;
TransactionStatus transaction = transactionManager
.getTransaction(new DefaultTransactionDefinition());
.getTransaction(new DefaultTransactionDefinition());
try {
@@ -304,8 +286,7 @@ public class ItemOrientedStep extends AbstractStep {
// minimum).
try {
synchronizer.lock(stepExecution);
}
catch (InterruptedException e) {
} catch (InterruptedException e) {
stepExecution.setStatus(BatchStatus.STOPPED);
Thread.currentThread().interrupt();
}
@@ -317,36 +298,30 @@ public class ItemOrientedStep extends AbstractStep {
stream.update(stepExecution.getExecutionContext());
try {
jobRepository.saveOrUpdateExecutionContext(stepExecution);
}
catch (Exception e) {
} catch (Exception e) {
fatalException.setException(e);
stepExecution.setStatus(BatchStatus.UNKNOWN);
throw new CommitFailedException(
"Fatal error detected during save of step execution context", e);
"Fatal error detected during save of step execution context", e);
}
try {
itemHandler.mark();
itemHandler.flush();
transactionManager.commit(transaction);
}
catch (Exception e) {
} catch (Exception e) {
fatalException.setException(e);
stepExecution.setStatus(BatchStatus.UNKNOWN);
throw new CommitFailedException("Fatal error detected during commit", e);
}
}
catch (CommitFailedException e) {
} catch (CommitFailedException e) {
throw e;
}
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.
*/
stepExecution.rollback();
@@ -354,21 +329,18 @@ public class ItemOrientedStep extends AbstractStep {
itemHandler.reset();
itemHandler.clear();
transactionManager.rollback(transaction);
}
catch (Exception e) {
} catch (Exception e) {
fatalException.setException(e);
stepExecution.setStatus(BatchStatus.UNKNOWN);
}
if (t instanceof RuntimeException) {
throw (RuntimeException) t;
}
else {
throw new RuntimeException(t);
} else {
throw new RuntimeException(t);
}
}
finally {
} finally {
synchronizer.release(stepExecution);
}
@@ -383,12 +355,10 @@ public class ItemOrientedStep extends AbstractStep {
});
fatalException.setException(updateStatus(stepExecution, BatchStatus.COMPLETED));
}
catch (CommitFailedException e) {
} catch (CommitFailedException e) {
logger.error("Fatal error detected during commit.");
throw e;
}
catch (RuntimeException e) {
} catch (RuntimeException e) {
// classify exception so an exit code can be stored.
status = exceptionClassifier.classifyForExitCode(e);
@@ -396,29 +366,24 @@ public class ItemOrientedStep extends AbstractStep {
if (e.getCause() instanceof JobInterruptedException) {
updateStatus(stepExecution, BatchStatus.STOPPED);
throw (JobInterruptedException) e.getCause();
}
else if (!fatalException.hasException()) {
} else if (!fatalException.hasException()) {
try {
status = status.and(listener.onErrorInStep(stepExecution, e));
}
catch (RuntimeException ex) {
} catch (RuntimeException ex) {
logger.error("Unexpected error in listener on error in step.", ex);
}
updateStatus(stepExecution, BatchStatus.FAILED);
throw e;
}
else {
} else {
logger.error("Fatal error detected during rollback caused by underlying exception: ", e);
throw e;
}
}
finally {
} finally {
try {
status = status.and(listener.afterStep(stepExecution));
}
catch (RuntimeException e) {
} catch (RuntimeException e) {
logger.error("Unexpected error in listener after step.", e);
}
@@ -427,8 +392,7 @@ public class ItemOrientedStep extends AbstractStep {
try {
jobRepository.saveOrUpdate(stepExecution);
}
catch (RuntimeException e) {
} catch (RuntimeException e) {
String msg = "Fatal error detected during final save of meta data";
logger.error(msg, e);
if (!fatalException.hasException()) {
@@ -439,10 +403,9 @@ public class ItemOrientedStep extends AbstractStep {
try {
stream.close(stepExecution.getExecutionContext());
}
catch (RuntimeException e) {
} catch (RuntimeException e) {
String msg = "Fatal error detected during close of streams. "
+ "The job execution completed (possibly unsuccessfully but with consistent meta-data).";
+ "The job execution completed (possibly unsuccessfully but with consistent meta-data).";
logger.error(msg, e);
if (!fatalException.hasException()) {
fatalException.setException(e);
@@ -452,7 +415,7 @@ public class ItemOrientedStep extends AbstractStep {
if (fatalException.hasException()) {
throw new InfrastructureException("Encountered an error saving batch meta data.", fatalException
.getException());
.getException());
}
}
@@ -460,13 +423,11 @@ public class ItemOrientedStep extends AbstractStep {
}
/**
* Execute a bunch of identical business logic operations all within a
* transaction. The transaction is programmatically started and stopped
* outside this method, so subclasses that override do not need to create a
* Execute a bunch of identical business logic operations all within a transaction. The transaction is
* programmatically started and stopped outside this method, so subclasses that override do not need to create a
* transaction.
*
* @param step the current step containing the {@link Tasklet} with the
* business logic.
* @param step the current step containing the {@link Tasklet} with the business logic.
* @return true if there is more data to process.
*/
protected ExitStatus processChunk(final StepContribution contribution) {
@@ -499,8 +460,7 @@ public class ItemOrientedStep extends AbstractStep {
try {
jobRepository.saveOrUpdate(stepExecution);
return null;
}
catch (Exception e) {
} catch (Exception e) {
return e;
}
@@ -534,4 +494,11 @@ public class ItemOrientedStep extends AbstractStep {
}
private class CommitFailedException extends RuntimeException {
public CommitFailedException(String string, Exception e) {
super(string, e);
}
}
}

View File

@@ -17,12 +17,12 @@ package org.springframework.batch.execution.step.support;
import org.springframework.batch.core.StepContribution;
import org.springframework.batch.execution.step.ItemHandler;
import org.springframework.batch.item.ClearFailedException;
import org.springframework.batch.item.FlushFailedException;
import org.springframework.batch.item.ItemReader;
import org.springframework.batch.item.ItemWriter;
import org.springframework.batch.item.exception.ClearFailedException;
import org.springframework.batch.item.exception.FlushFailedException;
import org.springframework.batch.item.exception.MarkFailedException;
import org.springframework.batch.item.exception.ResetFailedException;
import org.springframework.batch.item.MarkFailedException;
import org.springframework.batch.item.ResetFailedException;
import org.springframework.batch.repeat.ExitStatus;
/**

View File

@@ -21,9 +21,9 @@ import java.util.List;
import org.apache.commons.logging.Log;
import org.apache.commons.logging.LogFactory;
import org.springframework.batch.item.ClearFailedException;
import org.springframework.batch.item.FlushFailedException;
import org.springframework.batch.item.ItemWriter;
import org.springframework.batch.item.exception.ClearFailedException;
import org.springframework.batch.item.exception.FlushFailedException;
import org.springframework.batch.support.transaction.TransactionAwareProxyFactory;
import org.springframework.beans.factory.InitializingBean;

View File

@@ -44,9 +44,9 @@ import org.springframework.batch.item.ExecutionContext;
import org.springframework.batch.item.ItemReader;
import org.springframework.batch.item.ItemStream;
import org.springframework.batch.item.ItemWriter;
import org.springframework.batch.item.exception.MarkFailedException;
import org.springframework.batch.item.exception.ResetFailedException;
import org.springframework.batch.item.exception.StreamException;
import org.springframework.batch.item.MarkFailedException;
import org.springframework.batch.item.ResetFailedException;
import org.springframework.batch.item.ItemStreamException;
import org.springframework.batch.item.reader.AbstractItemReader;
import org.springframework.batch.item.reader.ListItemReader;
import org.springframework.batch.item.stream.ItemStreamSupport;
@@ -380,7 +380,7 @@ public class ItemOrientedStepTests extends TestCase {
public void beforeStep(StepExecution stepExecution) {
list.add("foo");
}
public void open(ExecutionContext executionContext) throws StreamException {
public void open(ExecutionContext executionContext) throws ItemStreamException {
assertEquals(1, list.size());
}
};
@@ -633,7 +633,7 @@ public class ItemOrientedStepTests extends TestCase {
public void testStatusForCloseFailedException() throws Exception {
MockRestartableItemReader itemReader = new MockRestartableItemReader() {
public void close(ExecutionContext executionContext) throws StreamException {
public void close(ExecutionContext executionContext) throws ItemStreamException {
super.close(executionContext);
// Simulate failure on rollback when stream resets
throw new RuntimeException("Bar");

View File

@@ -16,8 +16,8 @@
package org.springframework.batch.execution.step.support;
import org.springframework.batch.item.ItemReader;
import org.springframework.batch.item.exception.MarkFailedException;
import org.springframework.batch.item.exception.ResetFailedException;
import org.springframework.batch.item.MarkFailedException;
import org.springframework.batch.item.ResetFailedException;
public class MockItemReader implements ItemReader {