单元1 · Vert.x 快速入门
Reactive 工具包、事件循环、Hello World
- Vert.x:事件驱动、非阻塞、多语言响应式工具包。
- 创建:Vertx.vertx() + createHttpServer。
- Verticle:AbstractVerticle 生命周期,deployVerticle 部署。
实训1.1 Vert.x 特性
说明 Vert.x 的核心特性(事件驱动、非阻塞、多语言)。
Vert.x 是基于 JVM 的响应式工具包:事件循环、非阻塞 IO、支持 Java/Kotlin/JS。
// Vert.x 核心概念
// 1. Event Loop:事件循环线程处理 IO 事件
// 2. Non-blocking:非阻塞 IO,线程不等待
// 3. Polyglot:Java/Kotlin/JavaScript 多语言
// 4. 模块化:vertx-core、vertx-web 等模块
public class VertxIntro {
// 以上为要点
}
实训1.2 创建 Vert.x 应用
创建 Vert.x 实例并启动 HTTP 服务器。
Vertx.vertx() 创建实例;vertx.createHttpServer() 创建 HTTP 服务器。
import io.vertx.core.Vertx;
import io.vertx.core.http.HttpServer;
public class Main {
public static void main(String[] args) {
Vertx vertx = Vertx.vertx();
HttpServer server = vertx.createHttpServer();
server.requestHandler(req -> req.response()
.putHeader("content-type", "text/plain")
.end("Hello Vert.x"))
.listen(8080, http -> {
if (http.succeeded()) {
System.out.println("服务器启动:8080");
}
});
}
}
实训1.3 Verticle 组件
编写并部署一个 Verticle 组件。
AbstractVerticle 提供 start/stop 生命周期;vertx.deployVerticle 部署。
import io.vertx.core.AbstractVerticle;
import io.vertx.core.Vertx;
public class HelloVerticle extends AbstractVerticle {
@Override
public void start() {
vertx.createHttpServer()
.requestHandler(req -> req.response().end("Hello from Verticle"))
.listen(8080);
System.out.println("HelloVerticle 已部署");
}
public static void main(String[] args) {
Vertx vertx = Vertx.vertx();
vertx.deployVerticle(new HelloVerticle());
}
}
单元2 · 事件循环与线程模型
Event Loop、阻塞处理、Worker Verticle
- 事件循环:Event Loop 串行处理,禁止阻塞。
- 阻塞任务:executeBlocking 或 Worker Verticle。
实训2.1 事件循环
说明 Event Loop 线程的工作方式与注意事项。
Event Loop 串行处理事件;事件处理器中禁止阻塞,否则拖慢全部连接。
// Event Loop 特点
// 1. 每个 Event Loop 线程串行处理事件
// 2. 处理器必须快速返回,不能阻塞
// 3. 阻塞操作应交给 Worker 线程或 executeBlocking
public class EventLoopNote {
public void bad(io.vertx.core.http.HttpServerRequest req) throws Exception {
// 错误:在 Event Loop 上阻塞
// Thread.sleep(1000);
req.response().end("ok");
}
}
实训2.2 executeBlocking
使用 executeBlocking 将阻塞任务交给工作线程。
vertx.executeBlocking 执行阻塞任务,回调在 Event Loop 返回。
import io.vertx.core.Vertx;
import io.vertx.core.Future;
public class BlockingDemo {
public static void main(String[] args) {
Vertx vertx = Vertx.vertx();
vertx.executeBlocking(promise -> {
try {
Thread.sleep(500); // 模拟阻塞任务
promise.complete("任务完成");
} catch (InterruptedException e) {
promise.fail(e);
}
}).onSuccess(result -> System.out.println(result));
}
}
实训2.3 Worker Verticle
部署 Worker Verticle 处理阻塞任务。
deployVerticle 配置 worker=true,Verticle 运行在 Worker 线程池。
import io.vertx.core.AbstractVerticle;
import io.vertx.core.Promise;
import io.vertx.core.Vertx;
public class WorkerVerticle extends AbstractVerticle {
@Override
public void start(Promise<Void> startPromise) {
vertx.eventBus().consumer("task.queue", msg -> {
try {
Thread.sleep(300); // 阻塞任务
msg.reply("worker 处理:" + msg.body());
} catch (InterruptedException e) {
msg.fail(1, e.getMessage());
}
});
startPromise.complete();
}
public static void main(String[] args) {
Vertx vertx = Vertx.vertx();
// worker=true 部署为 Worker Verticle
vertx.deployVerticle(new WorkerVerticle(),
new io.vertx.core.DeploymentOptions().setWorker(true));
}
}
单元3 · Future 与 Promise
异步编程模型、组合与链式调用
- Future:onSuccess/onFailure 回调,succeededFuture 创建。
- 组合:compose 串行、all 并行。
- Promise:complete/fail 手动完成异步。
实训3.1 Future 基本使用
创建 Future 并链式处理成功失败。
Future.succeededFuture/failedFuture 创建;onSuccess/onFailure 回调。
import io.vertx.core.Future;
public class FutureDemo {
public static Future<String> fetch() {
return Future.succeededFuture("数据");
}
public static void main(String[] args) {
fetch()
.onSuccess(data -> System.out.println("成功:" + data))
.onFailure(err -> System.err.println("失败:" + err.getMessage()));
}
}
实训3.2 Future 组合
组合多个异步操作:串行与并行。
compose 串行依赖;all 并行组合多个 Future。
import io.vertx.core.Future;
import java.util.List;
public class ComposeDemo {
static Future<String> step1() { return Future.succeededFuture("A"); }
static Future<String> step2(String s) { return Future.succeededFuture(s + "B"); }
public static void main(String[] args) {
// 串行:step2 依赖 step1 结果
step1().compose(ComposeDemo::step2)
.onSuccess(r -> System.out.println("串行结果:" + r));
// 并行:两个 Future 同时执行
Future.all(List.of(step1(), step2("X")))
.onSuccess(cf -> System.out.println("并行完成:" + cf.list()));
}
}
实训3.3 Promise 手动完成
使用 Promise 手动控制异步结果完成时机。
Promise 提供 complete/fail 手动完成,future() 获取 Future。
import io.vertx.core.Promise;
import io.vertx.core.Future;
public class PromiseDemo {
public static Future<String> delayed(String value, long ms) {
Promise<String> promise = Promise.promise();
new Thread(() -> {
try {
Thread.sleep(ms);
promise.complete(value);
} catch (InterruptedException e) {
promise.fail(e);
}
}).start();
return promise.future();
}
public static void main(String[] args) {
delayed("Hello", 500).onSuccess(System.out::println);
}
}
单元4 · Event Bus
发布订阅、点对点、编解码
- Event Bus:send 点对点、publish 广播。
- 编解码:MessageCodec 自定义消息类型。
实训4.1 点对点消息
使用 Event Bus 发送点对点消息。
eventBus().send 点对点,consumer 消费;reply 回复。
import io.vertx.core.AbstractVerticle;
import io.vertx.core.Vertx;
public class PointToPoint extends AbstractVerticle {
@Override
public void start() {
vertx.eventBus().consumer("news", msg -> {
System.out.println("收到:" + msg.body());
msg.reply("已收到:" + msg.body());
});
vertx.setTimer(100, id -> {
vertx.eventBus().request("news", "Hello EventBus", reply -> {
if (reply.succeeded()) {
System.out.println("回复:" + reply.result().body());
}
});
});
}
}
实训4.2 发布订阅消息
使用 publish 广播消息给多个订阅者。
eventBus().publish 广播;所有订阅该地址的 consumer 都会收到。
import io.vertx.core.AbstractVerticle;
import io.vertx.core.Vertx;
public class PubSub extends AbstractVerticle {
@Override
public void start() {
// 订阅者 A
vertx.eventBus().consumer("alerts", msg ->
System.out.println("A 收到:" + msg.body()));
// 订阅者 B
vertx.eventBus().consumer("alerts", msg ->
System.out.println("B 收到:" + msg.body()));
vertx.setTimer(200, id -> {
// 广播:A、B 都会收到
vertx.eventBus().publish("alerts", "系统告警");
});
}
}
实训4.3 编解码器
注册自定义消息编解码器发送对象。
codec() 注册编解码器,send 时指定 codec 名称。
import io.vertx.core.AbstractVerticle;
import io.vertx.core.eventbus.MessageCodec;
import io.vertx.core.buffer.Buffer;
import java.io.*;
public class CodecDemo extends AbstractVerticle {
@Override
public void start() {
vertx.eventBus().registerDefaultCodec(User.class,
new MessageCodec<User, User>() {
@Override
public void encodeToWire(Buffer buffer, User user) {
buffer.appendString(user.name);
}
@Override
public User decodeFromWire(int pos, Buffer buffer) {
return new User(buffer.getString(pos));
}
@Override
public User transform(User user) { return user; }
@Override
public String name() { return "user-codec"; }
@Override
public byte systemCodecID() { return -1; }
});
vertx.eventBus().consumer("users", msg ->
System.out.println("收到对象:" + ((User) msg.body()).name));
vertx.setTimer(100, id ->
vertx.eventBus().send("users", new User("Tom")));
}
}
class User {
String name;
User(String name) { this.name = name; }
}
单元5 · Vert.x Web
Router、路由、静态资源
- Router:route().path() 路由匹配。
- 参数:pathParam/queryParam 提取。
- 静态:StaticHandler 提供静态资源。
实训5.1 Router 路由
创建 Router 处理不同路径的请求。
Router.router(vertx);route().path() 匹配路径,handler 处理。
import io.vertx.core.AbstractVerticle;
import io.vertx.ext.web.Router;
public class WebVerticle extends AbstractVerticle {
@Override
public void start() {
Router router = Router.router(vertx);
router.get("/hello").handler(ctx -> ctx.response().end("Hello"));
router.get("/api/users").handler(ctx -> ctx.response()
.putHeader("content-type", "application/json")
.end("[{"id":1,"name":"Tom"}]"));
vertx.createHttpServer()
.requestHandler(router)
.listen(8080);
}
}
实训5.2 路径参数
从请求路径中提取参数。
ctx.pathParam 获取 :id 参数;ctx.queryParam 获取查询参数。
import io.vertx.core.AbstractVerticle;
import io.vertx.ext.web.Router;
public class ParamVerticle extends AbstractVerticle {
@Override
public void start() {
Router router = Router.router(vertx);
router.get("/users/:id").handler(ctx -> {
String id = ctx.pathParam("id");
String verbose = ctx.queryParam("verbose").isEmpty()
? "false" : ctx.queryParam("verbose").get(0);
ctx.response().end("用户 id=" + id + ", verbose=" + verbose);
});
vertx.createHttpServer().requestHandler(router).listen(8080);
}
}
实训5.3 静态资源
配置静态资源目录提供 HTML/CSS 文件。
StaticHandler 提供静态资源;webroot 指定目录。
import io.vertx.core.AbstractVerticle;
import io.vertx.ext.web.Router;
import io.vertx.ext.web.handler.StaticHandler;
public class StaticVerticle extends AbstractVerticle {
@Override
public void start() {
Router router = Router.router(vertx);
// webroot 默认为 classpath:webroot
router.route("/static/*").handler(StaticHandler.create()
.setWebRoot("webroot")
.setCachingEnabled(true));
vertx.createHttpServer().requestHandler(router).listen(8080);
}
}
单元6 · HTTP 客户端
WebClient、异步请求
- WebClient:异步 HTTP 客户端,Future 回调。
- JSON:sendJsonObject 发送请求体。
- 超时:timeout 毫秒 + recover 重试。
实训6.1 WebClient 基本请求
使用 WebClient 发起 GET 请求。
WebClient.create(vertx);get 返回 Future
import io.vertx.core.Vertx;
import io.vertx.ext.web.client.WebClient;
import io.vertx.ext.web.client.WebClientOptions;
public class ClientDemo {
public static void main(String[] args) {
Vertx vertx = Vertx.vertx();
WebClient client = WebClient.create(vertx,
new WebClientOptions().setDefaultHost("api.helpme.com")
.setDefaultPort(443)
.setSsl(true));
client.get("/users/1")
.send()
.onSuccess(resp -> {
System.out.println("状态码:" + resp.statusCode());
System.out.println("响应体:" + resp.bodyAsString());
})
.onFailure(err -> System.err.println("失败:" + err.getMessage()));
}
}
实训6.2 POST JSON 请求
使用 WebClient 发送 JSON 请求体。
sendJson 序列化对象;Content-Type 自动设置。
import io.vertx.core.Vertx;
import io.vertx.ext.web.client.WebClient;
import io.vertx.core.json.JsonObject;
public class PostDemo {
public static void main(String[] args) {
Vertx vertx = Vertx.vertx();
WebClient client = WebClient.create(vertx,
new WebClientOptions().setDefaultHost("api.helpme.com").setSsl(true));
JsonObject user = new JsonObject()
.put("name", "Tom")
.put("age", 25);
client.post("/api/users")
.sendJsonObject(user)
.onSuccess(resp -> System.out.println("创建成功:" + resp.bodyAsString()))
.onFailure(err -> System.err.println("失败:" + err.getMessage()));
}
}
实训6.3 超时与重试
为客户端请求配置超时时间。
timeout 毫秒超时;send 返回 Future 可组合重试。
import io.vertx.core.Vertx;
import io.vertx.ext.web.client.WebClient;
import io.vertx.ext.web.client.WebClientOptions;
import io.vertx.ext.web.client.HttpRequest;
public class RetryDemo {
public static void main(String[] args) {
Vertx vertx = Vertx.vertx();
WebClient client = WebClient.create(vertx,
new WebClientOptions().setDefaultHost("api.helpme.com").setSsl(true));
HttpRequest<String> req = client.get("/users/1").timeout(3000);
req.send()
.recover(err -> {
System.out.println("第一次失败,重试");
return req.send();
})
.onSuccess(resp -> System.out.println(resp.bodyAsString()));
}
}
单元7 · 数据库访问
JDBC 客户端、异步 SQL
- JDBC:JDBCClient 异步 SQL。
- 参数化:queryWithParams 防注入。
- 事务:setAutoCommit(false) + commit/rollback。
实训7.1 JDBC 客户端
使用 JDBCClient 执行异步 SQL 查询。
JDBCClient.create(vertx, config);query 返回 Future
import io.vertx.core.Vertx;
import io.vertx.ext.jdbc.JDBCClient;
import io.vertx.ext.sql.SQLConnection;
public class JdbcDemo {
public static void main(String[] args) {
Vertx vertx = Vertx.vertx();
io.vertx.core.json.JsonObject config = new io.vertx.core.json.JsonObject()
.put("url", "jdbc:mysql://127.0.0.1:3306/helpme")
.put("driver_class", "com.mysql.cj.jdbc.Driver")
.put("user", "root")
.put("password", "123456");
JDBCClient client = JDBCClient.createShared(vertx, config);
client.getConnection(conn -> {
if (conn.failed()) {
System.err.println("连接失败");
return;
}
conn.result().query("SELECT id, name FROM users", rs -> {
rs.result().getRows().forEach(row ->
System.out.println(row.encode()));
conn.result().close();
});
});
}
}
实训7.2 参数化查询
使用参数化查询防止 SQL 注入。
queryWithParams 传入 ? 参数数组。
import io.vertx.core.Vertx;
import io.vertx.ext.jdbc.JDBCClient;
public class ParamQuery {
public static void main(String[] args) {
Vertx vertx = Vertx.vertx();
JDBCClient client = JDBCClient.createShared(vertx, config());
client.getConnection(conn -> {
String sql = "SELECT * FROM users WHERE age > ? AND name LIKE ?";
io.vertx.core.json.JsonArray params = new io.vertx.core.json.JsonArray()
.add(18).add("%T%");
conn.result().queryWithParams(sql, params, rs -> {
rs.result().getRows().forEach(System.out::println);
conn.result().close();
});
});
}
static io.vertx.core.json.JsonObject config() {
return new io.vertx.core.json.JsonObject()
.put("url", "jdbc:mysql://127.0.0.1:3306/helpme")
.put("driver_class", "com.mysql.cj.jdbc.Driver")
.put("user", "root")
.put("password", "123456");
}
}
实训7.3 事务处理
执行数据库事务保证数据一致性。
setAutoCommit(false) + commit/rollback 手动管理事务。
import io.vertx.ext.sql.SQLConnection;
public void transaction(SQLConnection conn) {
conn.setAutoCommit(false, res -> {
if (res.failed()) { return; }
conn.execute("UPDATE account SET balance = balance - 100 WHERE id = 1", u1 -> {
conn.execute("UPDATE account SET balance = balance + 100 WHERE id = 2", u2 -> {
if (u2.succeeded()) {
conn.commit(c -> {
System.out.println("事务提交成功");
conn.close();
});
} else {
conn.rollback(r -> {
System.out.println("事务回滚");
conn.close();
});
}
});
});
});
}
单元8 · WebSocket
WebSocket 服务端与客户端
- WebSocket:toWebSocket 升级协议。
- SockJS:浏览器兼容降级方案。
实训8.1 WebSocket 服务端
创建 WebSocket 服务端处理连接与消息。
server.webSocketHandler 接收连接;frameHandler 处理消息帧。
import io.vertx.core.AbstractVerticle;
import io.vertx.ext.web.Router;
import io.vertx.ext.web.handler.sockjs.SockJSHandler;
public class WsServerVerticle extends AbstractVerticle {
@Override
public void start() {
Router router = Router.router(vertx);
router.route("/ws/*").handler(ctx -> {
if (ctx.request().isUpgrade()) {
ctx.request().toWebSocket(ws -> {
if (ws.succeeded()) {
io.vertx.core.http.ServerWebSocket socket = ws.result();
System.out.println("连接:" + socket.remoteAddress());
socket.handler(buffer ->
socket.writeTextMessage("echo: " + buffer));
socket.closeHandler(v ->
System.out.println("连接关闭"));
}
});
} else {
ctx.response().end("请使用 WebSocket 协议");
}
});
vertx.createHttpServer().requestHandler(router).listen(8080);
}
}
实训8.2 WebSocket 客户端
编写 WebSocket 客户端连接服务端。
client.webSocket 发起连接;writeTextMessage 发送消息。
import io.vertx.core.Vertx;
import io.vertx.core.http.WebSocket;
public class WsClientDemo {
public static void main(String[] args) {
Vertx vertx = Vertx.vertx();
vertx.createHttpClient()
.webSocket(8080, "127.0.0.1", "/ws", ws -> {
if (ws.failed()) {
System.err.println("连接失败");
return;
}
WebSocket socket = ws.result();
socket.handler(buffer ->
System.out.println("收到:" + buffer));
socket.writeTextMessage("Hello WebSocket");
});
}
}
实训8.3 SockJS 兼容
使用 SockJSHandler 提供浏览器兼容的 WebSocket。
SockJS 支持降级轮询,浏览器兼容性好。
import io.vertx.core.AbstractVerticle;
import io.vertx.ext.web.Router;
import io.vertx.ext.web.handler.sockjs.SockJSHandler;
import io.vertx.ext.web.handler.sockjs.SockJSHandlerOptions;
public class SockJSVerticle extends AbstractVerticle {
@Override
public void start() {
SockJSHandler sockJS = SockJSHandler.create(vertx,
new SockJSHandlerOptions().setHeartbeatInterval(2000));
sockJS.socketHandler(sock -> {
System.out.println("SockJS 连接");
sock.handler(buffer -> sock.write(buffer));
sock.endHandler(v -> System.out.println("SockJS 断开"));
});
Router router = Router.router(vertx);
router.route("/chat/*").subRouter(sockJS);
vertx.createHttpServer().requestHandler(router).listen(8080);
}
}
单元9 · 集群与高可用
集群管理器、事件总线集群
- 集群:Hazelcast 集群管理器 + clusteredVertx。
- HA:setHAEnabled 故障转移。
实训9.1 集群模式
以集群模式启动 Vert.x 应用。
setClusterManager + vertx.clusteredVertx 启动集群;事件总线跨节点。
import io.vertx.core.Vertx;
import io.vertx.core.VertxOptions;
import io.vertx.spi.cluster.hazelcast.HazelcastClusterManager;
public class ClusterDemo {
public static void main(String[] args) {
VertxOptions options = new VertxOptions()
.setClusterManager(new HazelcastClusterManager());
Vertx.clusteredVertx(options, res -> {
if (res.succeeded()) {
Vertx vertx = res.result();
System.out.println("集群模式启动,节点数:" +
vertx.clusterVertxCount().result());
}
});
}
}
实训9.2 集群事件总线
在集群模式下跨节点发送消息。
集群模式下 send 可路由到其他节点的 consumer。
import io.vertx.core.AbstractVerticle;
import io.vertx.core.Vertx;
import io.vertx.core.VertxOptions;
import io.vertx.spi.cluster.hazelcast.HazelcastClusterManager;
public class ClusterMessaging extends AbstractVerticle {
@Override
public void start() {
vertx.eventBus().consumer("cluster.msg", msg ->
System.out.println("节点收到:" + msg.body() + " 来自 " + msg.replyAddress()));
vertx.setTimer(1000, id ->
vertx.eventBus().publish("cluster.msg", "跨节点消息"));
}
}
实训9.3 高可用部署
启用 Vert.x 高可用模式自动故障转移。
HA 模式:节点宕机后其他节点重新部署其 Verticle。
import io.vertx.core.Vertx;
import io.vertx.core.VertxOptions;
import io.vertx.core.DeploymentOptions;
import io.vertx.spi.cluster.hazelcast.HazelcastClusterManager;
public class HADemo {
public static void main(String[] args) {
VertxOptions options = new VertxOptions()
.setClusterManager(new HazelcastClusterManager())
.setHAEnabled(true);
Vertx.clusteredVertx(options, res -> {
if (res.succeeded()) {
res.result().deployVerticle("com.helpme.HelloVerticle",
new DeploymentOptions().setHa(true));
}
});
}
}
单元10 · 定时任务与 RxJava
setTimer、RxJava 集成
- 定时器:setTimer 一次性、setPeriodic 周期。
- RxJava:rxjava3 响应式 API。
实训10.1 setTimer 定时器
使用 setTimer 与 setPeriodic 定时执行任务。
setTimer 一次性延迟;setPeriodic 周期执行,返回定时器 ID。
import io.vertx.core.Vertx;
public class TimerDemo {
public static void main(String[] args) {
Vertx vertx = Vertx.vertx();
// 一次性延迟 2 秒
vertx.setTimer(2000, id ->
System.out.println("一次性定时器触发"));
// 周期任务:每 1 秒执行一次
long timerId = vertx.setPeriodic(1000, id ->
System.out.println("周期任务:" + System.currentTimeMillis()));
// 5 秒后取消周期任务
vertx.setTimer(5000, id -> vertx.cancelTimer(timerId));
}
}
实训10.2 RxJava 响应式
使用 RxJava 3 编写响应式代码。
io.vertx.rxjava3 包提供 Rx 化 API;Observable 订阅事件。
import io.vertx.rxjava3.core.Vertx;
import io.vertx.rxjava3.core.http.HttpServer;
public class RxDemo {
public static void main(String[] args) {
Vertx vertx = Vertx.rxVertx();
HttpServer server = vertx.createHttpServer();
server.requestHandler(req -> req.response().end("Rx Hello"))
.rxListen(8080)
.subscribe(
s -> System.out.println("Rx 服务器启动"),
err -> System.err.println("启动失败:" + err));
}
}
实训10.3 Rx 组合操作
使用 RxJava 组合多个异步结果。
zip 并行合并、flatMap 串行依赖。
import io.vertx.rxjava3.core.Vertx;
import io.reactivex.rxjava3.core.Single;
public class RxCompose {
public static void main(String[] args) {
Vertx vertx = Vertx.rxVertx();
Single<String> a = Single.just("A");
Single<String> b = Single.just("B");
// zip 并行合并
Single.zip(a, b, (x, y) -> x + y)
.subscribe(r -> System.out.println("zip 结果:" + r));
}
}
单元11 · 部署与配置
配置文件、环境变量、容器化
- 配置:fileSystem 读 JSON 配置。
- 部署:DeploymentOptions.setConfig。
- 容器:fat jar + Dockerfile。
实训11.1 配置文件读取
读取 JSON 配置文件中的参数。
vertx.fileSystem().readFile 读取;JsonObject 解析。
import io.vertx.core.AbstractVerticle;
import io.vertx.core.json.JsonObject;
public class ConfigVerticle extends AbstractVerticle {
@Override
public void start() {
vertx.fileSystem().readFile("config.json", res -> {
if (res.succeeded()) {
JsonObject config = new JsonObject(res.result());
System.out.println("app.name=" + config.getString("app.name"));
System.out.println("server.port=" + config.getInteger("server.port"));
}
});
}
}
实训11.2 部署配置
通过部署选项传递配置参数。
DeploymentOptions().setConfig() 传入 JSON 配置,config() 读取。
import io.vertx.core.AbstractVerticle;
import io.vertx.core.Vertx;
import io.vertx.core.DeploymentOptions;
import io.vertx.core.json.JsonObject;
public class DeployConfig extends AbstractVerticle {
@Override
public void start() {
JsonObject config = config();
System.out.println("name=" + config.getString("name", "default"));
}
public static void main(String[] args) {
Vertx vertx = Vertx.vertx();
DeploymentOptions options = new DeploymentOptions()
.setConfig(new JsonObject().put("name", "Helpme"));
vertx.deployVerticle(new DeployConfig(), options);
}
}
实训11.3 容器化部署
将 Vert.x 应用打包为 Docker 镜像。
fat jar + Dockerfile 多阶段构建,镜像轻量。
# 打包 fat jar
mvn package
# Dockerfile
FROM eclipse-temurin:17-jre-alpine
WORKDIR /app
COPY target/vertx-demo-1.0.0-fat.jar app.jar
EXPOSE 8080
ENTRYPOINT ["java", "-jar", "app.jar"]
# 构建并运行
docker build -t vertx-demo .
docker run -p 8080:8080 vertx-demo
单元12 · 综合项目实训
综合运用所学知识完成项目
- 综合应用:WebSocket + Event Bus + JDBC 组合。
实训12.1 聊天室 - 服务端
实现多房间 WebSocket 聊天室服务端。
综合运用 Router、WebSocket、Event Bus 实现广播。
import io.vertx.core.AbstractVerticle;
import io.vertx.core.eventbus.EventBus;
import io.vertx.ext.web.Router;
public class ChatVerticle extends AbstractVerticle {
@Override
public void start() {
EventBus eventBus = vertx.eventBus();
Router router = Router.router(vertx);
router.route("/chat/*").handler(ctx -> {
ctx.request().toWebSocket(ws -> {
if (ws.succeeded()) {
io.vertx.core.http.ServerWebSocket socket = ws.result();
String room = socket.query().substring(1);
eventBus.consumer("room." + room, msg -> {
if (msg.body().equals(socket.textHandlerID())) return;
socket.writeTextMessage((String) msg.body());
}, res -> {});
socket.handler(buffer -> {
String text = buffer.toString();
// 广播给房间内其他人
eventBus.publish("room." + room, text);
});
}
});
});
vertx.createHttpServer().requestHandler(router).listen(8080);
}
}
实训12.2 聊天室 - HTTP 网关
为聊天室提供 HTTP 网关接口。
综合运用 Router 路由、参数校验、JSON 响应。
import io.vertx.core.AbstractVerticle;
import io.vertx.core.json.JsonObject;
import io.vertx.ext.web.Router;
import io.vertx.ext.web.handler.BodyHandler;
public class ChatGateway extends AbstractVerticle {
@Override
public void start() {
Router router = Router.router(vertx);
router.route().handler(BodyHandler.create());
// 创建房间
router.post("/api/rooms").handler(ctx -> {
JsonObject body = ctx.body().asJsonObject();
String roomName = body.getString("name");
if (roomName == null || roomName.isBlank()) {
ctx.response().setStatusCode(400).end("{"error":"name required"}");
return;
}
ctx.response().putHeader("content-type", "application/json")
.end(new JsonObject().put("room", roomName).encode());
});
// 房间列表
router.get("/api/rooms").handler(ctx ->
ctx.response().putHeader("content-type", "application/json")
.end("["java","web","game"]"));
vertx.createHttpServer().requestHandler(router).listen(8080);
}
}
实训12.3 聊天室 - 消息持久化
将聊天消息异步写入数据库。
综合运用 Event Bus + JDBC 异步写入。
import io.vertx.core.AbstractVerticle;
import io.vertx.core.eventbus.Message;
import io.vertx.ext.jdbc.JDBCClient;
import io.vertx.core.json.JsonObject;
public class MessageStore extends AbstractVerticle {
@Override
public void start() {
JDBCClient client = JDBCClient.createShared(vertx,
new JsonObject()
.put("url", "jdbc:mysql://127.0.0.1:3306/chat")
.put("driver_class", "com.mysql.cj.jdbc.Driver")
.put("user", "root")
.put("password", "123456"));
vertx.eventBus().consumer("chat.save", (Message<JsonObject> msg) -> {
JsonObject data = msg.body();
String sql = "INSERT INTO messages(room, content, ts) VALUES(?, ?, NOW())";
io.vertx.core.json.JsonArray params = new io.vertx.core.json.JsonArray()
.add(data.getString("room"))
.add(data.getString("content"));
client.getConnection(conn -> {
conn.result().queryWithParams(sql, params, rs -> {
System.out.println("消息已保存");
conn.result().close();
});
});
});
}
}