diff --git a/spring-beans/src/main/java/org/springframework/beans/factory/support/DisposableBeanAdapter.java b/spring-beans/src/main/java/org/springframework/beans/factory/support/DisposableBeanAdapter.java index 420e3eb9a37..ee2ed6752c9 100644 --- a/spring-beans/src/main/java/org/springframework/beans/factory/support/DisposableBeanAdapter.java +++ b/spring-beans/src/main/java/org/springframework/beans/factory/support/DisposableBeanAdapter.java @@ -23,6 +23,7 @@ import java.util.ArrayList; import java.util.List; import java.util.concurrent.CountDownLatch; import java.util.concurrent.ExecutionException; +import java.util.concurrent.ExecutorService; import java.util.concurrent.Future; import org.apache.commons.logging.Log; @@ -412,14 +413,29 @@ class DisposableBeanAdapter implements DisposableBean, Runnable, Serializable { String destroyMethodName = beanDefinition.resolvedDestroyMethodName; if (destroyMethodName == null) { destroyMethodName = beanDefinition.getDestroyMethodName(); - boolean autoCloseable = (AutoCloseable.class.isAssignableFrom(target)); + boolean autoCloseable = AutoCloseable.class.isAssignableFrom(target); + boolean executorService = ExecutorService.class.isAssignableFrom(target); if (AbstractBeanDefinition.INFER_METHOD.equals(destroyMethodName) || - (destroyMethodName == null && autoCloseable)) { + (destroyMethodName == null && (autoCloseable || executorService))) { // Only perform destroy method inference in case of the bean // not explicitly implementing the DisposableBean interface destroyMethodName = null; if (!(DisposableBean.class.isAssignableFrom(target))) { - if (autoCloseable) { + if (executorService) { + destroyMethodName = SHUTDOWN_METHOD_NAME; + try { + // On JDK 19+, avoid the ExecutorService-level AutoCloseable default implementation + // which awaits task termination for 1 day, even for delayed tasks such as cron jobs. + // Custom close() implementations in ExecutorService subclasses are still accepted. + if (target.getMethod(CLOSE_METHOD_NAME).getDeclaringClass() != ExecutorService.class) { + destroyMethodName = CLOSE_METHOD_NAME; + } + } + catch (NoSuchMethodException ex) { + // Ignore - stick with shutdown() + } + } + else if (autoCloseable) { destroyMethodName = CLOSE_METHOD_NAME; } else { diff --git a/spring-beans/src/test/java/org/springframework/beans/factory/support/RootBeanDefinitionTests.java b/spring-beans/src/test/java/org/springframework/beans/factory/support/RootBeanDefinitionTests.java index 056597442ba..3d71b6d2420 100644 --- a/spring-beans/src/test/java/org/springframework/beans/factory/support/RootBeanDefinitionTests.java +++ b/spring-beans/src/test/java/org/springframework/beans/factory/support/RootBeanDefinitionTests.java @@ -17,6 +17,7 @@ package org.springframework.beans.factory.support; import java.lang.reflect.Method; +import java.util.concurrent.ExecutorService; import org.junit.jupiter.api.Test; @@ -59,13 +60,38 @@ class RootBeanDefinitionTests { } @Test - void resolveDestroyMethodWithMatchingCandidateReplacedInferredVaue() { + void resolveDestroyMethodWithMatchingCandidateReplacedForCloseMethod() { RootBeanDefinition beanDefinition = new RootBeanDefinition(BeanWithCloseMethod.class); beanDefinition.setDestroyMethodName(AbstractBeanDefinition.INFER_METHOD); beanDefinition.resolveDestroyMethodIfNecessary(); assertThat(beanDefinition.getDestroyMethodNames()).containsExactly("close"); } + @Test + void resolveDestroyMethodWithMatchingCandidateReplacedForShutdownMethod() { + RootBeanDefinition beanDefinition = new RootBeanDefinition(BeanWithShutdownMethod.class); + beanDefinition.setDestroyMethodName(AbstractBeanDefinition.INFER_METHOD); + beanDefinition.resolveDestroyMethodIfNecessary(); + assertThat(beanDefinition.getDestroyMethodNames()).containsExactly("shutdown"); + } + + @Test + void resolveDestroyMethodWithMatchingCandidateReplacedForExecutorService() { + RootBeanDefinition beanDefinition = new RootBeanDefinition(BeanImplementingExecutorService.class); + beanDefinition.setDestroyMethodName(AbstractBeanDefinition.INFER_METHOD); + beanDefinition.resolveDestroyMethodIfNecessary(); + assertThat(beanDefinition.getDestroyMethodNames()).containsExactly("shutdown"); + // even on JDK 19+ where the ExecutorService interface declares a default AutoCloseable implementation + } + + @Test + void resolveDestroyMethodWithMatchingCandidateReplacedForAutoCloseableExecutorService() { + RootBeanDefinition beanDefinition = new RootBeanDefinition(BeanImplementingExecutorServiceAndAutoCloseable.class); + beanDefinition.setDestroyMethodName(AbstractBeanDefinition.INFER_METHOD); + beanDefinition.resolveDestroyMethodIfNecessary(); + assertThat(beanDefinition.getDestroyMethodNames()).containsExactly("close"); + } + @Test void resolveDestroyMethodWithNoCandidateSetDestroyMethodNameToNull() { RootBeanDefinition beanDefinition = new RootBeanDefinition(BeanWithNoDestroyMethod.class); @@ -90,6 +116,25 @@ class RootBeanDefinitionTests { } + static class BeanWithShutdownMethod { + + public void shutdown() { + } + } + + + abstract static class BeanImplementingExecutorService implements ExecutorService { + } + + + abstract static class BeanImplementingExecutorServiceAndAutoCloseable implements ExecutorService, AutoCloseable { + + @Override + public void close() { + } + } + + static class BeanWithNoDestroyMethod { } diff --git a/spring-core/src/main/java/org/springframework/core/task/SimpleAsyncTaskExecutor.java b/spring-core/src/main/java/org/springframework/core/task/SimpleAsyncTaskExecutor.java index 39663c32ea2..ad598741187 100644 --- a/spring-core/src/main/java/org/springframework/core/task/SimpleAsyncTaskExecutor.java +++ b/spring-core/src/main/java/org/springframework/core/task/SimpleAsyncTaskExecutor.java @@ -87,6 +87,8 @@ public class SimpleAsyncTaskExecutor extends CustomizableThreadCreator private @Nullable Set activeThreads; + private boolean cancelRemainingTasksOnClose = false; + private boolean rejectTasksWhenLimitReached = false; private volatile boolean active = true; @@ -178,12 +180,33 @@ public class SimpleAsyncTaskExecutor extends CustomizableThreadCreator * @param timeout the timeout in milliseconds * @since 6.1 * @see #close() + * @see #setCancelRemainingTasksOnClose * @see org.springframework.scheduling.concurrent.ExecutorConfigurationSupport#setAwaitTerminationMillis */ public void setTaskTerminationTimeout(long timeout) { Assert.isTrue(timeout >= 0, "Timeout value must be >=0"); this.taskTerminationTimeout = timeout; - this.activeThreads = (timeout > 0 ? ConcurrentHashMap.newKeySet() : null); + trackActiveThreadsIfNecessary(); + } + + /** + * Specify whether to cancel remaining tasks on close: that is, whether to + * interrupt any active threads at the time of the {@link #close()} call. + *

The default is {@code false}, not tracking active threads at all or + * just interrupting any remaining threads that still have not finished after + * the specified {@link #setTaskTerminationTimeout taskTerminationTimeout}. + * Switch this to {@code true} for immediate interruption on close, either in + * combination with a subsequent termination timeout or without any waiting + * at all, depending on whether a {@code taskTerminationTimeout} has been + * specified as well. + * @since 6.2.11 + * @see #close() + * @see #setTaskTerminationTimeout + * @see org.springframework.scheduling.concurrent.ExecutorConfigurationSupport#setWaitForTasksToCompleteOnShutdown + */ + public void setCancelRemainingTasksOnClose(boolean cancelRemainingTasksOnClose) { + this.cancelRemainingTasksOnClose = cancelRemainingTasksOnClose; + trackActiveThreadsIfNecessary(); } /** @@ -243,6 +266,15 @@ public class SimpleAsyncTaskExecutor extends CustomizableThreadCreator return this.active; } + /** + * Track active threads only when a task termination timeout has been + * specified or interruption of remaining threads has been requested. + */ + private void trackActiveThreadsIfNecessary() { + this.activeThreads = (this.taskTerminationTimeout > 0 || this.cancelRemainingTasksOnClose ? + ConcurrentHashMap.newKeySet() : null); + } + /** * Executes the given task, within a concurrency throttle @@ -331,7 +363,7 @@ public class SimpleAsyncTaskExecutor extends CustomizableThreadCreator } /** - * This close methods tracks the termination of active threads if a concrete + * This close method tracks the termination of active threads if a concrete * {@link #setTaskTerminationTimeout task termination timeout} has been set. * Otherwise, it is not necessary to close this executor. * @since 6.1 @@ -342,17 +374,26 @@ public class SimpleAsyncTaskExecutor extends CustomizableThreadCreator this.active = false; Set threads = this.activeThreads; if (threads != null) { - synchronized (threads) { - try { - if (!threads.isEmpty()) { - threads.wait(this.taskTerminationTimeout); + if (this.cancelRemainingTasksOnClose) { + // Early interrupt for remaining tasks on close + threads.forEach(Thread::interrupt); + } + if (this.taskTerminationTimeout > 0) { + synchronized (threads) { + try { + if (!threads.isEmpty()) { + threads.wait(this.taskTerminationTimeout); + } + } + catch (InterruptedException ex) { + Thread.currentThread().interrupt(); } } - catch (InterruptedException ex) { - Thread.currentThread().interrupt(); + if (!this.cancelRemainingTasksOnClose) { + // Late interrupt for remaining tasks after timeout + threads.forEach(Thread::interrupt); } } - threads.forEach(Thread::interrupt); } } }