You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

Vert.x WorkerPool为何用shutdownNow而非shutdown?能否改为shutdown实现优雅停机?

Quarkus 3.8.6 优雅停机:替换Vert.x WorkerPool的shutdownNow为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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.13 12:04:50