Vert.x Event Bus阻塞回复问题:双Verticle消息收发延迟回复代码咨询
Hey,我帮你把这段Vert.x的Event Bus收发代码整理成规范格式,还补全了接收端的实现(符合你说的休眠后回复的逻辑),同时补充了完整的错误处理👇
Vert.x Event Bus 发送/接收 Verticle 示例
1. Sender 发送端 Verticle
这个Verticle会每隔1秒向Event Bus的test地址发送一个JsonObject:
public class Sender extends AbstractVerticle { @Override public void start() throws Exception { final EventBus eventBus = this.vertx.eventBus(); // 每隔1秒发送一次消息 this.vertx.setPeriodic(1000, handler -> { JsonObject message = new JsonObject().put("d", "data from sender"); eventBus.<JsonObject>send("test", message, this::handleReply); }); } // 处理接收端的回复 private void handleReply(AsyncResult<Message<JsonObject>> result) { if (result.succeeded()) { JsonObject reply = result.result().body(); System.out.println("Received reply: " + reply.encodePrettily()); } else { System.err.println("Failed to get reply: " + result.cause().getMessage()); } } }
2. Receiver 接收端 Verticle
这个Verticle监听test地址,收到消息后**异步延迟(模拟休眠)**再回复(划重点:绝对不能在Event Loop线程里用Thread.sleep(),会阻塞整个Vert.x实例):
public class Receiver extends AbstractVerticle { @Override public void start() throws Exception { final EventBus eventBus = this.vertx.eventBus(); // 监听test地址的消息 eventBus.<JsonObject>consumer("test", message -> { JsonObject receivedMsg = message.body(); System.out.println("Received message: " + receivedMsg.encodePrettily()); // 模拟1秒休眠(用异步定时器,不阻塞Event Loop) vertx.setTimer(1000, timerHandler -> { // 构造回复消息 JsonObject reply = new JsonObject().put("status", "success").put("replyData", "processed: " + receivedMsg.getString("d")); // 回复发送端 message.reply(reply); }); }); } }
关键注意事项
- 绝对不要在Vert.x的Event Loop线程中调用阻塞方法(比如
Thread.sleep()):这会导致整个Vert.x实例的事件处理停滞,要用vertx.setTimer()这类异步API来模拟延迟操作。 - Event Bus的地址是全局的,只要在同一个Vert.x实例中,不同Verticle可以通过地址互通。
- 消息发送的
send()方法是请求-响应模式,publish()是广播模式,这里用send()对应接收端的reply()刚好匹配。
内容的提问来源于stack exchange,提问作者Almas Abdrazak
相关产品推荐
相关产品推荐

