diff --git a/spring-batch-core/src/main/java/org/springframework/batch/core/job/flow/FlowExecutor.java b/spring-batch-core/src/main/java/org/springframework/batch/core/job/flow/FlowExecutor.java index 6b83034ed..893a54f66 100644 --- a/spring-batch-core/src/main/java/org/springframework/batch/core/job/flow/FlowExecutor.java +++ b/spring-batch-core/src/main/java/org/springframework/batch/core/job/flow/FlowExecutor.java @@ -66,4 +66,19 @@ public interface FlowExecutor { */ void updateStepExecutionStatus(); + /** + * Push a token onto a stack to indicate that the context is being nested. + */ + void nest(); + + /** + * Pop a token off a stack to indicate that the context is being un-nested. + */ + void unnest(); + + /** + * @return indicate whether the execution context is nested + */ + boolean isNested(); + } diff --git a/spring-batch-core/src/main/java/org/springframework/batch/core/job/flow/FlowJob.java b/spring-batch-core/src/main/java/org/springframework/batch/core/job/flow/FlowJob.java index 85aa80758..0d0212706 100644 --- a/spring-batch-core/src/main/java/org/springframework/batch/core/job/flow/FlowJob.java +++ b/spring-batch-core/src/main/java/org/springframework/batch/core/job/flow/FlowJob.java @@ -100,6 +100,8 @@ public class FlowJob extends AbstractJob { private final JobExecution execution; + private volatile boolean nested = false; + /** * @param execution */ @@ -138,6 +140,18 @@ public class FlowJob extends AbstractJob { stepExecutionHolder.set(null); } + public boolean isNested() { + return nested; + } + + public void nest() { + nested = true; + } + + public void unnest() { + nested = false; + } + } } diff --git a/spring-batch-core/src/main/java/org/springframework/batch/core/job/flow/support/state/EndState.java b/spring-batch-core/src/main/java/org/springframework/batch/core/job/flow/support/state/EndState.java index 478c47b5d..51934ad03 100644 --- a/spring-batch-core/src/main/java/org/springframework/batch/core/job/flow/support/state/EndState.java +++ b/spring-batch-core/src/main/java/org/springframework/batch/core/job/flow/support/state/EndState.java @@ -77,25 +77,25 @@ public class EndState extends AbstractState { /* * If there are step executions, then we are not at the * beginning of a restart. - * - * N.B. EndState has to be able to set the status directly, but - * only because the internal flows inside SplitStates contain - * EndState (which maybe they should not, since the JobExecution - * is not ending). */ - jobExecution.setStatus(status); - jobExecution.setExitStatus(exitStatus); + if (!executor.isNested()) { + jobExecution.setStatus(status); + jobExecution.setExitStatus(exitStatus); + } if (status == BatchStatus.STOPPED) { + if (abandon) { + /* + * Only if instructed to do so upgrade the status of + * last step execution so it is not replayed on a + * restart... + */ + executor.updateStepExecutionStatus(); + } /* * If we are in flight (not a restart) and we are supposed * to signal a stop, then make sure that happens * irrespective of the exit status. */ - if (abandon) { - // Only if instructed to do so upgrade the status of - // last step execution... - executor.updateStepExecutionStatus(); - } return FlowExecutionStatus.STOPPED; } } 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 e1f75ed63..efa0fe8e7 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 @@ -71,6 +71,10 @@ public class SplitState extends AbstractState { @Override public FlowExecutionStatus handle(final FlowExecutor executor) throws Exception { + try { + + executor.nest(); + // TODO: collect the last StepExecution from the flows as well, so they // can be abandoned if necessary Collection> tasks = new ArrayList>(); @@ -103,6 +107,10 @@ public class SplitState extends AbstractState { return aggregator.aggregate(results); + } finally { + executor.unnest(); + } + } /* diff --git a/spring-batch-core/src/test/java/org/springframework/batch/core/job/flow/support/JobFlowExecutorSupport.java b/spring-batch-core/src/test/java/org/springframework/batch/core/job/flow/support/JobFlowExecutorSupport.java index 55789ebd5..4f104f7ae 100644 --- a/spring-batch-core/src/test/java/org/springframework/batch/core/job/flow/support/JobFlowExecutorSupport.java +++ b/spring-batch-core/src/test/java/org/springframework/batch/core/job/flow/support/JobFlowExecutorSupport.java @@ -31,6 +31,8 @@ import org.springframework.batch.core.repository.JobRestartException; */ public class JobFlowExecutorSupport implements FlowExecutor { + private volatile boolean nested = false; + public String executeStep(Step step) throws JobInterruptedException, JobRestartException, StartLimitExceededException { return ExitStatus.COMPLETED.getExitCode(); @@ -50,4 +52,16 @@ public class JobFlowExecutorSupport implements FlowExecutor { public void updateStepExecutionStatus() { } + public boolean isNested() { + return nested; + } + + public void nest() { + nested= true; + } + + public void unnest() { + nested = false; + } + } 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 ae53b4667..a41e2a0f3 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 @@ -25,6 +25,7 @@ import org.junit.Test; import org.springframework.batch.core.job.flow.Flow; import org.springframework.batch.core.job.flow.FlowExecution; import org.springframework.batch.core.job.flow.FlowExecutionStatus; +import org.springframework.batch.core.job.flow.support.JobFlowExecutorSupport; import org.springframework.core.task.SimpleAsyncTaskExecutor; @@ -34,6 +35,8 @@ import org.springframework.core.task.SimpleAsyncTaskExecutor; */ public class SplitStateTests { + private JobFlowExecutorSupport executor = new JobFlowExecutorSupport(); + @Test public void testBasicHandling() throws Exception { @@ -45,11 +48,11 @@ public class SplitStateTests { SplitState state = new SplitState(flows, "foo"); - EasyMock.expect(flow1.start(null)).andReturn(new FlowExecution("step1", FlowExecutionStatus.COMPLETED)); - EasyMock.expect(flow2.start(null)).andReturn(new FlowExecution("step1", FlowExecutionStatus.COMPLETED)); + EasyMock.expect(flow1.start(executor)).andReturn(new FlowExecution("step1", FlowExecutionStatus.COMPLETED)); + EasyMock.expect(flow2.start(executor)).andReturn(new FlowExecution("step1", FlowExecutionStatus.COMPLETED)); EasyMock.replay(flow1, flow2); - FlowExecutionStatus result = state.handle(null); + FlowExecutionStatus result = state.handle(executor); assertEquals(FlowExecutionStatus.COMPLETED, result); EasyMock.verify(flow1, flow2); @@ -68,11 +71,11 @@ public class SplitStateTests { SplitState state = new SplitState(flows, "foo"); state.setTaskExecutor(new SimpleAsyncTaskExecutor()); - EasyMock.expect(flow1.start(null)).andReturn(new FlowExecution("step1", FlowExecutionStatus.COMPLETED)); - EasyMock.expect(flow2.start(null)).andReturn(new FlowExecution("step1", FlowExecutionStatus.COMPLETED)); + EasyMock.expect(flow1.start(executor)).andReturn(new FlowExecution("step1", FlowExecutionStatus.COMPLETED)); + EasyMock.expect(flow2.start(executor)).andReturn(new FlowExecution("step1", FlowExecutionStatus.COMPLETED)); EasyMock.replay(flow1, flow2); - FlowExecutionStatus result = state.handle(null); + FlowExecutionStatus result = state.handle(executor); assertEquals(FlowExecutionStatus.COMPLETED, result); EasyMock.verify(flow1, flow2);