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

如何基于NamedQuery结果,使用Hibernate Reactive实现条件POST操作

问题:基于Mutiny和Hibernate Reactive实现带容量校验的POST接口

需求

实现一个POST操作,仅当指定NamedQuery返回的记录行数小于设定的容量值时,才执行创建操作。

当前遇到的问题

  • 尝试调用NamedQuery的.await()获取布尔值时,报错:This method should exclusively be invoked from a Vert.x EventLoop thread; currently running on thread 'executor-thread-0
  • 使用transform或transformToUni处理逻辑时,明明所有分支都返回了Uni<Response>类型,但编译器仍提示必须返回该类型。

期望执行流程

  1. 先调用NamedQuery查询对应记录列表
  2. 比较列表长度与允许的容量值
  3. 有可用座位则执行POST创建操作,否则返回CONFLICT(409)状态码

当前代码

package org.electricaltrainingalliance.attendees.selections.boundary;

import java.util.UUID;

import javax.enterprise.context.ApplicationScoped;
import javax.inject.Inject;
import javax.ws.rs.Consumes;
import javax.ws.rs.DELETE;
import javax.ws.rs.GET;
import javax.ws.rs.POST;
import javax.ws.rs.PUT;
import javax.ws.rs.Path;
import javax.ws.rs.Produces;
import javax.ws.rs.QueryParam;
import javax.ws.rs.WebApplicationException;
import javax.ws.rs.core.Response;
import javax.ws.rs.core.Response.Status;

import org.electricaltrainingalliance.attendees.selections.entity.AttendeeSelection;
import org.hibernate.reactive.mutiny.Mutiny.SessionFactory;
import org.jboss.logging.Logger;
import org.jboss.logging.Logger.Level;
import org.jboss.resteasy.reactive.RestPath;

import io.smallrye.mutiny.Uni;

@Path("selections")
@ApplicationScoped
@Produces("application/json")
@Consumes("application/json")
public class AttendeeSelectionsService {

    private static final Logger LOGGER = Logger.getLogger(AttendeeSelectionsService.class.getName());

    @Inject
    SessionFactory sf;

    @POST
    public Uni<Response> create(AttendeeSelection selection) {
        if (selection == null || !selection.isValid()) {
            return Uni.createFrom()
                    .item(Response.status(400)
                            .entity("The Selection is not valid, please refer to the API Documentation.")::build);
        }

        Uni<Boolean> seatsAvailable = sf.withSession((s) -> s
                .createNamedQuery(AttendeeSelection.findByOfferingId, AttendeeSelection.class)
                .setParameter("offeringId", selection.getOffering().getOfferingId())
                .getResultList()).onItem()
                .transform(entity -> entity.size() >= selection.getOffering().getCapacity() ? Boolean.FALSE : Boolean.TRUE);

        seatsAvailable.onItem().transformToUni(result -> {
            if (!result) {
                return Uni.createFrom()
                        .item(Response.status(409)
                                .entity("The Offering is no longer available.")::build);
            } else {
                sf.withTransaction((s, t) -> s.persist(selection))
                        .replaceWith(Response.ok(selection).status(Status.CREATED)::build)
                        .onFailure().transform(failure -> {
                            if (failure.getMessage().toLowerCase().contains("uk_unique_selection")) {
                                throw new WebApplicationException(
                                        "A Selection with this information already exists.",
                                        Status.CONFLICT);
                            } else {
                                LOGGER.log(Level.DEBUG, failure.getMessage(), failure);
                                throw new WebApplicationException(failure.getMessage(),
                                        Status.INTERNAL_SERVER_ERROR);
                            }
                        });
            }
        });
        
    }
...
}

解决方案

问题1:EventLoop线程错误

.await()是阻塞式调用,会脱离Vert.x的EventLoop线程,这在响应式编程中是严格禁止的。必须全程使用Mutiny的链式非阻塞调用处理异步操作,不能用阻塞方法强制获取结果。

问题2:返回值类型错误

你的代码中,seatsAvailable.onItem().transformToUni(...)的处理结果没有作为方法返回值返回,导致编译器认为方法缺少Uni<Response>类型的返回值。另外else分支里的代码没有返回对应的Uni实例,需要补全返回语句。

修正后的完整代码

package org.electricaltrainingalliance.attendees.selections.boundary;

import java.util.UUID;

import javax.enterprise.context.ApplicationScoped;
import javax.inject.Inject;
import javax.ws.rs.Consumes;
import javax.ws.rs.DELETE;
import javax.ws.rs.GET;
import javax.ws.rs.POST;
import javax.ws.rs.PUT;
import javax.ws.rs.Path;
import javax.ws.rs.Produces;
import javax.ws.rs.QueryParam;
import javax.ws.rs.WebApplicationException;
import javax.ws.rs.core.Response;
import javax.ws.rs.core.Response.Status;

import org.electricaltrainingalliance.attendees.selections.entity.AttendeeSelection;
import org.hibernate.reactive.mutiny.Mutiny.SessionFactory;
import org.jboss.logging.Logger;
import org.jboss.logging.Logger.Level;
import org.jboss.resteasy.reactive.RestPath;

import io.smallrye.mutiny.Uni;

@Path("selections")
@ApplicationScoped
@Produces("application/json")
@Consumes("application/json")
public class AttendeeSelectionsService {

    private static final Logger LOGGER = Logger.getLogger(AttendeeSelectionsService.class.getName());

    @Inject
    SessionFactory sf;

    @POST
    public Uni<Response> create(AttendeeSelection selection) {
        if (selection == null || !selection.isValid()) {
            return Uni.createFrom()
                    .item(Response.status(400)
                            .entity("The Selection is not valid, please refer to the API Documentation.")::build);
        }

        // 1. 查询当前已选人数,判断是否有可用座位
        return sf.withSession(session -> session
                .createNamedQuery(AttendeeSelection.findByOfferingId, AttendeeSelection.class)
                .setParameter("offeringId", selection.getOffering().getOfferingId())
                .getResultList())
                .onItem().transform(list -> list.size() < selection.getOffering().getCapacity())
                // 2. 根据座位是否可用分支处理
                .onItem().transformToUni(seatsAvailable -> {
                    if (!seatsAvailable) {
                        // 无座位,返回409
                        return Uni.createFrom().item(
                                Response.status(Status.CONFLICT)
                                        .entity("The Offering is no longer available.")
                                        .build());
                    } else {
                        // 有座位,执行持久化操作
                        return sf.withTransaction((session, tx) -> session.persist(selection))
                                .replaceWith(Response.ok(selection).status(Status.CREATED).build())
                                .onFailure().transform(failure -> {
                                    if (failure.getMessage().toLowerCase().contains("uk_unique_selection")) {
                                        return new WebApplicationException(
                                                "A Selection with this information already exists.",
                                                Status.CONFLICT);
                                    } else {
                                        LOGGER.log(Level.DEBUG, failure.getMessage(), failure);
                                        return new WebApplicationException(failure.getMessage(),
                                                Status.INTERNAL_SERVER_ERROR);
                                    }
                                });
                    }
                });
    }

    // 其他方法...
}

关键优化点

  • 全程使用链式非阻塞调用:所有异步操作通过Mutiny的Uni链式处理,彻底摒弃.await()这类阻塞方法
  • 确保方法返回完整异步流:将整个校验+创建的异步流程作为方法返回值返回,避免中途丢弃结果
  • 简化座位判断逻辑:直接用list.size() < capacity生成布尔值,去掉冗余的三元表达式
  • 补全分支返回语句:确保每个逻辑分支都返回合法的Uni<Response>实例

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.22 17:14:56