RESOLVED - issue BATCH-1363: Job Stop issue Sequential Vs Flow Step
This commit is contained in:
@@ -143,6 +143,10 @@ public class SimpleFlow implements Flow, InitializingBean {
|
||||
logger.debug("Handling state="+stateName);
|
||||
status = state.handle(executor);
|
||||
}
|
||||
catch (FlowExecutionException e) {
|
||||
executor.close(new FlowExecution(stateName, status));
|
||||
throw e;
|
||||
}
|
||||
catch (Exception e) {
|
||||
executor.close(new FlowExecution(stateName, status));
|
||||
throw new FlowExecutionException(String.format("Ended flow=%s at state=%s with exception", name,
|
||||
|
||||
@@ -18,6 +18,7 @@ package org.springframework.batch.core.job.flow.support.state;
|
||||
import java.util.ArrayList;
|
||||
import java.util.Collection;
|
||||
import java.util.concurrent.Callable;
|
||||
import java.util.concurrent.ExecutionException;
|
||||
import java.util.concurrent.Future;
|
||||
import java.util.concurrent.FutureTask;
|
||||
|
||||
@@ -96,9 +97,20 @@ public class SplitState extends AbstractState {
|
||||
|
||||
Collection<FlowExecution> results = new ArrayList<FlowExecution>();
|
||||
|
||||
// TODO: could use a CompletionService
|
||||
// Could use a CompletionService here?
|
||||
for (Future<FlowExecution> task : tasks) {
|
||||
results.add(task.get());
|
||||
try {
|
||||
results.add(task.get());
|
||||
}
|
||||
catch (ExecutionException e) {
|
||||
// Unwrap the expected exceptions
|
||||
Throwable cause = e.getCause();
|
||||
if (cause instanceof Exception) {
|
||||
throw (Exception) cause;
|
||||
} else {
|
||||
throw e;
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
return aggregator.aggregate(results);
|
||||
|
||||
@@ -21,6 +21,7 @@ import static org.junit.Assert.assertNull;
|
||||
import static org.junit.Assert.fail;
|
||||
|
||||
import java.util.ArrayList;
|
||||
import java.util.Arrays;
|
||||
import java.util.List;
|
||||
|
||||
import org.junit.Before;
|
||||
@@ -37,6 +38,7 @@ import org.springframework.batch.core.job.flow.support.SimpleFlow;
|
||||
import org.springframework.batch.core.job.flow.support.StateTransition;
|
||||
import org.springframework.batch.core.job.flow.support.state.DecisionState;
|
||||
import org.springframework.batch.core.job.flow.support.state.EndState;
|
||||
import org.springframework.batch.core.job.flow.support.state.SplitState;
|
||||
import org.springframework.batch.core.job.flow.support.state.StepState;
|
||||
import org.springframework.batch.core.repository.JobRepository;
|
||||
import org.springframework.batch.core.repository.dao.JobExecutionDao;
|
||||
@@ -215,6 +217,41 @@ public class FlowJobTests {
|
||||
assertEquals(JobInterruptedException.class, jobExecution.getFailureExceptions().get(0).getClass());
|
||||
}
|
||||
|
||||
@Test
|
||||
public void testInterruptedSplit() throws Exception {
|
||||
SimpleFlow flow = new SimpleFlow("job");
|
||||
SimpleFlow flow1 = new SimpleFlow("flow1");
|
||||
SimpleFlow flow2 = new SimpleFlow("flow2");
|
||||
|
||||
List<StateTransition> transitions = new ArrayList<StateTransition>();
|
||||
transitions.add(StateTransition.createStateTransition(new StepState(new StubStep("step1") {
|
||||
@Override
|
||||
public void execute(StepExecution stepExecution) throws JobInterruptedException {
|
||||
stepExecution.setStatus(BatchStatus.STOPPING);
|
||||
jobRepository.update(stepExecution);
|
||||
}
|
||||
}), "end0"));
|
||||
transitions.add(StateTransition.createEndStateTransition(new EndState(FlowExecutionStatus.COMPLETED, "end0")));
|
||||
flow1.setStateTransitions(new ArrayList<StateTransition>(transitions));
|
||||
flow1.afterPropertiesSet();
|
||||
flow2.setStateTransitions(new ArrayList<StateTransition>(transitions));
|
||||
flow2.afterPropertiesSet();
|
||||
|
||||
transitions = new ArrayList<StateTransition>();
|
||||
transitions.add(StateTransition.createStateTransition(new SplitState(Arrays.<Flow>asList(flow1, flow2), "split"), "end0"));
|
||||
transitions.add(StateTransition.createEndStateTransition(new EndState(FlowExecutionStatus.COMPLETED, "end0")));
|
||||
flow.setStateTransitions(transitions);
|
||||
flow.afterPropertiesSet();
|
||||
|
||||
job.setFlow(flow);
|
||||
job.afterPropertiesSet();
|
||||
job.execute(jobExecution);
|
||||
assertEquals(BatchStatus.STOPPED, jobExecution.getStatus());
|
||||
checkRepository(BatchStatus.STOPPED, ExitStatus.STOPPED);
|
||||
assertEquals(1, jobExecution.getAllFailureExceptions().size());
|
||||
assertEquals(JobInterruptedException.class, jobExecution.getFailureExceptions().get(0).getClass());
|
||||
}
|
||||
|
||||
@Test
|
||||
public void testEndStateStopped() throws Exception {
|
||||
SimpleFlow flow = new SimpleFlow("job");
|
||||
|
||||
@@ -18,6 +18,7 @@ package org.springframework.batch.core.job.flow.support.state;
|
||||
import static org.junit.Assert.assertEquals;
|
||||
|
||||
import java.util.ArrayList;
|
||||
import java.util.Arrays;
|
||||
import java.util.Collection;
|
||||
|
||||
import org.easymock.EasyMock;
|
||||
@@ -62,13 +63,10 @@ public class SplitStateTests {
|
||||
@Test
|
||||
public void testConcurrentHandling() throws Exception {
|
||||
|
||||
Collection<Flow> flows = new ArrayList<Flow>();
|
||||
Flow flow1 = EasyMock.createMock(Flow.class);
|
||||
Flow flow2 = EasyMock.createMock(Flow.class);
|
||||
flows.add(flow1);
|
||||
flows.add(flow2);
|
||||
|
||||
SplitState state = new SplitState(flows, "foo");
|
||||
SplitState state = new SplitState(Arrays.asList(flow1, flow2), "foo");
|
||||
state.setTaskExecutor(new SimpleAsyncTaskExecutor());
|
||||
|
||||
EasyMock.expect(flow1.start(executor)).andReturn(new FlowExecution("step1", FlowExecutionStatus.COMPLETED));
|
||||
|
||||
Reference in New Issue
Block a user