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的
FlowAPI或RxJava扩展,控制流数据的发送速率。
内容的提问来源于stack exchange,提问作者firtydank
相关产品推荐
相关产品推荐

