From 2b0d0cdc33a524c473a235341ac0512a6808f0c5 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?=E9=83=91=E6=9C=9D=E9=92=A6?= <1357598741@qq.com> Date: Tue, 8 Sep 2026 13:48:10 +0800 Subject: [PATCH] Prevent step execution leakage between flow job executions MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Close the flow executor when transition lookup fails and keep the step execution holder local to each executor while preserving thread isolation. Add JDBC-backed regression coverage for restart, unrelated job execution, exceptional cleanup, and executor isolation. Fixes #5514 Signed-off-by: 郑朝钦 <1357598741@qq.com> --- .../batch/core/job/flow/JobFlowExecutor.java | 4 +- .../core/job/flow/support/SimpleFlow.java | 14 +-- .../core/job/flow/FlowJobFailureTests.java | 115 +++++++++++++++++- 3 files changed, 122 insertions(+), 11 deletions(-) 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 89ace2f924..e54860a4ae 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 @@ -1,5 +1,5 @@ /* - * Copyright 2006-2024 the original author or authors. + * Copyright 2006-2026 the original author or authors. * * Licensed under the Apache License, Version 2.0 (the "License"); * you may not use this file except in compliance with the License. @@ -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 434807d73f..570c0bb3d9 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 @@ -1,5 +1,5 @@ /* - * Copyright 2006-2024 the original author or authors. + * Copyright 2006-2026 the original author or authors. * * Licensed under the Apache License, Version 2.0 (the "License"); * you may not use this file except in compliance with the License. @@ -165,6 +165,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)); @@ -175,12 +181,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 a025c1cb69..dc43278943 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 @@ -1,5 +1,5 @@ /* - * Copyright 2010-2023 the original author or authors. + * Copyright 2010-2026 the original author or authors. * * Licensed under the Apache License, Version 2.0 (the "License"); * you may not use this file except in compliance with the License. @@ -16,6 +16,10 @@ package org.springframework.batch.core.job.flow; 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; @@ -27,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; @@ -56,6 +62,8 @@ class FlowJobFailureTests { private JobExecution execution; + private JobRepository jobRepository; + @BeforeEach void init() throws Exception { EmbeddedDatabase embeddedDatabase = new EmbeddedDatabaseBuilder() @@ -67,7 +75,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); @@ -91,6 +99,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");