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 old mode 100644 new mode 100755 index 8e19aa197..c74e1c562 --- 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 @@ -16,7 +16,9 @@ package org.springframework.batch.core.step.item; +import java.util.ArrayList; import java.util.Collections; +import java.util.List; import org.apache.commons.logging.Log; import org.apache.commons.logging.LogFactory; @@ -25,6 +27,7 @@ import org.springframework.batch.classify.Classifier; import org.springframework.batch.core.StepContribution; import org.springframework.batch.core.step.skip.LimitCheckingItemSkipPolicy; import org.springframework.batch.core.step.skip.NonSkippableProcessException; +import org.springframework.batch.core.step.skip.SkipListenerFailedException; import org.springframework.batch.core.step.skip.SkipPolicy; import org.springframework.batch.item.ItemProcessor; import org.springframework.batch.item.ItemWriter; @@ -48,18 +51,60 @@ public class FaultTolerantChunkProcessor extends SimpleChunkProcessor rollbackClassifier) { this.rollbackClassifier = rollbackClassifier; } + /** + * @param chunkMonitor + */ + public void setChunkMonitor(ChunkMonitor chunkMonitor) { + this.chunkMonitor = chunkMonitor; + } + + /** + * A flag to indicate that items have been buffered and therefore will + * always come back as a chunk after a rollback. Otherwise things are more + * complicated because after a rollback the new chunk might or moght not + * contain items from the previous failed chunk. + * + * @param buffering + */ public void setBuffering(boolean buffering) { this.buffering = buffering; } @@ -133,9 +178,8 @@ public class FaultTolerantChunkProcessor extends SimpleChunkProcessor extends SimpleChunkProcessor retryCallback = new RetryCallback() { public Object doWithRetry(RetryContext context) throws Exception { - doWrite(outputs.getItems()); - contribution.incrementWriteCount(outputs.size()); + + if (!inputs.isBusy()) { + chunkMonitor.setChunkSize(inputs.size()); + doWrite(outputs.getItems()); + contribution.incrementWriteCount(outputs.size()); + } + else { + scan(contribution, inputs, outputs, chunkMonitor); + } return null; + } }; - RecoveryCallback recoveryCallback = new RecoveryCallback() { + if (!buffering) { - public Object recover(RetryContext context) throws Exception { + RecoveryCallback batchRecoveryCallback = new RecoveryCallback() { - Exception le = (Exception) context.getLastThrowable(); - if (outputs.size() > 1 && !rollbackClassifier.classify(le)) { - throw new RetryException("Invalid retry state during write caused by " - + "exception that does not classify for rollback: ", le); - } + public Object recover(RetryContext context) throws Exception { - boolean singleton = outputs.size() == 1; - - Chunk.ChunkIterator inputIterator = inputs.iterator(); - for (Chunk.ChunkIterator outputIterator = outputs.iterator(); outputIterator.hasNext();) { - - inputIterator.next(); - O item = outputIterator.next(); - if (singleton) { - checkSkipPolicy(inputIterator, outputIterator, le, contribution); - return null; + Exception e = (Exception) context.getLastThrowable(); + if (outputs.size() > 1 && !rollbackClassifier.classify(e)) { + throw new RetryException("Invalid retry state during write caused by " + + "exception that does not classify for rollback: ", e); } - try { - writeItems(Collections.singletonList(item)); - } - catch (Exception e) { + Chunk.ChunkIterator inputIterator = inputs.iterator(); + for (Chunk.ChunkIterator outputIterator = outputs.iterator(); outputIterator.hasNext();) { + + inputIterator.next(); + outputIterator.next(); + checkSkipPolicy(inputIterator, outputIterator, e, contribution); - if (rollbackClassifier.classify(e)) { - throw e; - } - else { + if (!rollbackClassifier.classify(e)) { throw new RetryException( "Invalid retry state during recovery caused by exception that does not classify for rollback: ", e); } + } - } - - doAfterWrite(outputs.getItems()); - contribution.incrementWriteCount(outputs.size()); - return null; - - } - - }; - - RecoveryCallback batchRecoveryCallback = new RecoveryCallback() { - - public Object recover(RetryContext context) throws Exception { - - Exception e = (Exception) context.getLastThrowable(); - if (outputs.size() > 1 && !rollbackClassifier.classify(e)) { - throw new RetryException("Invalid retry state during write caused by " - + "exception that does not classify for rollback: ", e); - } - - Chunk.ChunkIterator inputIterator = inputs.iterator(); - for (Chunk.ChunkIterator outputIterator = outputs.iterator(); outputIterator.hasNext();) { - - inputIterator.next(); - outputIterator.next(); - - checkSkipPolicy(inputIterator, outputIterator, e, contribution); - if (!rollbackClassifier.classify(e)) { - throw new RetryException( - "Invalid retry state during recovery caused by exception that does not classify for rollback: ", - e); - } + return null; } - return null; + }; - } + batchRetryTemplate.execute(retryCallback, batchRecoveryCallback, BatchRetryTemplate.createState( + getInputKeys(inputs), rollbackClassifier)); - }; - - if (!buffering) { - batchRetryTemplate.execute(retryCallback, batchRecoveryCallback, BatchRetryTemplate.createState(inputs - .getItems(), rollbackClassifier)); } else { + + RecoveryCallback recoveryCallback = new RecoveryCallback() { + + public Object recover(RetryContext context) throws Exception { + + Exception le = (Exception) context.getLastThrowable(); + if (outputs.size() > 1 && !rollbackClassifier.classify(le)) { + throw new RetryException("Invalid retry state during write caused by " + + "exception that does not classify for rollback: ", le); + } + + boolean singleton = outputs.size() == 1; + + if (singleton && !inputs.isBusy()) { + 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; + + } + + }; + batchRetryTemplate.execute(retryCallback, recoveryCallback, new DefaultRetryState(inputs, rollbackClassifier)); + + } + + callSkipListeners(inputs, outputs); + + } + + private void callSkipListeners(final Chunk inputs, final Chunk outputs) { + + for (SkipWrapper wrapper : inputs.getSkips()) { + I item = wrapper.getItem(); + if (item == null) { + continue; + } + Exception e = wrapper.getException(); + try { + getListener().onSkipInProcess(item, e); + } + catch (RuntimeException ex) { + throw new SkipListenerFailedException("Fatal exception in SkipListener.", ex, e); + } } + for (SkipWrapper wrapper : outputs.getSkips()) { + Exception e = wrapper.getException(); + try { + getListener().onSkipInWrite(wrapper.getItem(), e); + } + catch (RuntimeException ex) { + throw new SkipListenerFailedException("Fatal exception in SkipListener.", ex, e); + } + } + + // Clear skips if we are possibly going to process this chunk again + outputs.clearSkips(); + inputs.clearSkips(); + + } + + private Object getInputKey(I item) { + if (keyGenerator == null) { + return item; + } + return keyGenerator.getKey(item); + } + + private List getInputKeys(final Chunk inputs) { + if (keyGenerator == null) { + return inputs.getItems(); + } + List keys = new ArrayList(); + for (I item : inputs.getItems()) { + keys.add(keyGenerator.getKey(item)); + } + return keys; } private void checkSkipPolicy(Chunk.ChunkIterator inputIterator, Chunk.ChunkIterator outputIterator, @@ -260,4 +348,43 @@ public class FaultTolerantChunkProcessor extends SimpleChunkProcessor inputs, final Chunk outputs, ChunkMonitor chunkMonitor) + throws Exception { + + if (outputs.isEmpty()) { + inputs.setBusy(false); + return; + } + + Chunk.ChunkIterator inputIterator = inputs.iterator(); + Chunk.ChunkIterator outputIterator = outputs.iterator(); + + List items = Collections.singletonList(outputIterator.next()); + try { + writeItems(items); + } + catch (Exception e) { + checkSkipPolicy(inputIterator, outputIterator, e, contribution); + if (rollbackClassifier.classify(e)) { + throw e; + } + else { + throw new RetryException( + "Invalid retry state during recovery caused by exception that does not classify for rollback: ", + e); + } + } + // If successful we are going to return and allow + // the driver to commit... + doAfterWrite(items); + contribution.incrementWriteCount(1); + inputIterator.remove(); + outputIterator.remove(); + chunkMonitor.incrementOffset(); + if (outputs.isEmpty()) { + inputs.setBusy(false); + chunkMonitor.resetOffset(); + } + } + } diff --git a/spring-batch-core/src/main/java/org/springframework/batch/core/step/item/FaultTolerantChunkProvider.java b/spring-batch-core/src/main/java/org/springframework/batch/core/step/item/FaultTolerantChunkProvider.java old mode 100644 new mode 100755 index 85b4e3ca9..9ad6430f7 --- a/spring-batch-core/src/main/java/org/springframework/batch/core/step/item/FaultTolerantChunkProvider.java +++ b/spring-batch-core/src/main/java/org/springframework/batch/core/step/item/FaultTolerantChunkProvider.java @@ -19,6 +19,7 @@ package org.springframework.batch.core.step.item; import org.springframework.batch.core.StepContribution; import org.springframework.batch.core.step.skip.LimitCheckingItemSkipPolicy; import org.springframework.batch.core.step.skip.NonSkippableReadException; +import org.springframework.batch.core.step.skip.SkipListenerFailedException; import org.springframework.batch.core.step.skip.SkipPolicy; import org.springframework.batch.item.ItemReader; import org.springframework.batch.repeat.RepeatOperations; @@ -31,6 +32,10 @@ public class FaultTolerantChunkProvider extends SimpleChunkProvider { super(itemReader, repeatOperations); } + /** + * The policy that determines whether exceptions can be skipped on read. + * @param SkipPolicy + */ public void setSkipPolicy(SkipPolicy SkipPolicy) { this.skipPolicy = SkipPolicy; } @@ -58,4 +63,16 @@ public class FaultTolerantChunkProvider extends SimpleChunkProvider { } } + @Override + public void postProcess(StepContribution contribution, Chunk chunk) { + for (Exception e : chunk.getErrors()) { + try { + getListener().onSkipInRead(e); + } + catch (RuntimeException ex) { + throw new SkipListenerFailedException("Fatal exception in SkipListener.", ex, e); + } + } + } + } diff --git a/spring-batch-core/src/main/java/org/springframework/batch/core/step/item/FaultTolerantStepFactoryBean.java b/spring-batch-core/src/main/java/org/springframework/batch/core/step/item/FaultTolerantStepFactoryBean.java old mode 100644 new mode 100755 index 3ac9cc4d1..240ad6c16 --- a/spring-batch-core/src/main/java/org/springframework/batch/core/step/item/FaultTolerantStepFactoryBean.java +++ b/spring-batch-core/src/main/java/org/springframework/batch/core/step/item/FaultTolerantStepFactoryBean.java @@ -29,6 +29,9 @@ import org.springframework.batch.core.step.skip.SkipLimitExceededException; import org.springframework.batch.core.step.skip.SkipListenerFailedException; import org.springframework.batch.core.step.skip.SkipPolicy; import org.springframework.batch.core.step.tasklet.TaskletStep; +import org.springframework.batch.item.ItemReader; +import org.springframework.batch.item.ItemStream; +import org.springframework.batch.item.support.CompositeItemStream; import org.springframework.batch.repeat.RepeatOperations; import org.springframework.batch.repeat.support.RepeatTemplate; import org.springframework.batch.retry.RetryException; @@ -40,6 +43,8 @@ import org.springframework.batch.retry.policy.MapRetryContextCache; import org.springframework.batch.retry.policy.NeverRetryPolicy; import org.springframework.batch.retry.policy.RetryContextCache; import org.springframework.batch.retry.policy.SimpleRetryPolicy; +import org.springframework.core.task.SyncTaskExecutor; +import org.springframework.core.task.TaskExecutor; /** * Factory bean for step that provides options for configuring skip behaviour. @@ -88,6 +93,22 @@ public class FaultTolerantStepFactoryBean extends SimpleStepFactoryBean extends SimpleStepFactoryBean configureChunkProvider() { - + protected SimpleChunkProvider configureChunkProvider() { + SkipPolicy readSkipPolicy = new LimitCheckingItemSkipPolicy(skipLimit, skippableExceptionClasses, fatalExceptionClasses); FaultTolerantChunkProvider chunkProvider = new FaultTolerantChunkProvider(getItemReader(), getChunkOperations()); chunkProvider.setSkipPolicy(readSkipPolicy); - + return chunkProvider; + } /** * @return {@link ChunkProcessor} configured for fault-tolerance. */ @Override - protected FaultTolerantChunkProcessor configureChunkProcessor() { + protected SimpleChunkProcessor configureChunkProcessor() { SkipPolicy writeSkipPolicy = new LimitCheckingItemSkipPolicy(skipLimit, skippableExceptionClasses, fatalExceptionClasses); @@ -246,10 +301,12 @@ public class FaultTolerantStepFactoryBean extends SimpleStepFactoryBean implements ChunkProcessor, Initializi this.listener.register(listener); } + /** + * @return the listener + */ + protected MulticasterBatchListener getListener() { + return listener; + } + /** * @param item the input item * @return the result of the processing @@ -122,21 +128,20 @@ public class SimpleChunkProcessor implements ChunkProcessor, Initializi try { listener.beforeWrite(items); writeItems(items); - listener.afterWrite(items); + doAfterWrite(items); } catch (Exception e) { listener.onWriteError(e, items); throw e; } } - + /** * Call the listener's after write method. * * @param items */ - protected final void doAfterWrite(List items) - { + protected final void doAfterWrite(List items) { listener.afterWrite(items); } @@ -146,17 +151,6 @@ public class SimpleChunkProcessor implements ChunkProcessor, Initializi public final void process(StepContribution contribution, Chunk inputs) throws Exception { - // If there is no input we don't have to do anything more - if (inputs.isEmpty()) { - return; - } - - int inputSize = inputs.size(); - - Chunk outputs = transform(contribution, inputs); - - contribution.incrementFilterCount(inputSize - outputs.size()); - /* * Need to remember the write skips across transactions, otherwise they * keep coming back. Since we register skips with the inputs they will @@ -172,37 +166,37 @@ public class SimpleChunkProcessor implements ChunkProcessor, Initializi skips = new Chunk(); } + // If there is no input we don't have to do anything more + if (inputs.isEmpty() && skips.getSkips().isEmpty()) { + return; + } + + int inputsSize = inputs.size(); + + Chunk outputs = transform(contribution, inputs); + + contribution.incrementFilterCount(inputsSize - outputs.size()); + outputs = new Chunk(outputs.getItems(), skips.getSkips()); + + // Remember for next time if there are skips accumulating inputs.setUserData(outputs); write(contribution, inputs, outputs); - for (SkipWrapper wrapper : inputs.getSkips()) { - I item = wrapper.getItem(); - if (item == null) { - continue; - } - Exception e = wrapper.getException(); - try { - listener.onSkipInProcess(item, e); - } - catch (RuntimeException ex) { - throw new SkipListenerFailedException("Fatal exception in SkipListener.", ex, e); - } - } - - for (SkipWrapper wrapper : outputs.getSkips()) { - Exception e = wrapper.getException(); - try { - listener.onSkipInWrite(wrapper.getItem(), e); - } - catch (RuntimeException ex) { - throw new SkipListenerFailedException("Fatal exception in SkipListener.", ex, e); - } - } - } + /** + * Simple implementation delegates to the {@link #doWrite(List)} method and + * increments the write count in the contribution. Subclasses can handle + * more complicated scenarios, e.g.with fault tolerance. If output items are + * skipped they should be removed from the inputs as well. + * + * @param contribution the current step contribution + * @param inputs the inputs that gave rise to the ouputs + * @param outputs the outputs to write + * @throws Exception if there is a problem + */ protected void write(StepContribution contribution, Chunk inputs, Chunk outputs) throws Exception { doWrite(outputs.getItems()); contribution.incrementWriteCount(outputs.size()); diff --git a/spring-batch-core/src/main/java/org/springframework/batch/core/step/item/SimpleChunkProvider.java b/spring-batch-core/src/main/java/org/springframework/batch/core/step/item/SimpleChunkProvider.java old mode 100644 new mode 100755 index d6dd035dd..c7e79a509 --- a/spring-batch-core/src/main/java/org/springframework/batch/core/step/item/SimpleChunkProvider.java +++ b/spring-batch-core/src/main/java/org/springframework/batch/core/step/item/SimpleChunkProvider.java @@ -23,7 +23,6 @@ import org.apache.commons.logging.LogFactory; import org.springframework.batch.core.StepContribution; import org.springframework.batch.core.StepListener; import org.springframework.batch.core.listener.MulticasterBatchListener; -import org.springframework.batch.core.step.skip.SkipListenerFailedException; import org.springframework.batch.item.ItemReader; import org.springframework.batch.repeat.RepeatCallback; import org.springframework.batch.repeat.RepeatContext; @@ -71,6 +70,13 @@ public class SimpleChunkProvider implements ChunkProvider { public void registerListener(StepListener listener) { this.listener.register(listener); } + + /** + * @return the listener + */ + protected MulticasterBatchListener getListener() { + return listener; + } /** * Surrounds the read call with listener callbacks. @@ -113,14 +119,7 @@ public class SimpleChunkProvider implements ChunkProvider { } public void postProcess(StepContribution contribution, Chunk chunk) { - for (Exception e : chunk.getErrors()) { - try { - listener.onSkipInRead(e); - } - catch (RuntimeException ex) { - throw new SkipListenerFailedException("Fatal exception in SkipListener.", ex, e); - } - } + // do nothing } protected I read(StepContribution contribution, Chunk chunk) throws Exception { diff --git a/spring-batch-core/src/main/java/org/springframework/batch/core/step/item/SimpleStepFactoryBean.java b/spring-batch-core/src/main/java/org/springframework/batch/core/step/item/SimpleStepFactoryBean.java old mode 100644 new mode 100755 index d87466406..baa88b56f --- a/spring-batch-core/src/main/java/org/springframework/batch/core/step/item/SimpleStepFactoryBean.java +++ b/spring-batch-core/src/main/java/org/springframework/batch/core/step/item/SimpleStepFactoryBean.java @@ -407,6 +407,14 @@ public class SimpleStepFactoryBean implements FactoryBean, BeanNameAware { public void setTaskExecutor(TaskExecutor taskExecutor) { this.taskExecutor = taskExecutor; } + + /** + * Mkae the {@link TaskExecutor} available to subclasses + * @return the taskExecutor to be used to execute chunks + */ + protected TaskExecutor getTaskExecutor() { + return taskExecutor; + } /** * Public setter for the throttle limit. This limits the number of tasks @@ -437,7 +445,7 @@ public class SimpleStepFactoryBean implements FactoryBean, BeanNameAware { step.setStartLimit(startLimit); step.setAllowStartIfComplete(allowStartIfComplete); - step.setStreams(streams); + registerStreams(step, streams); if (chunkOperations == null) { RepeatTemplate repeatTemplate = new RepeatTemplate(); @@ -468,6 +476,7 @@ public class SimpleStepFactoryBean implements FactoryBean, BeanNameAware { registerItemListeners(chunkProvider, chunkProcessor); registerStepListeners(step, chunkOperations); + registerStreams(step, itemReader, itemProcessor, itemWriter); ChunkOrientedTasklet tasklet = new ChunkOrientedTasklet(chunkProvider, chunkProcessor); tasklet.setBuffering(!isReaderTransactionalQueue()); @@ -476,6 +485,15 @@ public class SimpleStepFactoryBean implements FactoryBean, BeanNameAware { } + /** + * Register the streams with the step. + * @param step the {@link TaskletStep} + * @param streams the streams to register + */ + protected void registerStreams(TaskletStep step, ItemStream[] streams) { + step.setStreams(streams); + } + /** * Extension point for creating appropriate {@link ChunkProvider}. Return * value must subclass {@link SimpleChunkProvider} due to listener @@ -512,16 +530,21 @@ public class SimpleStepFactoryBean implements FactoryBean, BeanNameAware { } return new SimpleCompletionPolicy(commitInterval); } + + private void registerStreams(TaskletStep step, ItemReader itemReader, ItemProcessor itemProcessor, ItemWriter itemWriter) { + for (Object itemHandler : new Object[] { itemReader, itemWriter, itemProcessor }) { + if (itemHandler instanceof ItemStream) { + registerStreams(step, new ItemStream[] {(ItemStream) itemHandler}); + } + } + } /** * Register listeners with step and chunk. */ private void registerStepListeners(TaskletStep step, RepeatOperations chunkOperations) { - for (Object itemHandler : new Object[] { itemReader, itemWriter, itemProcessor }) { - if (itemHandler instanceof ItemStream) { - step.registerStream((ItemStream) itemHandler); - } + for (Object itemHandler : new Object[] { getItemReader(), itemWriter, itemProcessor }) { if (StepListenerFactoryBean.isListener(itemHandler)) { StepListener listener = StepListenerFactoryBean.getListener(itemHandler); if (listener instanceof StepExecutionListener) {