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

Vertx:基于Event Bus实现请求-订阅式响应流方案问询

问题

我想实现一个Verticle,通过Event Bus以长运行流的形式响应请求。Event Bus本身支持请求/响应和发布/订阅两种模式,我想把它们结合成一种「请求-订阅」模式——发起一次请求后,能持续接收多个响应。

常规的请求/响应模式是这样的(单次响应):

eventbus.request("address", myRequest).onSuccess(msg -> processResponse(msg)); //single response

我想要的效果是发起一次请求后,能持续处理后续的多个响应:

eventbus.request("address", myRequest).handle(msg -> processNewResponse(msg)); // many responses over time

我之前用两步法实现过:先通过请求/响应模式拿到临时流地址,再订阅这个地址,代码如下:

eventbus.request("address", myRequest)
    .onSuccess(msg -> eventbus.consumer(msg.body().getTempStreamAddress(), msg2 -> processNewResponse(msg2)))

但这种方式需要手动管理临时地址的创建和订阅,有点繁琐,想知道有没有更简便的实现方式?

解决方案

有两种更简洁的方式可以实现你要的「请求-订阅」模式,不需要手动管理临时地址:

1. 利用Message.reply()多次回复

Vert.x的Event Bus允许处理请求的Verticle对同一个请求消息进行多次回复,发起请求的一端可以通过Handler持续接收这些回复,直到发送端主动结束流。

实现示例:

请求端代码:

eventbus.request("stream-address", initialRequest, ar -> {
    if (ar.succeeded()) {
        Message<StreamData> replyMsg = ar.result();
        // 注册handler接收后续所有回复
        replyMsg.handler(streamMsg -> {
            processNewResponse(streamMsg.body());
            // 收到结束信号后清理资源
            if (streamMsg.body().isEndOfStream()) {
                replyMsg.endHandler(v -> {});
            }
        });
    }
});

响应端(Verticle)代码:

eventbus.consumer("stream-address", msg -> {
    // 首次回复确认连接建立
    msg.reply(new StreamData("initial response"));
    
    // 模拟持续发送流数据
    vertx.setPeriodic(1000, timerId -> {
        msg.reply(new StreamData("stream data " + System.currentTimeMillis()));
        
        // 模拟流结束条件
        if (someEndCondition()) {
            msg.reply(new StreamData("end", true));
            vertx.cancelTimer(timerId);
        }
    });
});

这种方式直接复用请求的消息通道,不需要额外创建临时地址,所有回复都会通过同一个Message的handler传递。

2. 使用EventBus.send()结合一次性订阅

如果需要更灵活的流控制,可以在请求时附带一个自动生成的临时消费者地址,请求失败或流结束时自动清理订阅:

实现示例:

请求端代码:

// 生成唯一临时地址并创建消费者
String tempStreamAddr = eventbus.generateUniqueAddress();
Consumer<Message<StreamData>> streamConsumer = eventbus.consumer(tempStreamAddr, msg -> {
    processNewResponse(msg.body());
    // 流结束时自动注销消费者
    if (msg.body().isEndOfStream()) {
        streamConsumer.unregister();
    }
});

// 发送请求并附带临时地址
eventbus.send("stream-address", new RequestWithStreamAddr(tempStreamAddr), ar -> {
    if (ar.failed()) {
        streamConsumer.unregister(); // 请求失败时清理资源
    }
});

响应端(Verticle)代码:

eventbus.consumer("stream-address", msg -> {
    RequestWithStreamAddr request = msg.body();
    String streamAddr = request.getStreamAddr();
    
    // 持续向临时地址发送流数据
    vertx.setPeriodic(1000, timerId -> {
        eventbus.send(streamAddr, new StreamData("stream data " + System.currentTimeMillis()));
        
        if (someEndCondition()) {
            eventbus.send(streamAddr, new StreamData("end", true));
            vertx.cancelTimer(timerId);
        }
    });
});

这里用eventbus.generateUniqueAddress()自动生成临时地址,省去手动创建和管理的繁琐,且能自动清理资源。

注意事项

  • 无论采用哪种方式,都要在流结束时及时清理资源(取消handler、注销消费者),避免内存泄漏。
  • 如果需要处理背压,可以结合Vert.x的Flow API或RxJava扩展,控制流数据的发送速率。

内容的提问来源于stack exchange,提问作者firtydank

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.21 18:22:49