You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.05.22 08:43:09