diff --git a/spring-batch-core/src/main/java/org/springframework/batch/core/job/flow/JobFlowExecutor.java b/spring-batch-core/src/main/java/org/springframework/batch/core/job/flow/JobFlowExecutor.java index 96d358742a..c2f4cd2c8e 100644 --- a/spring-batch-core/src/main/java/org/springframework/batch/core/job/flow/JobFlowExecutor.java +++ b/spring-batch-core/src/main/java/org/springframework/batch/core/job/flow/JobFlowExecutor.java @@ -42,7 +42,7 @@ */ public class JobFlowExecutor implements FlowExecutor { - private static final ThreadLocal stepExecutionHolder = new ThreadLocal<>(); + private final ThreadLocal stepExecutionHolder = new ThreadLocal<>(); private final JobExecution execution; 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 3e4174722c..b9d1c4672d 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 @@ -166,6 +166,12 @@ public FlowExecution resume(String stateName, FlowExecutor executor) throws Flow } status = state.handle(executor); stepExecution = executor.getStepExecution(); + + if (logger.isDebugEnabled()) { + logger.debug("Completed state=" + stateName + " with status=" + status); + } + + state = nextState(stateName, status, stepExecution); } catch (FlowExecutionException e) { executor.close(new FlowExecution(stateName, status)); @@ -176,12 +182,6 @@ public FlowExecution resume(String stateName, FlowExecutor executor) throws Flow throw new FlowExecutionException( String.format("Ended flow=%s at state=%s with exception", name, stateName), e); } - - if (logger.isDebugEnabled()) { - logger.debug("Completed state=" + stateName + " with status=" + status); - } - - state = nextState(stateName, status, stepExecution); } FlowExecution result = new FlowExecution(stateName, status); diff --git a/spring-batch-core/src/test/java/org/springframework/batch/core/job/flow/FlowJobFailureTests.java b/spring-batch-core/src/test/java/org/springframework/batch/core/job/flow/FlowJobFailureTests.java index 506ed8e87e..3e72f6c87d 100644 --- a/spring-batch-core/src/test/java/org/springframework/batch/core/job/flow/FlowJobFailureTests.java +++ b/spring-batch-core/src/test/java/org/springframework/batch/core/job/flow/FlowJobFailureTests.java @@ -15,8 +15,11 @@ */ package org.springframework.batch.core.job.flow; -import org.springframework.batch.infrastructure.support.DatabaseType; import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertNull; +import static org.junit.jupiter.api.Assertions.assertSame; +import static org.junit.jupiter.api.Assertions.assertThrows; +import static org.junit.jupiter.api.Assertions.assertTrue; import java.util.ArrayList; import java.util.List; @@ -28,8 +31,10 @@ import org.springframework.batch.core.job.JobExecution; import org.springframework.batch.core.job.JobInstance; import org.springframework.batch.core.job.JobInterruptedException; +import org.springframework.batch.core.job.SimpleStepHandler; import org.springframework.batch.core.job.parameters.JobParameters; import org.springframework.batch.core.step.StepExecution; +import org.springframework.batch.core.step.Step; import org.springframework.batch.core.job.UnexpectedJobExecutionException; import org.springframework.batch.core.job.flow.support.SimpleFlow; import org.springframework.batch.core.job.flow.support.StateTransition; @@ -39,6 +44,7 @@ import org.springframework.batch.core.repository.support.JdbcJobRepositoryFactoryBean; import org.springframework.batch.core.step.StepSupport; import org.springframework.batch.infrastructure.item.ExecutionContext; +import org.springframework.batch.infrastructure.support.DatabaseType; import org.springframework.jdbc.support.JdbcTransactionManager; import org.springframework.jdbc.datasource.embedded.EmbeddedDatabase; import org.springframework.jdbc.datasource.embedded.EmbeddedDatabaseBuilder; @@ -57,6 +63,8 @@ class FlowJobFailureTests { private JobExecution execution; + private JobRepository jobRepository; + @BeforeEach void init() throws Exception { EmbeddedDatabase embeddedDatabase = new EmbeddedDatabaseBuilder() @@ -68,7 +76,7 @@ void init() throws Exception { factory.setDataSource(embeddedDatabase); factory.setTransactionManager(new JdbcTransactionManager(embeddedDatabase)); factory.afterPropertiesSet(); - JobRepository jobRepository = factory.getObject(); + jobRepository = factory.getObject(); job.setJobRepository(jobRepository); JobParameters jobParameters = new JobParameters(); JobInstance jobInstance = jobRepository.createJobInstance("job", jobParameters); @@ -92,6 +100,109 @@ void testStepFailure() throws Exception { assertEquals(BatchStatus.FAILED, execution.getStatus()); } + @Test + void testUnmatchedTransitionDoesNotAffectOtherJob() throws Exception { + job.setFlow(failingFlow()); + job.afterPropertiesSet(); + job.execute(execution); + assertEquals(BatchStatus.FAILED, execution.getStatus()); + StepExecution failedStep = jobRepository.getLastStepExecution(execution.getJobInstance(), "step"); + assertEquals(BatchStatus.FAILED, failedStep.getStatus()); + + FlowJob otherJob = new FlowJob("otherJob"); + otherJob.setJobRepository(jobRepository); + SimpleFlow otherFlow = new SimpleFlow("otherFlow"); + StepState otherStep = new StepState(new StepSupport("otherStep") { + @Override + public void execute(StepExecution stepExecution) { + stepExecution.setStatus(BatchStatus.COMPLETED); + stepExecution.setExitStatus(ExitStatus.COMPLETED); + jobRepository.update(stepExecution); + } + }); + otherFlow.setStateTransitions(List.of(StateTransition.createEndStateTransition(otherStep))); + otherJob.setFlow(otherFlow); + otherJob.afterPropertiesSet(); + JobParameters parameters = new JobParameters(); + JobInstance instance = jobRepository.createJobInstance("otherJob", parameters); + JobExecution otherExecution = jobRepository.createJobExecution(instance, parameters, new ExecutionContext()); + otherJob.execute(otherExecution); + assertEquals(BatchStatus.COMPLETED, otherExecution.getStatus()); + + StepExecution reloaded = jobRepository.getLastStepExecution(execution.getJobInstance(), "step"); + assertEquals(BatchStatus.FAILED, reloaded.getStatus()); + assertEquals(failedStep.getVersion(), reloaded.getVersion()); + } + + @Test + void testUnmatchedTransitionDoesNotPreventRestart() throws Exception { + job.setFlow(failingFlow()); + job.afterPropertiesSet(); + job.execute(execution); + assertEquals(BatchStatus.FAILED, execution.getStatus()); + + JobExecution restart = jobRepository.createJobExecution(execution.getJobInstance(), new JobParameters(), + new ExecutionContext()); + job.execute(restart); + + assertEquals(BatchStatus.FAILED, restart.getStatus()); + assertEquals(1, restart.getStepExecutions().size()); + assertEquals(2, jobRepository.getStepExecutionCount(execution.getJobInstance(), "step")); + assertEquals(BatchStatus.FAILED, + jobRepository.getLastStepExecution(execution.getJobInstance(), "step").getStatus()); + } + + @Test + void testUnmatchedTransitionClosesExecutor() throws Exception { + JobFlowExecutor executor = new JobFlowExecutor(jobRepository, new SimpleStepHandler(jobRepository), execution); + SimpleFlow flow = failingFlow(); + try { + FlowExecutionException exception = assertThrows(FlowExecutionException.class, () -> flow.start(executor)); + assertTrue(exception.getMessage().contains("Next state not found")); + assertNull(executor.getStepExecution()); + } + finally { + executor.close(new FlowExecution("step", FlowExecutionStatus.FAILED)); + } + } + + @Test + void testExecutorsKeepSeparateStepExecutions() throws Exception { + JobExecution otherExecution = jobRepository.createJobExecution(execution.getJobInstance(), new JobParameters(), + new ExecutionContext()); + JobFlowExecutor executor = new JobFlowExecutor(jobRepository, new SimpleStepHandler(jobRepository), execution); + JobFlowExecutor otherExecutor = new JobFlowExecutor(jobRepository, new SimpleStepHandler(jobRepository), + otherExecution); + try { + executor.executeStep(failingStep()); + StepExecution stepExecution = executor.getStepExecution(); + assertNull(otherExecutor.getStepExecution()); + otherExecutor.close(new FlowExecution("other", FlowExecutionStatus.COMPLETED)); + assertSame(stepExecution, executor.getStepExecution()); + } + finally { + executor.close(new FlowExecution("step", FlowExecutionStatus.FAILED)); + } + } + + private SimpleFlow failingFlow() { + SimpleFlow flow = new SimpleFlow("job"); + flow.setStateTransitions( + List.of(StateTransition.createEndStateTransition(new StepState(failingStep()), "COMPLETED"))); + return flow; + } + + private Step failingStep() { + return new StepSupport("step") { + @Override + public void execute(StepExecution stepExecution) { + stepExecution.setStatus(BatchStatus.FAILED); + stepExecution.setExitStatus(ExitStatus.FAILED); + jobRepository.update(stepExecution); + } + }; + } + @Test void testStepStatusUnknown() throws Exception { SimpleFlow flow = new SimpleFlow("job");