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

如何在Vertx中每秒调度阻塞操作且不阻塞主事件循环?

Vertx实现非阻塞的周期性阻塞任务与多Worker线程协作

核心思路

  • 主Verticle(EventLoop类型)仅负责启动WorkerVerticle,不处理任何阻塞逻辑,避免占用主事件循环
  • 用两个独立的WorkerVerticle分别承担「状态获取」和「同步处理」的职责,各自运行在Worker线程池的独立线程中
  • 状态获取WorkerVerticle内部使用setPeriodic调度每秒一次的阻塞操作,通过EventBus将结果发送给同步器WorkerVerticle

代码实现

1. 主Verticle(Verticle1)

这个Verticle运行在主事件循环,只做部署工作,完全不涉及阻塞逻辑:

import io.vertx.core.AbstractVerticle;
import io.vertx.core.DeploymentOptions;

public class MainVerticle extends AbstractVerticle {
    @Override
    public void start() {
        // 部署状态获取WorkerVerticle,指定为worker模式
        vertx.deployVerticle(StateFetcherVerticle.class.getName(),
                new DeploymentOptions().setWorker(true));
        
        // 部署同步器WorkerVerticle,指定为worker模式
        vertx.deployVerticle(SynchronizerVerticle.class.getName(),
                new DeploymentOptions().setWorker(true));
    }
}

2. 状态获取WorkerVerticle

负责每秒执行阻塞的状态获取操作,并将结果通过EventBus发送给同步器:

import io.vertx.core.AbstractVerticle;

public class StateFetcherVerticle extends AbstractVerticle {
    private static final String EVENT_BUS_ADDRESS = "state.sync";

    @Override
    public void start() {
        // 在Worker线程中设置周期定时器,每秒触发一次
        vertx.setPeriodic(1000, timerId -> {
            // 执行阻塞的状态获取操作(比如数据库查询、外部API调用等)
            String state = fetchStateBlocking();
            
            // 将状态发送到EventBus,同步器会接收处理
            vertx.eventBus().send(EVENT_BUS_ADDRESS, state);
        });
    }

    // 模拟阻塞的状态获取方法
    private String fetchStateBlocking() {
        try {
            // 模拟阻塞操作,比如耗时的IO
            Thread.sleep(500);
        } catch (InterruptedException e) {
            Thread.currentThread().interrupt();
        }
        return "Current State: " + System.currentTimeMillis();
    }
}

3. 同步器WorkerVerticle

监听EventBus的状态消息,处理同步逻辑:

import io.vertx.core.AbstractVerticle;
import io.vertx.core.eventbus.Message;

public class SynchronizerVerticle extends AbstractVerticle {
    private static final String EVENT_BUS_ADDRESS = "state.sync";

    @Override
    public void start() {
        // 注册EventBus消费者,接收状态消息
        vertx.eventBus().consumer(EVENT_BUS_ADDRESS, this::handleStateMessage);
    }

    // 处理状态同步逻辑
    private void handleStateMessage(Message<String> message) {
        String state = message.body();
        // 这里执行同步操作,比如更新缓存、同步到外部系统等
        System.out.println("Synchronizing state: " + state);
    }
}

关键说明

  • WorkerVerticle的作用:每个WorkerVerticle默认运行在Worker线程池的独立线程中,所以状态获取和同步处理会分别占用两个Worker线程,符合你需要的「两个工作线程」的要求
  • 非阻塞保证:主事件循环只负责部署Verticle,所有阻塞操作都在Worker线程中执行,完全不会阻塞主事件循环
  • EventBus通信:Vertx的EventBus是线程安全的,WorkerVerticle之间通过EventBus传递消息,避免了手动线程同步的复杂问题

内容的提问来源于stack exchange,提问作者Zygmuntix

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.25 15:55:32