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

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上正在发布的随机奖品。


解决方案

问题根源

  1. EventBus语义混淆:publish是广播模式,无回复机制;request是请求-响应模式,需要接收方主动回复,但你的发布者仅做广播,不会处理request的回复请求。
  2. 重复注册消费者:每次HTTP请求都会新增一个prizes消费者,逻辑混乱且浪费资源。
  3. 请求地址错误:request发送到广播地址prizes,该地址的消费者并非处理请求的角色,无法正确回复。

修正方案

方案1:缓存最新奖品(适配定时广播场景)

这是最贴合你需求的方案,直接缓存定时发布的最新奖品,HTTP请求直接取缓存值:

  1. 修改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;
    }
    // 原有代码保留
}
  1. 修改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请求触发一次奖品生成,而非定时广播,调整如下:

  1. 修改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);
        });
    }
    // 原有代码保留
}
  1. 修改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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.03 19:14:53