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

Quarkus:基于Kafka写入确认(Ack)与拒绝(Nack)返回对应HTTP状态码

解决Kafka推送接口提前返回状态码的问题

嘿,我明白你的困扰了!你现在的问题核心在于,Emitter.send是异步执行的,你的方法直接返回响应,根本没等Kafka的ack或nack回调执行完毕。而sendAndAwait只是阻塞到消息“发送出去”,但不会等待Kafka broker的确认,所以确实没法判断最终是否写入成功。

下面给你两种可行的解决方案,根据你用的Emitter类型来选:

方案一:使用CompletableFuture包装普通Emitter的回调

如果你的priceEmitter是普通的Emitter(不是Mutiny版本),可以用CompletableFuture来捕获ack和nack的结果,再转换成Quarkus支持的Uni<Response>来实现异步响应:

@Path("/prices")
public class PriceResource {
    @Inject
    @Channel("price-create")
    Emitter<Double> priceEmitter;

    @POST
    @Consumes(MediaType.TEXT_PLAIN)
    public Uni<Response> addPrice(Double price) {
        // 创建一个Future来承载最终的响应结果
        CompletableFuture<Response> responseFuture = new CompletableFuture<>();

        priceEmitter.send(Message.of(price)
                .withAck(() -> {
                    // Kafka确认消息写入成功,返回200
                    responseFuture.complete(Response.ok().build());
                    return CompletableFuture.completedFuture(null);
                })
                .withNack(throwable -> {
                    // Kafka拒绝消息,抛出异常让后续处理返回5xx
                    responseFuture.completeExceptionally(throwable);
                    return CompletableFuture.completedFuture(null);
                }));

        // 将Future转为Uni,处理异常返回500
        return Uni.createFrom().future(responseFuture)
                .onFailure().recoverWithItem(throwable -> 
                    Response.status(Response.Status.INTERNAL_SERVER_ERROR)
                            .entity("写入Kafka失败:" + throwable.getMessage())
                            .build()
                );
    }
}

方案二:用MutinyEmitter简化代码(推荐)

如果你的项目里用的是Mutiny风格的MutinyEmitter(大多数Quarkus新项目默认都是这个),那代码可以更简洁——它的send方法直接返回一个Uni,这个Uni会在Kafka确认消息时完成,在消息被拒绝时失败:

@Path("/prices")
public class PriceResource {
    // 注意这里是MutinyEmitter
    @Inject
    @Channel("price-create")
    MutinyEmitter<Double> priceEmitter;

    @POST
    @Consumes(MediaType.TEXT_PLAIN)
    public Uni<Response> addPrice(Double price) {
        return priceEmitter.send(price)
                // 消息确认成功,返回200
                .onItem().transform(ignored -> Response.ok().build())
                // 消息写入失败,返回500并携带错误信息
                .onFailure().recoverWithItem(throwable -> 
                    Response.status(Response.Status.INTERNAL_SERVER_ERROR)
                            .entity("Kafka写入失败:" + throwable.getMessage())
                            .build()
                );
    }
}

关键要点解释:

  • 把接口返回类型改成Uni<Response>:Quarkus的JAX-RS支持异步响应,会自动等待Uni完成后再给客户端返回状态码,不会提前返回。
  • 绑定Kafka的确认/拒绝事件:无论是用回调还是Mutiny的Uni,核心都是要等到Kafka broker给出明确的ack/nack信号,再生成对应的HTTP响应。
  • 异常处理:一定要捕获nack的异常,转成5xx状态码返回,这样客户端才能明确知道写入失败了。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.04.30 15:04:07