Vert.x WorkerPool为何用shutdownNow而非shutdown?能否改为shutdown实现优雅停机?
问题背景
调试Quarkus 3.8.6应用的优雅停机问题时,发现io.vertx.core.impl.WorkerPool在关闭阶段调用java.util.concurrent.ExecutorService#shutdownNow()终止SmallRye自定义工作线程池。由于该线程池处理Oracle数据库任务,Oracle JDBC驱动会响应shutdownNow()发送的中断信号,直接终止Socket IO并抛出异常,无法完成本应快速结束的数据库调用——这符合中断信号的语义(要求立即终止),但不符合优雅停机的需求。
解决方案
直接修改Vert.x源码不现实,可通过以下几种方式实现Vert.x使用shutdown()而非shutdownNow()关闭线程池:
1. 包装自定义ExecutorService
创建一个ExecutorService包装类,将shutdownNow()的逻辑替换为调用shutdown(),再把这个包装后的线程池配置给SmallRye。
示例代码:
import java.util.Collections; import java.util.List; import java.util.concurrent.*; public class GracefulShutdownExecutorWrapper implements ExecutorService { private final ExecutorService delegate; public GracefulShutdownExecutorWrapper(ExecutorService delegate) { this.delegate = delegate; } @Override public void shutdown() { delegate.shutdown(); } @Override public List<Runnable> shutdownNow() { // 替换为优雅关闭,不发送中断信号 delegate.shutdown(); return Collections.emptyList(); } // 以下所有ExecutorService接口方法均直接委托给delegate实现 @Override public boolean isShutdown() { return delegate.isShutdown(); } @Override public boolean isTerminated() { return delegate.isTerminated(); } @Override public boolean awaitTermination(long timeout, TimeUnit unit) throws InterruptedException { return delegate.awaitTermination(timeout, unit); } @Override public <T> Future<T> submit(Callable<T> task) { return delegate.submit(task); } @Override public <T> Future<T> submit(Runnable task, T result) { return delegate.submit(task, result); } @Override public Future<?> submit(Runnable task) { return delegate.submit(task); } @Override public <T> List<Future<T>> invokeAll(Collection<? extends Callable<T>> tasks) throws InterruptedException { return delegate.invokeAll(tasks); } @Override public <T> List<Future<T>> invokeAll(Collection<? extends Callable<T>> tasks, long timeout, TimeUnit unit) throws InterruptedException { return delegate.invokeAll(tasks, timeout, unit); } @Override public <T> T invokeAny(Collection<? extends Callable<T>> tasks) throws InterruptedException, ExecutionException { return delegate.invokeAny(tasks); } @Override public <T> T invokeAny(Collection<? extends Callable<T>> tasks, long timeout, TimeUnit unit) throws InterruptedException, ExecutionException, TimeoutException { return delegate.invokeAny(tasks, timeout, unit); } @Override public void execute(Runnable command) { delegate.execute(command); } }
之后通过Quarkus的CDI生产者,将包装后的ExecutorService注入为SmallRye的工作线程池:
import jakarta.enterprise.context.ApplicationScoped; import jakarta.enterprise.inject.Produces; import java.util.concurrent.ExecutorService; import java.util.concurrent.Executors; @ApplicationScoped public class WorkerPoolProducer { @Produces public ExecutorService customWorkerPool() { ExecutorService originalPool = Executors.newFixedThreadPool(10); // 根据你的配置调整线程池参数 return new GracefulShutdownExecutorWrapper(originalPool); } }
2. 利用Quarkus停机钩子提前处理
通过监听Quarkus的ShutdownEvent,在Vert.x触发WorkerPool关闭前,手动调用自定义线程池的shutdown()并等待任务完成。
示例代码:
import jakarta.enterprise.context.ApplicationScoped; import jakarta.enterprise.event.Observes; import io.quarkus.runtime.ShutdownEvent; import java.util.concurrent.ExecutorService; import java.util.concurrent.TimeUnit; @ApplicationScoped public class GracefulShutdownHandler { private final ExecutorService customWorkerPool; public GracefulShutdownHandler(ExecutorService customWorkerPool) { this.customWorkerPool = customWorkerPool; } void onShutdown(@Observes ShutdownEvent event) { customWorkerPool.shutdown(); try { // 等待任务完成,可根据业务场景调整超时时间 if (!customWorkerPool.awaitTermination(30, TimeUnit.SECONDS)) { // 超时后强制关闭(可选,根据业务需求决定) customWorkerPool.shutdownNow(); } } catch (InterruptedException e) { Thread.currentThread().interrupt(); customWorkerPool.shutdownNow(); } } }
3. 自定义Vert.x WorkerPool实现
如果上述方法无法满足需求,可以通过Quarkus的VertxConfigurationCustomizer扩展点,替换默认的WorkerPool创建逻辑,使用自定义的WorkerPool实现(将close()方法中的shutdownNow()改为shutdown())。
注意事项
- 为
shutdown()设置合理的超时时间,避免应用停机过程无限等待。 - 确保数据库任务具备幂等性,防止停机过程中任务重复执行(若存在重试机制)。
- 测试多种场景下的停机行为,比如长耗时数据库操作、并发任务等,验证优雅停机效果。
内容的提问来源于stack exchange,提问作者aldius

