IN PROGRESS - BATCH-896: "DRY" FaultTolerantTasklet implementations
pulled up readSkipPolicy and attribute keys
This commit is contained in:
@@ -36,6 +36,12 @@ import org.springframework.core.AttributeAccessor;
|
||||
* @author Robert Kasanicky
|
||||
*/
|
||||
public abstract class AbstractFaultTolerantChunkOrientedTasklet<I, O> extends AbstractItemOrientedTasklet<I, O> {
|
||||
|
||||
final static protected String SKIPPED_INPUTS_KEY = "SKIPPED_INPUTS_KEY";
|
||||
|
||||
final static protected String SKIPPED_OUTPUTS_KEY = "SKIPPED_OUTPUTS_KEY";
|
||||
|
||||
final static protected String SKIPPED_READS_KEY = "SKIPPED_READS_KEY";
|
||||
|
||||
final private RetryOperations retryOperations;
|
||||
|
||||
@@ -43,20 +49,27 @@ public abstract class AbstractFaultTolerantChunkOrientedTasklet<I, O> extends Ab
|
||||
|
||||
final private ItemSkipPolicy processSkipPolicy;
|
||||
|
||||
final private ItemSkipPolicy readSkipPolicy;
|
||||
|
||||
final private Classifier<Throwable, Boolean> rollbackClassifier;
|
||||
|
||||
public AbstractFaultTolerantChunkOrientedTasklet(ItemReader<? extends I> itemReader,
|
||||
ItemProcessor<? super I, ? extends O> itemProcessor, ItemWriter<? super O> itemWriter,
|
||||
RetryOperations retryOperations, ItemSkipPolicy processSkipPolicy, ItemSkipPolicy writeSkipPolicy,
|
||||
Classifier<Throwable, Boolean> rollbackClassifier) {
|
||||
RetryOperations retryOperations, ItemSkipPolicy readSkipPolicy, ItemSkipPolicy processSkipPolicy,
|
||||
ItemSkipPolicy writeSkipPolicy, Classifier<Throwable, Boolean> rollbackClassifier) {
|
||||
|
||||
super(itemReader, itemProcessor, itemWriter);
|
||||
this.retryOperations = retryOperations;
|
||||
this.readSkipPolicy = readSkipPolicy;
|
||||
this.processSkipPolicy = processSkipPolicy;
|
||||
this.writeSkipPolicy = writeSkipPolicy;
|
||||
this.rollbackClassifier = rollbackClassifier;
|
||||
}
|
||||
|
||||
protected ItemSkipPolicy getReadSkipPolicy() {
|
||||
return readSkipPolicy;
|
||||
}
|
||||
|
||||
/**
|
||||
* Call all skip listeners in read-process-write order
|
||||
* @param skippedReads read exceptions
|
||||
|
||||
@@ -42,7 +42,7 @@ import org.springframework.core.AttributeAccessor;
|
||||
* listener is invoked and the skip count incremented. A retryable exception is
|
||||
* thus also effectively also implicitly skippable.
|
||||
*
|
||||
* ItemProcessor is assumed to be transactional. In case of rollback caused by
|
||||
* <code>ItemProcessor</code> is assumed to be transactional. In case of rollback caused by
|
||||
* error on write the processing phase will be repeated.
|
||||
*
|
||||
* @author Dave Syer
|
||||
@@ -50,17 +50,9 @@ import org.springframework.core.AttributeAccessor;
|
||||
*/
|
||||
public class FaultTolerantChunkOrientedTasklet<I, S> extends AbstractFaultTolerantChunkOrientedTasklet<I, S> {
|
||||
|
||||
private static final String INPUT_BUFFER_KEY = "INPUT_BUFFER_KEY";
|
||||
final static private String INPUT_BUFFER_KEY = "INPUT_BUFFER_KEY";
|
||||
|
||||
private final RepeatOperations repeatOperations;
|
||||
|
||||
final private ItemSkipPolicy readSkipPolicy;
|
||||
|
||||
private static final String SKIPPED_OUTPUTS_KEY = "SKIPPED_OUTPUTS_BUFFER_KEY";
|
||||
|
||||
private static final String SKIPPED_INPUTS_KEY = "SKIPPED_INPUTS_BUFFER_KEY";
|
||||
|
||||
private static final String SKIPPED_READS_KEY = "SKIPPED_READS_BUFFER_KEY";
|
||||
final private RepeatOperations repeatOperations;
|
||||
|
||||
public FaultTolerantChunkOrientedTasklet(ItemReader<? extends I> itemReader,
|
||||
ItemProcessor<? super I, ? extends S> itemProcessor, ItemWriter<? super S> itemWriter,
|
||||
@@ -68,10 +60,9 @@ public class FaultTolerantChunkOrientedTasklet<I, S> extends AbstractFaultTolera
|
||||
Classifier<Throwable, Boolean> rollbackClassifier, ItemSkipPolicy readSkipPolicy,
|
||||
ItemSkipPolicy writeSkipPolicy, ItemSkipPolicy processSkipPolicy) {
|
||||
|
||||
super(itemReader, itemProcessor, itemWriter, retryTemplate, processSkipPolicy, writeSkipPolicy,
|
||||
super(itemReader, itemProcessor, itemWriter, retryTemplate, readSkipPolicy, processSkipPolicy, writeSkipPolicy,
|
||||
rollbackClassifier);
|
||||
this.repeatOperations = chunkOperations;
|
||||
this.readSkipPolicy = readSkipPolicy;
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -156,7 +147,7 @@ public class FaultTolerantChunkOrientedTasklet<I, S> extends AbstractFaultTolera
|
||||
}
|
||||
catch (Exception e) {
|
||||
|
||||
if (readSkipPolicy.shouldSkip(e, contribution.getStepSkipCount())) {
|
||||
if (getReadSkipPolicy().shouldSkip(e, contribution.getStepSkipCount())) {
|
||||
// increment skip count and try again
|
||||
contribution.incrementReadSkipCount();
|
||||
skippedReads.add(e);
|
||||
|
||||
@@ -34,15 +34,7 @@ import org.springframework.core.AttributeAccessor;
|
||||
public class NonbufferingFaultTolerantChunkOrientedTasklet<I, O> extends
|
||||
AbstractFaultTolerantChunkOrientedTasklet<I, O> {
|
||||
|
||||
private static final String SKIPPED_INPUTS_KEY = "SKIPPED_INPUTS_KEY";
|
||||
|
||||
private static final String SKIPPED_OUTPUTS_KEY = "SKIPPED_OUTPUTS_KEY";
|
||||
|
||||
private static final String SKIPPED_READS_KEY = "SKIPPED_READS_KEY";
|
||||
|
||||
private final RepeatOperations repeatOperations;
|
||||
|
||||
private final ItemSkipPolicy readSkipPolicy;
|
||||
final private RepeatOperations repeatOperations;
|
||||
|
||||
public NonbufferingFaultTolerantChunkOrientedTasklet(ItemReader<? extends I> itemReader,
|
||||
ItemProcessor<? super I, ? extends O> itemProcessor, ItemWriter<? super O> itemWriter,
|
||||
@@ -50,10 +42,9 @@ public class NonbufferingFaultTolerantChunkOrientedTasklet<I, O> extends
|
||||
Classifier<Throwable, Boolean> rollbackClassifier, ItemSkipPolicy readSkipPolicy,
|
||||
ItemSkipPolicy writeSkipPolicy, ItemSkipPolicy processSkipPolicy) {
|
||||
|
||||
super(itemReader, itemProcessor, itemWriter, retryTemplate, processSkipPolicy, writeSkipPolicy,
|
||||
super(itemReader, itemProcessor, itemWriter, retryTemplate, readSkipPolicy, processSkipPolicy, writeSkipPolicy,
|
||||
rollbackClassifier);
|
||||
this.repeatOperations = chunkOperations;
|
||||
this.readSkipPolicy = readSkipPolicy;
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -118,7 +109,7 @@ public class NonbufferingFaultTolerantChunkOrientedTasklet<I, O> extends
|
||||
}
|
||||
catch (Exception e) {
|
||||
|
||||
if (readSkipPolicy.shouldSkip(e, contribution.getStepSkipCount())) {
|
||||
if (getReadSkipPolicy().shouldSkip(e, contribution.getStepSkipCount())) {
|
||||
// increment skip count and try again
|
||||
contribution.incrementReadSkipCount();
|
||||
skipped.add(e);
|
||||
|
||||
@@ -161,7 +161,7 @@ public class FaultTolerantChunkOrientedTaskletTests {
|
||||
catch (Exception e) {
|
||||
assertEquals("Barf!", e.getMessage());
|
||||
}
|
||||
assertTrue(attributes.hasAttribute("SKIPPED_OUTPUTS_BUFFER_KEY"));
|
||||
assertTrue(attributes.hasAttribute("SKIPPED_OUTPUTS_KEY"));
|
||||
tasklet.execute(contribution, attributes);
|
||||
assertEquals(1, contribution.getReadCount());
|
||||
assertEquals(1, contribution.getWriteSkipCount());
|
||||
@@ -191,10 +191,10 @@ public class FaultTolerantChunkOrientedTaskletTests {
|
||||
catch (Exception e) {
|
||||
assertEquals("Barf!", e.getMessage());
|
||||
}
|
||||
assertTrue(attributes.hasAttribute("SKIPPED_OUTPUTS_BUFFER_KEY"));
|
||||
assertTrue(attributes.hasAttribute("SKIPPED_OUTPUTS_KEY"));
|
||||
}
|
||||
@SuppressWarnings("unchecked")
|
||||
Map<String, Exception> skips = (Map<String, Exception>) attributes.getAttribute("SKIPPED_OUTPUTS_BUFFER_KEY");
|
||||
Map<String, Exception> skips = (Map<String, Exception>) attributes.getAttribute("SKIPPED_OUTPUTS_KEY");
|
||||
assertEquals(1, skips.size());
|
||||
// The last recovery for this chunk...
|
||||
tasklet.execute(contribution, attributes);
|
||||
@@ -214,7 +214,7 @@ public class FaultTolerantChunkOrientedTaskletTests {
|
||||
catch (SkipLimitExceededException e) {
|
||||
// expected
|
||||
}
|
||||
assertTrue(attributes.hasAttribute("SKIPPED_OUTPUTS_BUFFER_KEY"));
|
||||
assertTrue(attributes.hasAttribute("SKIPPED_OUTPUTS_KEY"));
|
||||
assertEquals(3, contribution.getReadCount());
|
||||
assertEquals(0, contribution.getFilterCount());
|
||||
assertEquals(2, contribution.getWriteSkipCount());
|
||||
@@ -285,7 +285,7 @@ public class FaultTolerantChunkOrientedTaskletTests {
|
||||
assertTrue(attributes.hasAttribute("INPUT_BUFFER_KEY"));
|
||||
}
|
||||
@SuppressWarnings("unchecked")
|
||||
Map<Integer, Exception > skips = (Map<Integer, Exception>) attributes.getAttribute("SKIPPED_INPUTS_BUFFER_KEY");
|
||||
Map<Integer, Exception > skips = (Map<Integer, Exception>) attributes.getAttribute("SKIPPED_INPUTS_KEY");
|
||||
assertEquals(1, skips.size());
|
||||
|
||||
// The last recovery for this chunk...
|
||||
|
||||
Reference in New Issue
Block a user