文章
Vert.x 教程(三)
Vertx核心组件Verticle的线程模型
在前两个教程中对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());
}
}
运行结果如下:
你可能会疑惑为什么没有端口占用冲突的错误。原因如下:
- **Vert.x 默认启用了 **
multiple event-loop threadsVert.x 运行在一个 异步、事件驱动的模型 上,它使用了 多线程的事件循环。当你调用:
vertx.deployVerticle(new HelloVerticle());
vertx.deployVerticle(new HelloVerticle());
```
你实际上启动了两个 HelloVerticle 实例,并且 它们在不同的 event loop 线程上运行。
2. Vert.x 会自动处理端口共享
在 Vert.x 的 HttpServer 实现中,多个 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或许更好。
📎 参考文章#
- vertx官网
- 《Vert.x in action》