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 index 8f7d17dfb..19e9b905d 100644 --- 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 @@ -207,6 +207,21 @@ public class FaultTolerantStepFactoryBean extends SimpleStepFactoryBean 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. @@ -235,20 +250,6 @@ public class FaultTolerantStepFactoryBean extends SimpleStepFactoryBean configureChunkProvider() { - - SkipPolicy readSkipPolicy = new LimitCheckingItemSkipPolicy(skipLimit, skippableExceptionClasses, - fatalExceptionClasses); - FaultTolerantChunkProvider chunkProvider = new FaultTolerantChunkProvider(getItemReader(), - getChunkOperations()); - chunkProvider.setSkipPolicy(readSkipPolicy); - - return chunkProvider; - } /** * @return fully configured retry template for item processing phase. 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 index 80f2fb9f9..51a3bb2f4 100644 --- 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 @@ -459,28 +459,8 @@ public class SimpleStepFactoryBean implements FactoryBean, BeanNameAware { SimpleChunkProcessor chunkProcessor = configureChunkProcessor(); - registerExplicitItemListeners(chunkProvider, chunkProcessor); - registerImplicitItemListeners(chunkProvider, chunkProcessor); - - // Since we are going to wrap these things with listener callbacks we - // need to register them here because the step will not know we did - // that. - List chunkListeners = new ArrayList(Arrays.asList(getListeners())); - for (Object itemHandler : new Object[] { itemReader, itemWriter, itemProcessor }) { - if (itemHandler instanceof ItemStream) { - step.registerStream((ItemStream) itemHandler); - } - if (itemHandler instanceof StepExecutionListener) { - step.registerStepExecutionListener((StepExecutionListener) itemHandler); - } - if (itemHandler instanceof ChunkListener) { - chunkListeners.add((StepListener) itemHandler); - } - } - - BatchListenerFactoryHelper.addChunkListeners(chunkOperations, chunkListeners.toArray(new StepListener[] {})); - step.setStepExecutionListeners(BatchListenerFactoryHelper.getListeners(listeners, StepExecutionListener.class) - .toArray(new StepExecutionListener[] {})); + registerItemListeners(chunkProvider, chunkProcessor); + registerStepListeners(step, chunkOperations); ChunkOrientedTasklet tasklet = new ChunkOrientedTasklet(chunkProvider, chunkProcessor); tasklet.setBuffering(!isReaderTransactionalQueue()); @@ -489,10 +469,20 @@ public class SimpleStepFactoryBean implements FactoryBean, BeanNameAware { } + /** + * Extension point for creating appropriate {@link ChunkProvider}. Return + * value must subclass {@link SimpleChunkProvider} due to listener + * registration. + */ protected SimpleChunkProvider configureChunkProvider() { return new SimpleChunkProvider(itemReader, chunkOperations); } + /** + * Extension point for creating appropriate {@link ChunkProcessor}. Return + * value must subclass {@link SimpleChunkProcessor} due to listener + * registration. + */ protected SimpleChunkProcessor configureChunkProcessor() { return new SimpleChunkProcessor(itemProcessor, itemWriter); } @@ -517,12 +507,35 @@ public class SimpleStepFactoryBean implements FactoryBean, BeanNameAware { } /** - * Register explicitly set ({@link #setListeners(StepListener[])}) item - * listeners. + * Register listeners with step and chunk. */ - private void registerExplicitItemListeners(SimpleChunkProvider chunkProvider, - SimpleChunkProcessor chunkProcessor) { + private void registerStepListeners(TaskletStep step, RepeatOperations chunkOperations) { + List chunkListeners = new ArrayList(Arrays.asList(getListeners())); + for (Object itemHandler : new Object[] { itemReader, itemWriter, itemProcessor }) { + if (itemHandler instanceof ItemStream) { + step.registerStream((ItemStream) itemHandler); + } + if (itemHandler instanceof StepExecutionListener) { + step.registerStepExecutionListener((StepExecutionListener) itemHandler); + } + if (itemHandler instanceof ChunkListener) { + chunkListeners.add((StepListener) itemHandler); + } + } + + BatchListenerFactoryHelper.addChunkListeners(chunkOperations, chunkListeners.toArray(new StepListener[] {})); + step.setStepExecutionListeners(BatchListenerFactoryHelper.getListeners(listeners, StepExecutionListener.class) + .toArray(new StepExecutionListener[] {})); + } + + /** + * Register explicitly set ({@link #setListeners(StepListener[])}) item + * listeners and auto-register reader, processor and writer if applicable + */ + private void registerItemListeners(SimpleChunkProvider chunkProvider, SimpleChunkProcessor chunkProcessor) { + + // explicitly set item listeners chunkProvider.setListeners(BatchListenerFactoryHelper.> getListeners(getListeners(), ItemReadListener.class)); chunkProvider.setListeners(BatchListenerFactoryHelper.> getListeners(getListeners(), @@ -534,14 +547,8 @@ public class SimpleStepFactoryBean implements FactoryBean, BeanNameAware { ItemWriteListener.class)); chunkProcessor.setListeners(BatchListenerFactoryHelper.> getListeners(getListeners(), SkipListener.class)); - } - /** - * Auto-register reader, processor and writer as item listeners if - * applicable. - */ - private void registerImplicitItemListeners(SimpleChunkProvider chunkProvider, - SimpleChunkProcessor chunkProcessor) { + // auto-register reader, processor and writer for (Object itemHandler : new Object[] { getItemReader(), getItemWriter(), getItemProcessor() }) { if (itemHandler instanceof SkipListener) {