Refactored setting of completion policy to not overwrite JSR defaults

This commit is contained in:
Michael Minella
2013-12-17 14:30:36 -06:00
parent 2f63bd0664
commit 3e5597717e
3 changed files with 49 additions and 24 deletions

View File

@@ -80,9 +80,7 @@ import org.springframework.batch.item.ItemReader;
import org.springframework.batch.item.ItemStream;
import org.springframework.batch.item.ItemWriter;
import org.springframework.batch.repeat.CompletionPolicy;
import org.springframework.batch.repeat.policy.CompositeCompletionPolicy;
import org.springframework.batch.repeat.policy.SimpleCompletionPolicy;
import org.springframework.batch.repeat.policy.TimeoutTerminationPolicy;
import org.springframework.batch.repeat.support.TaskExecutorRepeatTemplate;
import org.springframework.beans.factory.BeanNameAware;
import org.springframework.beans.factory.FactoryBean;
@@ -215,8 +213,6 @@ public class StepParserStepFactoryBean<I, O> implements FactoryBean, BeanNameAwa
private ItemWriter<? super O> itemWriter;
private Integer timeout;
//
// Chunk Elements
//
@@ -422,7 +418,7 @@ public class StepParserStepFactoryBean<I, O> implements FactoryBean, BeanNameAwa
return new FaultTolerantStepBuilder<I, O>(new StepBuilder(stepName));
}
private void registerItemListeners(SimpleStepBuilder<I, O> builder) {
protected void registerItemListeners(SimpleStepBuilder<I, O> builder) {
for (ItemReadListener<I> listener : readListeners) {
builder.listener(listener);
}
@@ -436,22 +432,10 @@ public class StepParserStepFactoryBean<I, O> implements FactoryBean, BeanNameAwa
@SuppressWarnings("unchecked")
protected Step createSimpleStep() {
SimpleStepBuilder builder = getSimpleStepBuilder(this.name);
SimpleStepBuilder builder = new SimpleStepBuilder(new StepBuilder(name));
if(timeout != null && commitInterval != null) {
CompositeCompletionPolicy completionPolicy = new CompositeCompletionPolicy();
CompletionPolicy [] policies = new CompletionPolicy[2];
policies[0] = new SimpleCompletionPolicy(commitInterval);
policies[1] = new TimeoutTerminationPolicy(timeout * 1000);
completionPolicy.setPolicies(policies);
builder.chunk(completionPolicy);
} else if(timeout != null) {
builder.chunk(new TimeoutTerminationPolicy(timeout * 1000));
} else if(commitInterval != null) {
builder.chunk(commitInterval);
}
setChunk(builder);
builder.chunk(chunkCompletionPolicy);
enhanceTaskletStepBuilder(builder);
registerItemListeners(builder);
builder.reader(itemReader);
@@ -460,6 +444,17 @@ public class StepParserStepFactoryBean<I, O> implements FactoryBean, BeanNameAwa
return builder.build();
}
protected void setChunk(SimpleStepBuilder builder) {
if (commitInterval != null) {
builder.chunk(commitInterval);
}
builder.chunk(chunkCompletionPolicy);
}
protected CompletionPolicy getCompletionPolicy() {
return this.chunkCompletionPolicy;
}
@SuppressWarnings("unchecked")
protected SimpleStepBuilder getSimpleStepBuilder(String stepName) {
return new SimpleStepBuilder(new StepBuilder(stepName));
@@ -983,6 +978,10 @@ public class StepParserStepFactoryBean<I, O> implements FactoryBean, BeanNameAwa
this.commitInterval = commitInterval;
}
protected Integer getCommitInterval() {
return this.commitInterval;
}
/**
* Flag to signal that the reader is transactional (usually a JMS consumer) so that items are re-presented after a
* rollback. The default is false and readers are assumed to be forward-only.
@@ -1116,10 +1115,6 @@ public class StepParserStepFactoryBean<I, O> implements FactoryBean, BeanNameAwa
this.streams = streams;
}
public void setTimeout(Integer timeout) {
this.timeout = timeout;
}
// =========================================================
// Additional
// =========================================================

View File

@@ -43,6 +43,9 @@ import org.springframework.batch.jsr.item.ItemReaderAdapter;
import org.springframework.batch.jsr.item.ItemWriterAdapter;
import org.springframework.batch.jsr.repeat.CheckpointAlgorithmAdapter;
import org.springframework.batch.repeat.CompletionPolicy;
import org.springframework.batch.repeat.policy.CompositeCompletionPolicy;
import org.springframework.batch.repeat.policy.SimpleCompletionPolicy;
import org.springframework.batch.repeat.policy.TimeoutTerminationPolicy;
import org.springframework.beans.factory.FactoryBean;
import org.springframework.util.Assert;
@@ -62,6 +65,8 @@ public class StepFactoryBean extends StepParserStepFactoryBean {
private PartitionReducer reducer;
private Integer timeout;
public void setPartitionReducer(PartitionReducer reducer) {
this.reducer = reducer;
}
@@ -117,6 +122,27 @@ public class StepFactoryBean extends StepParserStepFactoryBean {
return builder.build();
}
@Override
protected void setChunk(SimpleStepBuilder builder) {
if(timeout != null && getCommitInterval() != null) {
CompositeCompletionPolicy completionPolicy = new CompositeCompletionPolicy();
CompletionPolicy [] policies = new CompletionPolicy[2];
policies[0] = new SimpleCompletionPolicy(getCommitInterval());
policies[1] = new TimeoutTerminationPolicy(timeout * 1000);
completionPolicy.setPolicies(policies);
builder.chunk(completionPolicy);
} else if(timeout != null) {
builder.chunk(new TimeoutTerminationPolicy(timeout * 1000));
} else if(getCommitInterval() != null) {
builder.chunk(getCommitInterval());
}
if(getCompletionPolicy() != null) {
builder.chunk(getCompletionPolicy());
}
}
@Override
protected Step createPartitionStep() {
// Creating a partitioned step for the JSR needs to create two steps...the partitioned step and the step being executed.
@@ -263,4 +289,8 @@ public class StepFactoryBean extends StepParserStepFactoryBean {
jsrSimpleStepBuilder.setBatchPropertyContext(batchPropertyContext);
return jsrSimpleStepBuilder;
}
public void setTimeout(Integer timeout) {
this.timeout = timeout;
}
}

View File

@@ -64,7 +64,7 @@ public class SimpleItemBasedJobParsingTests {
assertEquals(4, execution.getStepExecutions().size());
assertEquals(27, processor.count);
assertEquals(1, policy.checkpointCount);
assertEquals(8, writer.writeCount);
assertEquals(7, writer.writeCount);
assertEquals(27, writer.itemCount);
}