如何在Vert.x Verticle中使用自定义ThreadPoolExecutor?
我刚从Play框架转用Vert.x,之前在Play中通过自定义MessageDispatcherConfigurator搭配CustomThreadPoolExecutor实现了Actor系统的线程上下文传播(切换执行器或线程池内线程时保留上下文)。现在想在Vert.x的所有Verticle中用自己的CustomThreadPoolExecutor替代默认线程池,已经基于SPI的ExecutorServiceFactory实现了CustomExecutorServiceFactory,但不清楚部署Verticle时如何启用它。
附上已实现的代码:
CustomDispatcherConfigurator(Play环境下的实现)
public class CustomDispatcherConfigurator extends MessageDispatcherConfigurator { private final CustomDispatcher instance; public CustomDispatcherConfigurator(Config config, DispatcherPrerequisites prerequisites) { super(config, prerequisites); Config threadPoolConfig = config.getConfig("thread-pool-executor"); int fixedPoolSize = threadPoolConfig.getInt("fixed-pool-size"); instance = new CustomDispatcher( this, config.getString("id"), config.getInt("throughput"), Duration.create(config.getDuration("throughput-deadline-time", TimeUnit.NANOSECONDS), TimeUnit.NANOSECONDS), (id, threadFactory) -> () -> new CustomThreadPoolExecutor(fixedPoolSize, fixedPoolSize, threadPoolConfig.getDuration("keep-alive-time", TimeUnit.MILLISECONDS), TimeUnit.MILLISECONDS, new LinkedBlockingDeque<>(), new ThreadFactory() { private int threadId = 1; @Override public Thread newThread(@NotNull Runnable r) { Thread thread = new Thread(r); thread.setName(config.getString("name") + "-" + threadId++); return thread; } }), Duration.create(config.getDuration("shutdown-timeout", TimeUnit.MILLISECONDS), TimeUnit.MILLISECONDS) ); } @Override public MessageDispatcher dispatcher() { return instance; } } class CustomDispatcher extends Dispatcher { public CustomDispatcher(MessageDispatcherConfigurator _configurator, String id, int throughput, Duration throughputDeadlineTime, ExecutorServiceFactoryProvider executorServiceFactoryProvider, scala.concurrent.duration.FiniteDuration shutdownTimeout) { super(_configurator, id, throughput, throughputDeadlineTime, executorServiceFactoryProvider, shutdownTimeout); } }
CustomThreadPoolExecutor(上下文传播核心实现)
public class CustomThreadPoolExecutor extends ThreadPoolExecutor { public CustomThreadPoolExecutor(int corePoolSize, int maximumPoolSize, long keepAliveTime, @NotNull TimeUnit unit, @NotNull BlockingQueue<Runnable> workQueue) { super(corePoolSize, maximumPoolSize, keepAliveTime, unit, workQueue); } public CustomThreadPoolExecutor(int corePoolSize, int maximumPoolSize, long keepAliveTime, @NotNull TimeUnit unit, @NotNull BlockingQueue<Runnable> workQueue, @NotNull ThreadFactory threadFactory) { super(corePoolSize, maximumPoolSize, keepAliveTime, unit, workQueue, threadFactory); } public CustomThreadPoolExecutor(int corePoolSize, int maximumPoolSize, long keepAliveTime, @NotNull TimeUnit unit, @NotNull BlockingQueue<Runnable> workQueue, @NotNull RejectedExecutionHandler handler) { super(corePoolSize, maximumPoolSize, keepAliveTime, unit, workQueue, handler); } public CustomThreadPoolExecutor(int corePoolSize, int maximumPoolSize, long keepAliveTime, @NotNull TimeUnit unit, @NotNull BlockingQueue<Runnable> workQueue, @NotNull ThreadFactory threadFactory, @NotNull RejectedExecutionHandler handler) { super(corePoolSize, maximumPoolSize, keepAliveTime, unit, workQueue, threadFactory, handler); } @Override public <T> @NotNull Future<T> submit(@NotNull Callable<T> task) { return super.submit(ContextUtility.wrapWithContext(task)); } @Override public <T> @NotNull Future<T> submit(@NotNull Runnable task, T result) { return super.submit(ContextUtility.wrapWithContext(task), result); } @Override public @NotNull Future<?> submit(@NotNull Runnable task) { return super.submit(ContextUtility.wrapWithContext(task)); } @Override public void execute(@NotNull Runnable task) { super.execute(ContextUtility.wrapWithContext(task)); } }
CustomExecutorServiceFactory(Vert.x SPI实现)
public class CustomExecutorServiceFactory implements ExecutorServiceFactory { @Override public void init(VertxBuilder builder) { ExecutorServiceFactory.super.init(builder); } @Override public ExecutorService createExecutor(ThreadFactory threadFactory, Integer concurrency, Integer maxConcurrency) { return new CustomThreadPoolExecutor(concurrency, maxConcurrency, 10L, TimeUnit.MICROSECONDS, new LinkedBlockingQueue<>(), threadFactory); } }
要让Vert.x全局使用你的CustomExecutorServiceFactory,需通过SPI机制加载并配合部署配置生效:
1. 配置SPI加载文件
在项目的src/main/resources/META-INF/services目录下,创建名为io.vertx.core.spi.ExecutorServiceFactory的文件,文件内容为你的工厂类全限定名:
com.your.package.CustomExecutorServiceFactory
(替换为实际的包路径)
2. 全局生效配置(所有Verticle默认使用)
编程式启动Vert.x
创建Vertx实例时,通过VertxOptions配置线程池参数,Vert.x会自动加载SPI工厂:
public static void main(String[] args) { VertxOptions options = new VertxOptions() // 配置工作线程池大小,会传给createExecutor方法的concurrency参数 .setWorkerPoolSize(16) // 可选:配置事件循环线程池大小 .setEventLoopPoolSize(4); // Vert.x通过SPI加载自定义工厂创建线程池 Vertx vertx = Vertx.vertx(options); // 部署Verticle,默认使用全局配置的线程池 vertx.deployVerticle(new YourVerticle()); }
配置文件方式启动
若使用Vert.x CLI或配置文件启动,在vertx-options.json中配置线程池参数,SPI工厂会自动生效:
{ "workerPoolSize": 16, "eventLoopPoolSize": 4 }
启动命令:
vertx run com.your.package.YourVerticle -conf vertx-options.json
3. 单个Verticle单独配置(可选)
如果需要给特定Verticle配置独立线程池,部署时指定DeploymentOptions:
DeploymentOptions deploymentOptions = new DeploymentOptions() .setWorker(true) // 标记为工作线程Verticle .setWorkerPoolName("custom-pool") .setWorkerPoolSize(8); vertx.deployVerticle(new YourSpecificVerticle(), deploymentOptions);
此时自定义工厂会为该线程池创建CustomThreadPoolExecutor实例。
4. 验证生效
可在CustomThreadPoolExecutor的构造方法或execute方法中添加日志,启动后查看日志是否有自定义线程池的输出,或通过Thread.currentThread().getName()验证线程名称是否符合预期。
内容的提问来源于stack exchange,提问作者Anuj Singh

