Vertx EventBus+SpringBoot:HTTP请求超时问题及需求实现
Vert.x + SpringBoot EventBus 微服务问题:HTTP请求超时,无法返回发布的随机奖品
问题背景
首次尝试用Vert.x结合SpringBoot开发微服务,需求是每2秒通过EventBus发布随机奖品。目前发布消息的Verticle部署成功,日志显示消息正常发布、消费者能接收,但发送HTTP请求时返回超时错误。
错误信息
io.vertx.core.eventbus.ReplyException: Timed out after waiting 30000(ms) for a reply. address: __vertx.reply.1, repliedAddress: prizes
相关代码
PrizeBrokerVerticle(消息发布者)
@Component public class PrizeBrokerVerticle extends AbstractVerticle { public static final Logger LOGGER = Logger.getLogger(PrizeBrokerVerticle.class.getName()); private final EventBus eventBus; private final List<String> availablePrizes = List.of("Apple", "HP", "Dell"); public PrizeBrokerVerticle(EventBus eventBus) { this.eventBus = eventBus; } public void start() { vertx.setPeriodic(2000, timerId -> { String randomPrize = generateRandomPrize(); eventBus.publish("prizes", randomPrize); LOGGER.info("Prize published to EventBus"); }); } private String generateRandomPrize() { Random random = new Random(); int index = random.nextInt(availablePrizes.size()); return availablePrizes.get(index); } }
PrizeController(HTTP控制器)
@Controller public class PrizeController implements Function<RoutingContext, Uni<Void>> { public static final Logger LOGGER = Logger.getLogger(PrizeController.class.getName()); private final EventBus eventBus; public PrizeController(Router router, EventBus eventBus) { this.eventBus = eventBus; router.get("/api/prizes") .respond(this); } @Override public Uni<Void> apply(RoutingContext ctx) { eventBus.consumer("prizes", message -> { String prize = (String) message.body(); LOGGER.info("Received prize from EventBus: " + prize); }); return eventBus.request("prizes", "") .flatMap(reply -> { Object body = reply.body(); if (body != null) { return ctx.response() .setStatusCode(200) .putHeader("content-type", "application/json") .end((String) body); } else { LOGGER.warning("Received null response from EventBus"); return ctx.response() .setStatusCode(500) .end("Internal Server Error"); } }); } }
运行日志(正常发布与接收)
: Started Main in 0.546 seconds (process running for 0.801) 2023-12-19T16:56:46.291+01:00 INFO 39191 --- [ main] org.example.runner.VertxRunner : [] 2023-12-19T16:56:46.306+01:00 INFO 39191 --- [ntloop-thread-3] org.example.model.ConsumerPrize : ConsumerPrize verticle started 2023-12-19T16:56:46.307+01:00 INFO 39191 --- [ntloop-thread-1] org.example.verticle.MainVerticle : PrizeBrokerVerticle deployed with ID: 0fceae1f-5106-49da-9aa9-28dc70b14907 2023-12-19T16:56:46.307+01:00 INFO 39191 --- [ntloop-thread-1] org.example.verticle.MainVerticle : ConsumerPrize deployed with ID: a34b18d6-f621-4067-8a7f-6b19d667a3b8 2023-12-19T16:56:46.347+01:00 INFO 39191 --- [ntloop-thread-1] org.example.verticle.MainVerticle : Server Running on http://localhost8081 2023-12-19T16:56:46.347+01:00 INFO 39191 --- [ntloop-thread-0] org.example.runner.VertxRunner : 3d63fe6d-3186-4598-9111-2e626b7008f0 2023-12-19T16:56:52.306+01:00 INFO 39191 --- [ntloop-thread-2] o.example.verticle.PrizeBrokerVerticle : Prize published to EventBus 2023-12-19T16:56:52.306+01:00 INFO 39191 --- [ntloop-thread-1] org.example.controller.PrizeController : Received prize from EventBus: HP 2023-12-19T16:56:54.305+01:00 INFO 39191 --- [ntloop-thread-2] o.example.verticle.PrizeBrokerVerticle : Prize published to EventBus 2023-12-19T16:56:54.305+01:00 INFO 39191 --- [ntloop-thread-1] org.example.controller.PrizeController : Received prize from EventBus: Dell
尝试调整后的代码(仍有问题)
给消费者添加reply逻辑后,eventBus.request返回空内容,还是无法拿到发布的随机奖品:
eventBus.consumer("prizes", message -> { String prize = (String) message.body(); message.reply(prize); });
核心需求
实现HTTP请求/api/prizes返回EventBus上正在发布的随机奖品。
解决方案
问题根源
- EventBus语义混淆:
publish是广播模式,无回复机制;request是请求-响应模式,需要接收方主动回复,但你的发布者仅做广播,不会处理request的回复请求。 - 重复注册消费者:每次HTTP请求都会新增一个
prizes消费者,逻辑混乱且浪费资源。 - 请求地址错误:
request发送到广播地址prizes,该地址的消费者并非处理请求的角色,无法正确回复。
修正方案
方案1:缓存最新奖品(适配定时广播场景)
这是最贴合你需求的方案,直接缓存定时发布的最新奖品,HTTP请求直接取缓存值:
- 修改
PrizeBrokerVerticle,添加缓存逻辑:
@Component public class PrizeBrokerVerticle extends AbstractVerticle { // 原有代码保留 private String latestPrize; public void start() { vertx.setPeriodic(2000, timerId -> { String randomPrize = generateRandomPrize(); this.latestPrize = randomPrize; // 缓存最新奖品 eventBus.publish("prizes", randomPrize); LOGGER.info("Prize published to EventBus"); }); } // 提供获取最新奖品的方法 public String getLatestPrize() { return latestPrize; } // 原有代码保留 }
- 修改
PrizeController,直接读取缓存:
@Controller public class PrizeController implements Function<RoutingContext, Uni<Void>> { private final PrizeBrokerVerticle prizeBrokerVerticle; public PrizeController(Router router, PrizeBrokerVerticle prizeBrokerVerticle) { this.prizeBrokerVerticle = prizeBrokerVerticle; router.get("/api/prizes").respond(this); } @Override public Uni<Void> apply(RoutingContext ctx) { String latestPrize = prizeBrokerVerticle.getLatestPrize(); if (latestPrize != null) { return ctx.response() .setStatusCode(200) .putHeader("content-type", "application/json") .end("\"" + latestPrize + "\""); // 包装为JSON格式 } else { return ctx.response() .setStatusCode(204) // 无内容状态码 .end(); } } }
方案2:改用请求-响应模式(适配按需生成场景)
如果希望每次HTTP请求触发一次奖品生成,而非定时广播,调整如下:
- 修改
PrizeBrokerVerticle为请求处理器:
@Component public class PrizeBrokerVerticle extends AbstractVerticle { // 原有代码保留 public void start() { // 注册请求处理器,处理获取奖品的请求 eventBus.consumer("prize-request", message -> { String randomPrize = generateRandomPrize(); message.reply(randomPrize); // 回复请求 LOGGER.info("Prize generated and replied: " + randomPrize); }); } // 原有代码保留 }
- 修改
PrizeController,发送请求到指定地址:
@Controller public class PrizeController implements Function<RoutingContext, Uni<Void>> { private final EventBus eventBus; public PrizeController(Router router, EventBus eventBus) { this.eventBus = eventBus; router.get("/api/prizes").respond(this); } @Override public Uni<Void> apply(RoutingContext ctx) { return eventBus.request("prize-request", "") .flatMap(reply -> { String prize = (String) reply.body(); return ctx.response() .setStatusCode(200) .putHeader("content-type", "application/json") .end("\"" + prize + "\""); }) .onFailure(err -> { ctx.response() .setStatusCode(500) .end("Failed to get prize"); return Uni.createFrom().voidItem(); }); } }
关键说明
- 定时广播场景优先用缓存方案,避免EventBus语义混淆。
- 请求-响应模式适合按需生成资源的场景,需单独指定请求地址,不要和广播地址混用。
内容的提问来源于stack exchange,提问作者Plus Ultra
相关产品推荐
相关产品推荐

