diff --git a/spring-batch-core/src/main/java/org/springframework/batch/core/job/flow/support/SimpleFlow.java b/spring-batch-core/src/main/java/org/springframework/batch/core/job/flow/support/SimpleFlow.java index b0677b67e..5165230b6 100644 --- a/spring-batch-core/src/main/java/org/springframework/batch/core/job/flow/support/SimpleFlow.java +++ b/spring-batch-core/src/main/java/org/springframework/batch/core/job/flow/support/SimpleFlow.java @@ -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, diff --git a/spring-batch-core/src/main/java/org/springframework/batch/core/job/flow/support/state/SplitState.java b/spring-batch-core/src/main/java/org/springframework/batch/core/job/flow/support/state/SplitState.java index c04655321..340f9a0a5 100644 --- a/spring-batch-core/src/main/java/org/springframework/batch/core/job/flow/support/state/SplitState.java +++ b/spring-batch-core/src/main/java/org/springframework/batch/core/job/flow/support/state/SplitState.java @@ -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 results = new ArrayList(); - // TODO: could use a CompletionService + // Could use a CompletionService here? for (Future 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); diff --git a/spring-batch-core/src/test/java/org/springframework/batch/core/job/flow/FlowJobTests.java b/spring-batch-core/src/test/java/org/springframework/batch/core/job/flow/FlowJobTests.java index 721b51eac..4f2b3235b 100644 --- a/spring-batch-core/src/test/java/org/springframework/batch/core/job/flow/FlowJobTests.java +++ b/spring-batch-core/src/test/java/org/springframework/batch/core/job/flow/FlowJobTests.java @@ -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 transitions = new ArrayList(); + 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(transitions)); + flow1.afterPropertiesSet(); + flow2.setStateTransitions(new ArrayList(transitions)); + flow2.afterPropertiesSet(); + + transitions = new ArrayList(); + transitions.add(StateTransition.createStateTransition(new SplitState(Arrays.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"); diff --git a/spring-batch-core/src/test/java/org/springframework/batch/core/job/flow/support/state/SplitStateTests.java b/spring-batch-core/src/test/java/org/springframework/batch/core/job/flow/support/state/SplitStateTests.java index a41e2a0f3..5254507c3 100644 --- a/spring-batch-core/src/test/java/org/springframework/batch/core/job/flow/support/state/SplitStateTests.java +++ b/spring-batch-core/src/test/java/org/springframework/batch/core/job/flow/support/state/SplitStateTests.java @@ -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 flows = new ArrayList(); 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));