diff --git a/spring-batch-core/src/main/java/org/springframework/batch/core/step/AbstractStep.java b/spring-batch-core/src/main/java/org/springframework/batch/core/step/AbstractStep.java index 2302901c4..fe5ad0998 100644 --- a/spring-batch-core/src/main/java/org/springframework/batch/core/step/AbstractStep.java +++ b/spring-batch-core/src/main/java/org/springframework/batch/core/step/AbstractStep.java @@ -211,7 +211,7 @@ public abstract class AbstractStep implements Step, InitializingBean, BeanNameAw logger.debug("Step execution success: id=" + stepExecution.getId()); } catch (Throwable e) { - logger.error("Encountered an error executing the step: " + e.getClass() + ": " + e.getMessage(), e); + logger.error("Encountered an error executing the step", e); stepExecution.setStatus(determineBatchStatus(e)); exitStatus = exitStatus.and(getDefaultExitStatusForFailure(e)); stepExecution.addFailureException(e); diff --git a/spring-batch-core/src/main/java/org/springframework/batch/core/step/item/ChunkOrientedTasklet.java b/spring-batch-core/src/main/java/org/springframework/batch/core/step/item/ChunkOrientedTasklet.java index 5d540b71b..632d4e6ea 100644 --- a/spring-batch-core/src/main/java/org/springframework/batch/core/step/item/ChunkOrientedTasklet.java +++ b/spring-batch-core/src/main/java/org/springframework/batch/core/step/item/ChunkOrientedTasklet.java @@ -73,8 +73,6 @@ public class ChunkOrientedTasklet implements Tasklet { // Allow a message coming back from the processor to say that we // are not done yet if (inputs.isBusy()) { - // TODO: update ExecutionContext with an offset if the - // ItemReader was stateful return RepeatStatus.CONTINUABLE; } diff --git a/spring-batch-core/src/main/java/org/springframework/batch/core/step/item/FaultTolerantChunkProcessor.java b/spring-batch-core/src/main/java/org/springframework/batch/core/step/item/FaultTolerantChunkProcessor.java index ab2cdccc8..d6a67ca18 100755 --- a/spring-batch-core/src/main/java/org/springframework/batch/core/step/item/FaultTolerantChunkProcessor.java +++ b/spring-batch-core/src/main/java/org/springframework/batch/core/step/item/FaultTolerantChunkProcessor.java @@ -277,26 +277,14 @@ public class FaultTolerantChunkProcessor extends SimpleChunkProcessor recoveryCallback = new RecoveryCallback() { public Object recover(RetryContext context) throws Exception { - - Exception le = (Exception) context.getLastThrowable(); - - boolean singleton = outputs.size() == 1 && outputs.getSkips().isEmpty(); - - if (singleton) { - Chunk.ChunkIterator inputIterator = inputs.iterator(); - Chunk.ChunkIterator outputIterator = outputs.iterator(); - checkSkipPolicy(inputIterator, outputIterator, le, contribution); - return null; - } - inputs.setBusy(true); scan(contribution, inputs, outputs, chunkMonitor); return null; - } }; + logger.debug("Attempting to write: "+inputs); batchRetryTemplate.execute(retryCallback, recoveryCallback, new DefaultRetryState(inputs, rollbackClassifier)); @@ -358,6 +346,7 @@ public class FaultTolerantChunkProcessor extends SimpleChunkProcessor.ChunkIterator inputIterator, Chunk.ChunkIterator outputIterator, Exception e, StepContribution contribution) { + logger.debug("Checking skip policy after failed write"); if (itemWriteSkipPolicy.shouldSkip(e, contribution.getStepSkipCount())) { contribution.incrementWriteSkipCount(); inputIterator.remove(); @@ -372,6 +361,7 @@ public class FaultTolerantChunkProcessor extends SimpleChunkProcessor inputs, final Chunk outputs, ChunkMonitor chunkMonitor) throws Exception { + logger.debug("Scanning for failed item on write."); if (outputs.isEmpty()) { inputs.setBusy(false); return; diff --git a/spring-batch-core/src/main/java/org/springframework/batch/core/step/item/SimpleRetryExceptionHandler.java b/spring-batch-core/src/main/java/org/springframework/batch/core/step/item/SimpleRetryExceptionHandler.java index b0d25e011..57aae53e0 100644 --- a/spring-batch-core/src/main/java/org/springframework/batch/core/step/item/SimpleRetryExceptionHandler.java +++ b/spring-batch-core/src/main/java/org/springframework/batch/core/step/item/SimpleRetryExceptionHandler.java @@ -79,10 +79,11 @@ public class SimpleRetryExceptionHandler extends RetryListenerSupport implements // Only bother to check the delegate exception handler if we know that // retry is exhausted if (fatalExceptionClassifier.classify(throwable) || context.hasAttribute(EXHAUSTED)) { + logger.debug("Handled fatal exception"); exceptionHandler.handleException(context, throwable); } else { - logger.debug("handled non-fatal exception", throwable); + logger.debug("Handled non-fatal exception", throwable); } } @@ -95,6 +96,7 @@ public class SimpleRetryExceptionHandler extends RetryListenerSupport implements */ public void close(RetryContext context, RetryCallback callback, Throwable throwable) { if (!retryPolicy.canRetry(context)) { + logger.debug("Marking retry as exhausted: "+context); getRepeatContext().setAttribute(EXHAUSTED, "true"); } } diff --git a/spring-batch-core/src/main/java/org/springframework/batch/core/step/tasklet/TaskletStep.java b/spring-batch-core/src/main/java/org/springframework/batch/core/step/tasklet/TaskletStep.java index 9522b10e5..941f26b42 100644 --- a/spring-batch-core/src/main/java/org/springframework/batch/core/step/tasklet/TaskletStep.java +++ b/spring-batch-core/src/main/java/org/springframework/batch/core/step/tasklet/TaskletStep.java @@ -306,6 +306,7 @@ public class TaskletStep extends AbstractStep { } catch (FatalException e) { try { + logger.debug("Rollback for FatalException: "+e.getClass().getName()+": "+e.getMessage()); rollback(stepExecution, transaction); } catch (Exception rollbackException) { @@ -317,6 +318,7 @@ public class TaskletStep extends AbstractStep { } catch (Error e) { try { + logger.debug("Rollback for Error: "+e.getClass().getName()+": "+e.getMessage()); rollback(stepExecution, transaction); } catch (Exception rollbackException) { @@ -328,6 +330,7 @@ public class TaskletStep extends AbstractStep { } catch (Exception e) { try { + logger.debug("Rollback for Exception: "+e.getClass().getName()+": "+e.getMessage()); rollback(stepExecution, transaction); } catch (Exception rollbackException) { diff --git a/spring-batch-core/src/test/java/org/springframework/batch/core/step/item/ExceptionThrowingItemHandlerStub.java b/spring-batch-core/src/test/java/org/springframework/batch/core/step/item/ExceptionThrowingItemHandlerStub.java index e92e9eaaf..7c4846ce0 100644 --- a/spring-batch-core/src/test/java/org/springframework/batch/core/step/item/ExceptionThrowingItemHandlerStub.java +++ b/spring-batch-core/src/test/java/org/springframework/batch/core/step/item/ExceptionThrowingItemHandlerStub.java @@ -19,11 +19,16 @@ import java.util.Arrays; import java.util.Collection; import java.util.Collections; +import org.apache.commons.logging.Log; +import org.apache.commons.logging.LogFactory; + /** * @author Dan Garrette * @since 2.0.1 */ public abstract class ExceptionThrowingItemHandlerStub { + + protected Log logger = LogFactory.getLog(getClass()); private Collection failures = Collections.emptyList(); @@ -44,10 +49,10 @@ public abstract class ExceptionThrowingItemHandlerStub { protected void checkFailure(T item) throws Exception { if (isFailure(item)) { if (runtimeException) { - throw new SkippableRuntimeException("Intended Failure"); + throw new SkippableRuntimeException("Intended Failure: "+item); } else { - throw new SkippableException("Intended Failure"); + throw new SkippableException("Intended Failure: "+item); } } } diff --git a/spring-batch-core/src/test/java/org/springframework/batch/core/step/item/FaultTolerantStepFactoryBeanRetryTests.java b/spring-batch-core/src/test/java/org/springframework/batch/core/step/item/FaultTolerantStepFactoryBeanRetryTests.java index 185489245..e4239afd3 100644 --- a/spring-batch-core/src/test/java/org/springframework/batch/core/step/item/FaultTolerantStepFactoryBeanRetryTests.java +++ b/spring-batch-core/src/test/java/org/springframework/batch/core/step/item/FaultTolerantStepFactoryBeanRetryTests.java @@ -314,13 +314,9 @@ public class FaultTolerantStepFactoryBeanRetryTests { List expectedOutput = Arrays.asList(StringUtils.commaDelimitedListToStringArray("a,c,e,f")); assertEquals(expectedOutput, written); - // [a, b, c, d, e, f, null] - assertEquals(7, provided.size()); - // [a, b, b, b, b, b, c, d, d, d, d, d, e, f] - System.err.println(processed); - assertEquals(14, processed.size()); - // [b, d] - assertEquals(2, recovered.size()); + assertEquals("[a, b, c, d, e, f, null]", provided.toString()); + assertEquals("[a, b, b, b, b, b, b, c, d, d, d, d, d, d, e, f]", processed.toString()); + assertEquals("[b, d]", recovered.toString()); } @Test @@ -423,8 +419,7 @@ public class FaultTolerantStepFactoryBeanRetryTests { assertEquals(1, provided.size()); // the failed items are tried up to the limit (but only precisely so if // the commit interval is 1) - // [b, b, b, b] - assertEquals(4, processed.size()); + assertEquals("[b, b, b, b, b]", processed.toString()); // [] assertEquals(0, recovered.size()); assertEquals(1, stepExecution.getReadCount()); @@ -474,9 +469,9 @@ public class FaultTolerantStepFactoryBeanRetryTests { assertEquals(0, stepExecution.getSkipCount()); // [b] - assertEquals(1, provided.size()); + assertEquals("[b]", provided.toString()); // [b] - assertEquals(1, processed.size()); + assertEquals("[b, b]", processed.toString()); // [] assertEquals(0, recovered.size()); assertEquals(1, stepExecution.getReadCount()); @@ -517,8 +512,7 @@ public class FaultTolerantStepFactoryBeanRetryTests { assertEquals(0, stepExecution.getSkipCount()); // [b] assertEquals(1, provided.size()); - // [b, b, b, b] - assertEquals(4, processed.size()); + assertEquals("[b, b, b, b, b]", processed.toString()); // [] assertEquals(0, recovered.size()); assertEquals(1, stepExecution.getReadCount()); diff --git a/spring-batch-core/src/test/java/org/springframework/batch/core/step/item/FaultTolerantStepFactoryBeanRollbackTests.java b/spring-batch-core/src/test/java/org/springframework/batch/core/step/item/FaultTolerantStepFactoryBeanRollbackTests.java index 4dcbea7ab..81e7c9fa2 100644 --- a/spring-batch-core/src/test/java/org/springframework/batch/core/step/item/FaultTolerantStepFactoryBeanRollbackTests.java +++ b/spring-batch-core/src/test/java/org/springframework/batch/core/step/item/FaultTolerantStepFactoryBeanRollbackTests.java @@ -12,6 +12,7 @@ import java.util.List; import org.apache.commons.logging.Log; import org.apache.commons.logging.LogFactory; import org.junit.Before; +import org.junit.Ignore; import org.junit.Test; import org.springframework.batch.core.BatchStatus; import org.springframework.batch.core.JobExecution; @@ -21,6 +22,7 @@ import org.springframework.batch.core.StepExecution; import org.springframework.batch.core.repository.JobRepository; import org.springframework.batch.core.repository.support.MapJobRepositoryFactoryBean; import org.springframework.batch.support.transaction.ResourcelessTransactionManager; +import org.springframework.core.task.SimpleAsyncTaskExecutor; import org.springframework.transaction.interceptor.RollbackRuleAttribute; import org.springframework.transaction.interceptor.RuleBasedTransactionAttribute; import org.springframework.transaction.interceptor.TransactionAttribute; @@ -355,6 +357,24 @@ public class FaultTolerantStepFactoryBeanRollbackTests { .toString()); } + @Test + @Ignore + public void testMultithreadedSkipInWriter() throws Exception { + writer.setFailures("1", "2", "3", "4", "5"); + factory.setCommitInterval(3); + factory.setSkipLimit(10); + factory.setTaskExecutor(new SimpleAsyncTaskExecutor()); + + Step step = (Step) factory.getObject(); + + step.execute(stepExecution); + assertEquals(BatchStatus.COMPLETED, stepExecution.getStatus()); + + assertEquals("[]", writer.getCommitted().toString()); + assertEquals("[]", processor.getCommitted().toString()); + assertEquals(5, stepExecution.getSkipCount()); + } + @Test public void testMultipleSkipsInWriter() throws Exception { writer.setFailures("2", "4"); diff --git a/spring-batch-core/src/test/java/org/springframework/batch/core/step/item/SkipWriterStub.java b/spring-batch-core/src/test/java/org/springframework/batch/core/step/item/SkipWriterStub.java index 7d6a3a64e..06df81f5f 100644 --- a/spring-batch-core/src/test/java/org/springframework/batch/core/step/item/SkipWriterStub.java +++ b/spring-batch-core/src/test/java/org/springframework/batch/core/step/item/SkipWriterStub.java @@ -26,7 +26,7 @@ import org.springframework.batch.support.transaction.TransactionAwareProxyFactor * @since 2.0.1 */ public class SkipWriterStub extends ExceptionThrowingItemHandlerStub implements ItemWriter { - + private List written = new ArrayList(); private List committed = TransactionAwareProxyFactory.createTransactionalList(); @@ -46,6 +46,7 @@ public class SkipWriterStub extends ExceptionThrowingItemHandlerStub imple } public void write(List items) throws Exception { + logger.debug("Writing: "+items); for (T item : items) { written.add(item); checkFailure(item); diff --git a/spring-batch-infrastructure/src/main/java/org/springframework/batch/repeat/support/RepeatTemplate.java b/spring-batch-infrastructure/src/main/java/org/springframework/batch/repeat/support/RepeatTemplate.java index a70f1c168..47855b360 100644 --- a/spring-batch-infrastructure/src/main/java/org/springframework/batch/repeat/support/RepeatTemplate.java +++ b/spring-batch-infrastructure/src/main/java/org/springframework/batch/repeat/support/RepeatTemplate.java @@ -110,8 +110,8 @@ public class RepeatTemplate implements RepeatOperations { /** * Setter for policy to decide when the batch is complete. The default is to - * complete normally when the callback returns a {@link RepeatStatus} which is - * not marked as continuable, and abnormally when the callback throws an + * complete normally when the callback returns a {@link RepeatStatus} which + * is not marked as continuable, and abnormally when the callback throws an * exception (but the decision to re-throw the exception is deferred to the * {@link ExceptionHandler}). * @@ -222,13 +222,15 @@ public class RepeatTemplate implements RepeatOperations { for (int i = listeners.length; i-- > 0;) { RepeatListener interceptor = listeners[i]; - interceptor.onError(context, unwrappedThrowable); // This is not an error - only log at debug // level. logger.debug("Exception intercepted (" + (i + 1) + " of " + listeners.length + ")", unwrappedThrowable); + interceptor.onError(context, unwrappedThrowable); } + logger.debug("Handling exception: " + throwable.getClass().getName() + ", caused by: " + + unwrappedThrowable.getClass().getName() + ": " + unwrappedThrowable.getMessage()); exceptionHandler.handleException(context, unwrappedThrowable); } @@ -261,7 +263,10 @@ public class RepeatTemplate implements RepeatOperations { try { if (!throwables.isEmpty()) { - rethrow((Throwable) throwables.iterator().next()); + Throwable throwable = (Throwable) throwables.iterator().next(); + logger.debug("Handling fatal exception explicitly (rethrowing first of " + throwables.size() + + "): " + throwable.getClass().getName() + ": " + throwable.getMessage()); + rethrow(throwable); } } @@ -358,8 +363,8 @@ public class RepeatTemplate implements RepeatOperations { * processes. By default does nothing and returns true. * * @param state the internal state. - * @return true if {@link #canContinue(RepeatStatus)} is true for all results - * retrieved. + * @return true if {@link #canContinue(RepeatStatus)} is true for all + * results retrieved. */ protected boolean waitForResults(RepeatInternalState state) { // no-op by default diff --git a/spring-batch-infrastructure/src/main/java/org/springframework/batch/retry/support/RetryTemplate.java b/spring-batch-infrastructure/src/main/java/org/springframework/batch/retry/support/RetryTemplate.java index 2a5e13e4c..a92087651 100644 --- a/spring-batch-infrastructure/src/main/java/org/springframework/batch/retry/support/RetryTemplate.java +++ b/spring-batch-infrastructure/src/main/java/org/springframework/batch/retry/support/RetryTemplate.java @@ -195,6 +195,7 @@ public class RetryTemplate implements RetryOperations { // Allow the retry policy to initialise itself... RetryContext context = open(retryPolicy, state); + logger.debug("RetryContext retrieved: "+context); // Make sure the context is available globally for clients who need // it... @@ -248,6 +249,7 @@ public class RetryTemplate implements RetryOperations { throw ex; } + logger.debug("Checking for rethrow: count=" + context.getRetryCount()); if (shouldRethrow(retryPolicy, context, state)) { logger.debug("Rethrow in retry for policy: count=" + context.getRetryCount()); throw e;