如何基于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>类型,但编译器仍提示必须返回该类型。
期望执行流程
- 先调用NamedQuery查询对应记录列表
- 比较列表长度与允许的容量值
- 有可用座位则执行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
相关产品推荐
相关产品推荐

