返回文章列表

文章

Vert.x 教程(三)

Vertx核心组件Verticle的线程模型

目录
  1. 📝 Verticle的执行细节
  2. 实验一
  3. 实验二
  4. 实验三
  5. 处理阻塞事件
  6. Worker Verticle示例
  7. executeBlocking示例
  8. 📎 参考文章

在前两个教程中对Vert.x开发Web应用已经有一个大概的了解,本章内容主要是讲述一些更深层次的内容,Vert.x的基本执行单元:Verticle。

📝 Verticle的执行细节#

我们从一个简单的Verticle入手,示例代码如下:

实验一#

public class HelloVerticle extends AbstractVerticle {

    private final Logger logger = LoggerFactory.getLogger(HelloVerticle.class);
    private long counter = 1L;

    @Override
    public void start() {
        // 创建一个定时作业
        vertx.setPeriodic(5000, id -> logger.info("tick"));
        // 创建一个http服务
        vertx.createHttpServer()
            .requestHandler(req -> {
                logger.info("Request #{} from {}", counter++, req.remoteAddress().host());
                req.response().end("Hello!");
            }).listen(8888)
            .onSuccess(server -> logger.info("HTTP server started on port: {}", server.actualPort()))
            // Print the problem on failure
            .onFailure(throwable -> logger.info("HTTP server start failure: {}", throwable.getMessage()));
    }

    public static void main(String[] args) {
        Vertx vertx = Vertx.vertx();
        //部署Verticle最简单的方式
        vertx.deployVerticle(new HelloVerticle());
    }
}

这段代码中我们创建了一个5秒的定时任务,打印一段日志,然后创建一个http的服务监听端口8888。运行main方法,运行结果如下: 从日志中可以看到,定时任务和Http服务的请求处理都是由vert.x-eventloop-thread-0完成的。

实验二#

我们继续实验,修改代码如下所示,连续部署两个HelloVeticle

public class HelloVerticle extends AbstractVerticle {

    private final Logger logger = LoggerFactory.getLogger(HelloVerticle.class);
    private long counter = 1L;

    @Override
    public void start() {
        // 创建一个定时作业
        vertx.setPeriodic(5000, id -> logger.info("tick"));
        // 创建一个http服务器
        vertx.createHttpServer()
            .requestHandler(req -> {
                logger.info("Request #{} from {}", counter++, req.remoteAddress().host());
                req.response().end("Hello!");
            }).listen(8888)
            .onSuccess(server -> logger.info("HTTP server started on port: {}", server.actualPort()))
            // Print the problem on failure
            .onFailure(throwable -> logger.info("HTTP server start failure: {}", throwable.getMessage()));
    }

    public static void main(String[] args) {
        Vertx vertx = Vertx.vertx();
        //部署Verticle最简单的方式
        vertx.deployVerticle(new HelloVerticle());
        vertx.deployVerticle(new HelloVerticle());
    }
}

运行结果如下: 你可能会疑惑为什么没有端口占用冲突的错误。原因如下:

  1. **Vert.x 默认启用了 **multiple event-loop threads Vert.x 运行在一个 异步、事件驱动的模型 上,它使用了 多线程的事件循环。当你调用:

vertx.deployVerticle(new HelloVerticle()); vertx.deployVerticle(new HelloVerticle()); ``` 你实际上启动了两个 HelloVerticle 实例,并且 它们在不同的 event loop 线程上运行。 2. Vert.x 会自动处理端口共享Vert.xHttpServer 实现中,多个 Verticle 可以 共享同一个端口,而不会报端口占用错误。这是因为: - Vert.x 的 HttpServer 默认支持端口复用(端口共享)。 - 在多个 HelloVerticle 监听同一个端口时,Vert.x 并不会为每个实例都创建独立的服务器,而是让它们共享 同一个底层的 HTTP 服务器。 - 这样的话,多个 HelloVerticle 会被轮询调度来处理不同的请求。 3. 请求会被多个 Verticle 处理 由于多个 HelloVerticle 共享同一个 HTTP 服务器,当有请求进来时,它们会被 不同的 Verticle 实例轮流处理。这就像一个负载均衡器一样,分摊了 HTTP 请求的处理负载。

实验三#

新建一个无限循环的阻塞Verticle

public class BlockEventLoopVerticle extends AbstractVerticle {
    @Override
    public void start() {
        vertx.setTimer(1000, id -> {
            while (true)
                ;
        });
    }

    public static void main(String[] args) {
        Vertx vertx = Vertx.vertx();
        vertx.deployVerticle(new BlockEventLoopVerticle());
    }
}

运行代码,运行的结果如下: 正如日志中可以看到的Vert.x内置了一个阻塞检测的线程,当事件处理程序陷入一个无限循环时,日志中开始出现警告信息。经过几轮检测(默认是5秒),从日志中基本也能看出,警告中开始打印详细的调用堆栈信息,这样就可以快速定位代码中的问题。需要注意的是,这些仅仅是警告信息,事件检测线程并不能终止那些执行时间过长的任务。

处理阻塞事件#

有时候线程阻塞是不可避免的,例如我们使用的一些第三方库自带了其他的线程模型(像一些网络服务程序,Java中socket默认是线程阻塞的),对于这类情况Vert.x提供了两种方式来处理:使用Worker Verticle或者executeBlocking方法。

Worker Verticle示例#

public class WorkerVerticle extends AbstractVerticle {
    private final Logger logger = LoggerFactory.getLogger(WorkerVerticle.class);
    @Override
    public void start(){
        vertx.setPeriodic(10_000, id ->{

            try {
                logger.info("Zzz...");
                Thread.sleep(8000);
                logger.info("Up");
            } catch (InterruptedException e) {
                logger.error("Wrong", e);
            }
        });
    }

    public static void main(String[] args) {
        Vertx vertx = Vertx.vertx();
        // 定义部署配置
        DeploymentOptions opts = new DeploymentOptions()
            //设置2个实例
            .setInstances(2)
            //设置为Worker线程
            .setThreadingModel(ThreadingModel.WORKER);
        vertx.deployVerticle("com.verticle.hello.WorkerVerticle",opts);
    }
}

代码的执行结果如下: WorkerVerticle的两个实例每隔10秒会阻塞8秒,然后再去执行下一个周期的事件。

executeBlocking示例#

public class OffloadVerticle extends AbstractVerticle {

    private final Logger logger = LoggerFactory.getLogger(OffloadVerticle.class);

    @Override
    public void start() {
        vertx.setPeriodic(5000, id -> {
           logger.info("Tick");
           vertx.executeBlocking(this::blockingCode, this::resultHandler);
        });
    }

    private void blockingCode(Promise<String> promise) {
        try {
            logger.info("Blocking code running");
            Thread.sleep(4000);
            logger.info("Done!");
            promise.complete("Ok!");
        } catch (InterruptedException e) {
            promise.fail(e);
        }

    }

    private void resultHandler(AsyncResult<String> asyncResult) {
        if (asyncResult.succeeded()) {
            logger.info("Blocking code result: {}", asyncResult.result());
        }else {
            logger.error("Wrong", asyncResult.cause());
        }
    }

    public static void main(String[] args) {
        Vertx vertx = Vertx.vertx();
        vertx.deployVerticle(new OffloadVerticle());
    }
}

代码的运行结果如下 阻塞代码运行在worker线程上,非阻塞代码运行在eventloop线程上。 什么时候用使用WorkerVerticle,什么时候使用executeBlocking方法并没有一个准则,通常情况下使用WorkerVerticle是一个明智的选择,但是有时候阻塞代码的并不足以组合成一个执行单元时选择executeBlocking或许更好。

📎 参考文章#