文章
Vert.x 教程(五)
事件总线概念及使用
# 📝 主旨内容
## 什么是事件总线?
**事件总线**是一种异步地发送和接收消息的机制。[**消息**](/1bc5c571bb7d807699e1ed547c75b5f1?pvs=25)的发送和接收都有一个相应的目标地址。目标地址只是一个任意形式的字符串,如:abc.def.ghi或者abc-def-ghi,通常我们选择带有点号的前一种格式。
## 事件总线怎么使用?
### 1、点对点(Point-to-Point)模式
**模式描述**:
- 发送端通过 `eventBus.send("address", message)` 发送消息。
- 只有 **一个消费者** 会接收并处理该消息(即使有多个消费者订阅同一地址,只有一个能收到)。
### **✅ 优势**
✔ 适用于 **一对一任务分配**(单个消费者处理)。
✔ 适用于 **负载均衡**(多个消费者可以轮流接收请求)。
✔ 保证消息不会被多个消费者重复处理。
### **❌ 劣势**
✖ 不能广播消息,多个消费者不会同时收到消息。
✖ 如果接收者崩溃,可能导致消息丢失(除非使用 `DeliveryOptions` 设置重试机制)。
### **📌 适用场景**
- **任务分配**(比如分布式任务队列,多个消费者轮流处理任务)。
- **负载均衡**(多个实例分担消息处理任务,避免单点压力)。
- **单个服务处理请求**(比如请求处理到唯一的后端实例)。
示例代码如下:
```java
public class SenderVerticle extends AbstractVerticle {
@Override
public void start() {
vertx.setPeriodic(2000, id -> {
JsonObject message = new JsonObject().put("content", "Hello from Sender!");
vertx.eventBus().send("address.message", message);
System.out.println("Sent message: " + message.encode());
});
}
}
```
receiver示例代码
```java
public class ReceiverVerticle extends AbstractVerticle {
@Override
public void start() {
vertx.eventBus().consumer("address.message", message -> {
JsonObject body = (JsonObject) message.body();
System.out.println("Received message: " + body.encode());
});
}
}
```
测试启动入口
```java
public class MainVerticle {
public static void main(String[] args) {
Vertx vertx = Vertx.vertx();
//vertx.deployVerticle(new MainVerticle());
vertx.deployVerticle(new SenderVerticle());
vertx.deployVerticle(new ReceiverVerticle());
}
}
### 2、发布-订阅(Publish-Subscribe)模式
**模式描述**:
- 发送端通过 `eventBus.publish("address", message)` 发送消息。
- **所有订阅该地址的消费者都会收到消息**(广播模式)。
### **✅ 优势**
✔ **支持广播**,多个消费者可以同时接收同一条消息。
✔ **适用于事件驱动架构**,让多个组件监听同一事件。
✔ **解耦生产者与消费者**,发送者不关心谁接收了消息。
### **❌ 劣势**
✖ 无法保证所有订阅者都能成功处理消息(如果某个消费者崩溃,它不会影响其他消费者)。
✖ 无法负载均衡(所有订阅者都收到相同消息,不能自动分配任务)。
✖ 可能导致消息泛滥,影响系统性能(尤其是高频发布的情况下)。
### **📌 适用场景**
- **日志/监控系统**(多个服务监听相同日志事件)。
- **物联网(IoT)系统**(传感器数据广播给多个服务)。
- **用户通知系统**(用户 A 发送消息,所有关注 A 的人都收到)。
示例代码如下:
```java
public class PublisherVerticle extends AbstractVerticle {
@Override
public void start() {
vertx.setPeriodic(3000, id ->{
JsonObject message = new JsonObject().put("temp", 25 + Math.random() * 10);
vertx.eventBus().publish("publish.update", message);
System.out.println("Published message: " + message.encode());
});
}
}
Subscriber1示例代码
public class Subscriber1Verticle extends AbstractVerticle {
@Override
public void start() {
vertx.eventBus().consumer("publish.update", message -> {
System.out.println("Subscriber 1 received: " + message.body());
});
}
}
Subscriber2示例代码
public class Subscriber2Verticle extends AbstractVerticle {
@Override
public void start() {
vertx.eventBus().consumer("publish.update", message -> {
System.out.println("Subscriber 2 received: " + message.body());
});
}
}
MainVerticle示例代码
public class MainVerticle extends AbstractVerticle {
@Override
public void start() {
vertx.deployVerticle(new PublisherVerticle());
vertx.deployVerticle(new Subscriber1Verticle());
vertx.deployVerticle(new Subscriber2Verticle());
}
public static void main(String[] args) {
Vertx vertx = Vertx.vertx();
vertx.deployVerticle(new MainVerticle());
}
}
3、请求-应答(request-response)模式#
模式描述:
- 发送端通过 eventBus.request("address", message, replyHandler) 发送消息,并 等待接收端回复。
- 接收端使用 message.reply(responseMessage) 进行回复。
- 适用于需要 双向通信 的场景(类似于 HTTP 请求)。
### ✅ 优势
✔ 适用于 服务调用(类似于 HTTP 请求/响应模式)。
✔ 发送者可以 收到处理结果(不像 send() 和 publish() 只是单向发送)。
✔ 适用于 异步远程调用(RPC),可以让服务之间进行高效通信。
### ❌ 劣势
✖ 需要额外处理 超时 和 错误情况(请求可能失败或超时)。
✖ 需要 保证消费者正确处理消息并回复,否则发送者会一直等待。
✖ 额外的消息开销(需要 reply() 进行双向通信)。
### 📌 适用场景
- 微服务调用(A 调用 B 服务并等待返回数据)。
- 数据库查询(请求某个服务查询数据库,并返回查询结果)。
- 任务执行后返回结果(请求计算任务,等待结果返回)。
ServiceVerticle示例代码
```java
public class ServiceVerticle extends AbstractVerticle {
@Override
public void start() {
vertx.eventBus().consumer("service.request", message -> {
JsonObject request = (JsonObject) message.body();
System.out.println("Received request: " + request.encode());
JsonObject response = new JsonObject()
.put("status", "success")
.put("message", "Processed: " + request.getString("data"));
message.reply(response);
});
}
} ``` ClientVerticle示例代码
public class ClientVerticle extends AbstractVerticle {
@Override
public void start() {
JsonObject request = new JsonObject().put("data", "Hello Service!");
vertx.eventBus().request("service.request", request, reply -> {
if (reply.succeeded()) {
System.out.println("Response received: " + reply.result().body());
} else {
System.out.println("Failed to receive response");
}
});
}
}
MainVerticle示例代码
public class MainVerticle extends AbstractVerticle {
@Override
public void start() {
vertx.deployVerticle(new ServiceVerticle());
vertx.deployVerticle(new ClientVerticle());
}
public static void main(String[] args) {
Vertx vertx = Vertx.vertx();
vertx.deployVerticle(new MainVerticle());
}
}