Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
@@ -0,0 +1,36 @@
/*
* Copyright 2026-present 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.
* You may obtain a copy of the License at
*
* https://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
* See the License for the specific language governing permissions and
* limitations under the License.
*/
package org.springframework.batch.core.launch;

/**
* Exception thrown when a request to stop a job execution does not complete within the
* configured timeout, i.e. its running step(s) did not reach a terminal state in time.
*
* @author Kyungrae Kim
* @since 6.0
*/
public class JobExecutionStopException extends RuntimeException {

/**
* Create a {@link JobExecutionStopException} with a message and a cause.
* @param msg the message to signal the cause of failure
* @param cause the underlying cause
*/
public JobExecutionStopException(String msg, Throwable cause) {
super(msg, cause);
}

}
Original file line number Diff line number Diff line change
Expand Up @@ -187,14 +187,17 @@ Long startNextInstance(String jobName) throws NoSuchJobException, JobParametersN
boolean stop(long executionId) throws NoSuchJobExecutionException, JobExecutionNotRunningException;

/**
* Send a stop signal to the supplied {@link JobExecution}. The signal is successfully
* sent if this method returns true, but that doesn't mean that the job has stopped.
* The only way to be sure of that is to poll the job execution status.
* Stop the supplied {@link JobExecution} and wait for its running step(s) to stop.
* The job execution is marked as
* {@link org.springframework.batch.core.BatchStatus#STOPPING STOPPING}, the running
* step(s) are signalled to stop, and this method blocks until they have terminated
* and persisted their stopped state, up to a configurable timeout.
* @param jobExecution the running {@link JobExecution}
* @return true if the message was successfully sent (does not guarantee that the job
* has stopped)
* @return {@code true} once the job execution has stopped
* @throws JobExecutionNotRunningException if the supplied {@link JobExecution} is not
* running (so cannot be stopped)
* @throws JobExecutionStopException if the running step(s) do not stop within the
* configured timeout
*/
boolean stop(JobExecution jobExecution) throws JobExecutionNotRunningException;

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -16,6 +16,7 @@
package org.springframework.batch.core.launch.support;

import java.lang.reflect.Method;
import java.time.Duration;

import io.micrometer.observation.ObservationRegistry;
import org.apache.commons.logging.Log;
Expand Down Expand Up @@ -77,6 +78,8 @@ public class JobOperatorFactoryBean implements FactoryBean<JobOperator>, Applica

private JobParametersConverter jobParametersConverter = new DefaultJobParametersConverter();

private Duration stopTimeout = Duration.ofSeconds(30);

@SuppressWarnings("NullAway.Init")
private TaskExecutor taskExecutor;

Expand Down Expand Up @@ -175,6 +178,18 @@ public void setTransactionManager(PlatformTransactionManager transactionManager)
this.transactionManager = transactionManager;
}

/**
* Set how long {@code stop(JobExecution)} waits for the running step(s) to actually
* stop before failing with a {@code JobExecutionStopException}. Defaults to 30
* seconds.
* @param stopTimeout the maximum time to wait for a job to stop
* @since 6.0
*/
public void setStopTimeout(Duration stopTimeout) {
Assert.notNull(stopTimeout, "stopTimeout must not be null");
this.stopTimeout = stopTimeout;
}

/**
* Set the transaction attributes source to use in the created proxy.
* @param transactionAttributeSource the transaction attributes source to use in the
Expand Down Expand Up @@ -217,6 +232,8 @@ private TaskExecutorJobOperator getTarget() throws Exception {
taskExecutorJobOperator.setObservationRegistry(this.observationRegistry);
}
taskExecutorJobOperator.setJobParametersConverter(this.jobParametersConverter);
taskExecutorJobOperator.setTransactionManager(this.transactionManager);
taskExecutorJobOperator.setStopTimeout(this.stopTimeout);
taskExecutorJobOperator.afterPropertiesSet();
return taskExecutorJobOperator;
}
Expand All @@ -226,10 +243,8 @@ private static class DefaultJobOperatorTransactionAttributeSource extends Method
public DefaultJobOperatorTransactionAttributeSource() {
DefaultTransactionAttribute transactionAttribute = new DefaultTransactionAttribute();
try {
Method stopMethod = TaskExecutorJobOperator.class.getMethod("stop", JobExecution.class);
Method abandonMethod = TaskExecutorJobOperator.class.getMethod("abandon", JobExecution.class);
Method recoverMethod = TaskExecutorJobOperator.class.getMethod("recover", JobExecution.class);
addTransactionalMethod(stopMethod, transactionAttribute);
addTransactionalMethod(abandonMethod, transactionAttribute);
addTransactionalMethod(recoverMethod, transactionAttribute);
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -15,6 +15,7 @@
*/
package org.springframework.batch.core.launch.support;

import java.time.Duration;
import java.time.LocalDateTime;
import java.util.ArrayList;
import java.util.LinkedHashMap;
Expand All @@ -24,13 +25,16 @@
import java.util.Properties;
import java.util.Set;
import java.util.TreeSet;
import java.util.concurrent.CompletableFuture;
import java.util.concurrent.ExecutionException;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.TimeoutException;

import org.apache.commons.logging.Log;
import org.apache.commons.logging.LogFactory;
import org.jspecify.annotations.NullUnmarked;

import org.springframework.batch.core.BatchStatus;
import org.springframework.batch.core.ExitStatus;
import org.springframework.batch.core.job.Job;
import org.springframework.batch.core.job.JobExecution;
import org.springframework.batch.core.job.JobInstance;
Expand All @@ -54,10 +58,13 @@
import org.springframework.batch.core.launch.JobInstanceAlreadyCompleteException;
import org.springframework.batch.core.repository.JobRepository;
import org.springframework.batch.core.launch.JobRestartException;
import org.springframework.batch.core.launch.JobExecutionStopException;
import org.springframework.batch.infrastructure.support.transaction.ResourcelessTransactionManager;
import org.springframework.transaction.PlatformTransactionManager;
import org.springframework.transaction.support.TransactionTemplate;
import org.springframework.batch.core.scope.context.StepSynchronizationManager;
import org.springframework.batch.core.step.StepLocator;
import org.springframework.batch.core.step.tasklet.StoppableTasklet;
import org.springframework.batch.core.step.tasklet.Tasklet;
import org.springframework.batch.core.step.tasklet.TaskletStep;
import org.springframework.batch.infrastructure.support.PropertiesConverter;
import org.springframework.beans.factory.InitializingBean;
Expand Down Expand Up @@ -102,6 +109,34 @@ public class SimpleJobOperator extends TaskExecutorJobLauncher implements JobOpe

private final Log logger = LogFactory.getLog(getClass());

private PlatformTransactionManager transactionManager = new ResourcelessTransactionManager();

private Duration stopTimeout = Duration.ofSeconds(30);

/**
* Set the transaction manager used to persist the stopping state of a job execution
* atomically before waiting for it to stop. Defaults to a
* {@link ResourcelessTransactionManager}.
* @param transactionManager the transaction manager to use
* @since 6.0
*/
public void setTransactionManager(PlatformTransactionManager transactionManager) {
Assert.notNull(transactionManager, "transactionManager must not be null");
this.transactionManager = transactionManager;
}

/**
* Set how long {@link #stop(JobExecution)} waits for the running step(s) to actually
* stop before giving up with a {@link JobExecutionStopException}. Defaults to 30
* seconds (aligned with the Spring graceful shutdown convention).
* @param stopTimeout the maximum time to wait for the job to stop
* @since 6.0
*/
public void setStopTimeout(Duration stopTimeout) {
Assert.notNull(stopTimeout, "stopTimeout must not be null");
this.stopTimeout = stopTimeout;
}

/**
* Check mandatory properties.
*
Expand Down Expand Up @@ -342,51 +377,78 @@ public boolean stop(JobExecution jobExecution) throws JobExecutionNotRunningExce
if (logger.isInfoEnabled()) {
logger.info("Stopping job execution: " + jobExecution);
}
jobExecution.setStatus(BatchStatus.STOPPING); // will be upgraded to STOPPED in
// JobRepository.update
jobExecution.setExitStatus(ExitStatus.STOPPED);
jobExecution.setEndTime(LocalDateTime.now());
jobRepository.update(jobExecution);
jobRepository.updateExecutionContext(jobExecution);

List<CompletableFuture<StepExecution>> terminations = new ArrayList<>();
List<Runnable> stopSignals = new ArrayList<>();
Job job = jobRegistry.getJob(jobExecution.getJobInstance().getJobName());
if (job != null) {
if (job instanceof StepLocator stepLocator) {
// can only process as StepLocator is the only way to get the step object
// get the current stepExecution
for (StepExecution stepExecution : jobExecution.getStepExecutions()) {
if (stepExecution.getStatus().isRunning()) {
// have the step execution that's running -> need to 'stop' it
Step step = stepLocator.getStep(stepExecution.getStepName());
if (step != null) {
if (step instanceof TaskletStep taskletStep) {
Tasklet tasklet = taskletStep.getTasklet();
if (tasklet instanceof StoppableTasklet stoppableTasklet) {
StepSynchronizationManager.register(stepExecution);
stoppableTasklet.stop(stepExecution);
jobRepository.update(stepExecution);
jobRepository.updateExecutionContext(stepExecution);
StepSynchronizationManager.release();
}
}
if (step instanceof StoppableStep stoppableStep) {
StepSynchronizationManager.register(stepExecution);
stoppableStep.stop(stepExecution);
jobRepository.update(stepExecution);
jobRepository.updateExecutionContext(stepExecution);
StepSynchronizationManager.release();
}
if (job instanceof StepLocator stepLocator) {
for (StepExecution stepExecution : jobExecution.getStepExecutions()) {
if (!stepExecution.getStatus().isRunning()) {
continue;
}
Step step = stepLocator.getStep(stepExecution.getStepName());
if (step instanceof StoppableStep stoppableStep) {
terminations.add(stoppableStep.subscribeToTermination(stepExecution));
stopSignals.add(() -> {
StepSynchronizationManager.register(stepExecution);
try {
stoppableStep.stop(stepExecution);
}
}
finally {
StepSynchronizationManager.release();
}
});
}
if (step instanceof TaskletStep taskletStep
&& taskletStep.getTasklet() instanceof StoppableTasklet stoppableTasklet) {
stopSignals.add(() -> {
StepSynchronizationManager.register(stepExecution);
try {
stoppableTasklet.stop(stepExecution);
}
finally {
StepSynchronizationManager.release();
}
});
}
}
// TODO what if the job is not a StepLocator? ie a job with no steps?
// FIXME Job should provide a stop() method

}

// Persist STOPPING in its own short transaction that commits at once,
// so it is durable and holds no lock during the wait.
new TransactionTemplate(this.transactionManager).executeWithoutResult(transactionStatus -> {
jobExecution.setStatus(BatchStatus.STOPPING);
jobRepository.update(jobExecution);
});

stopSignals.forEach(Runnable::run);
awaitStop(jobExecution, terminations);
return true;
}

private void awaitStop(JobExecution jobExecution, List<CompletableFuture<StepExecution>> terminations) {
try {
CompletableFuture.allOf(terminations.toArray(new CompletableFuture[0]))
.get(this.stopTimeout.toMillis(), TimeUnit.MILLISECONDS);
}
catch (TimeoutException e) {
throw new JobExecutionStopException("Timed out after " + this.stopTimeout
+ " while waiting for job execution " + jobExecution.getId() + " to stop", e);
}
catch (ExecutionException e) {
throw new JobExecutionStopException(
"Failure while waiting for job execution " + jobExecution.getId() + " to stop", e.getCause());
}
catch (InterruptedException e) {
Thread.currentThread().interrupt();
throw new JobExecutionStopException(
"Interrupted while waiting for job execution " + jobExecution.getId() + " to stop", e);
}
}

@Override
@Deprecated(since = "6.0", forRemoval = true)
public JobExecution abandon(long jobExecutionId)
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -164,7 +164,6 @@ public void update(StepExecution stepExecution) {
this.jobExecutionDao.synchronizeStatus(jobExecution);

if (jobExecution.isStopped() || jobExecution.isStopping()) {
this.stepExecutionDao.synchronizeStatus(stepExecution);
stepExecution.setTerminateOnly();
}

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -18,6 +18,10 @@
import java.time.Duration;
import java.time.LocalDateTime;
import java.util.List;
import java.util.Map;
import java.util.Set;
import java.util.concurrent.CompletableFuture;
import java.util.concurrent.ConcurrentHashMap;
import java.util.stream.Collectors;

import io.micrometer.observation.Observation;
Expand Down Expand Up @@ -74,6 +78,10 @@ public abstract class AbstractStep implements StoppableStep, InitializingBean, B

private final CompositeStepExecutionListener stepExecutionListener = new CompositeStepExecutionListener();

private final Set<Long> executingStepExecutions = ConcurrentHashMap.newKeySet();

private final Map<Long, CompletableFuture<StepExecution>> terminationSignals = new ConcurrentHashMap<>();

private JobRepository jobRepository;

protected ObservationRegistry observationRegistry;
Expand Down Expand Up @@ -212,6 +220,7 @@ public final void execute(StepExecution stepExecution)
throws JobInterruptedException, UnexpectedJobExecutionException {

Assert.notNull(stepExecution, "stepExecution must not be null");
this.executingStepExecutions.add(stepExecution.getId());
stepExecution.getExecutionContext().put(SpringBatchVersion.BATCH_VERSION_KEY, SpringBatchVersion.getVersion());

if (logger.isDebugEnabled()) {
Expand Down Expand Up @@ -355,7 +364,28 @@ public final void execute(StepExecution stepExecution)
if (logger.isDebugEnabled()) {
logger.debug("Step execution complete: " + stepExecution.getSummary());
}

// Notify any caller of stop(StepExecution) that this execution has terminated
// and its final metadata has been saved.
this.executingStepExecutions.remove(stepExecution.getId());
CompletableFuture<StepExecution> terminationSignal = this.terminationSignals.remove(stepExecution.getId());
if (terminationSignal != null) {
terminationSignal.complete(stepExecution);
}
}
}

@Override
public CompletableFuture<StepExecution> subscribeToTermination(StepExecution stepExecution) {
Long stepExecutionId = stepExecution.getId();
CompletableFuture<StepExecution> terminationSignal = this.terminationSignals
.computeIfAbsent(stepExecutionId, key -> new CompletableFuture<>());
// If the execution is not running in this JVM, it has already terminated.
if (!this.executingStepExecutions.contains(stepExecution.getId())) {
terminationSignal.complete(stepExecution);
this.terminationSignals.remove(stepExecutionId, terminationSignal);
}
return terminationSignal;
}

private void stopObservation(StepExecution stepExecution, Observation observation) {
Expand Down
Loading