Quarkus多线程运行时事务错误排查求助
问题分析与修复方案
问题根源
- POST接口的
copyOne方法带有@Transactional注解,且在请求线程内执行,CDI容器会自动激活事务上下文,因此数据库操作正常。 - GET接口的
streamHPO方法手动创建了新线程,该线程不继承原请求的CDI上下文:streamHPOThread中调用ComPocTtEntity.listAll()时,没有活跃的事务或CDI请求上下文,触发错误。- 后续手动创建线程调用
copyOne时,@Transactional注解无法生效——CDI事务上下文不会绑定到自定义线程上。
修复方案
方案1:使用Quarkus ManagedExecutor(推荐)
Quarkus提供的ManagedExecutor会自动传播CDI上下文,确保注解生效,同时避免手动线程管理的问题:
import jakarta.enterprise.concurrent.ManagedExecutor; import jakarta.inject.Inject; import jakarta.transaction.Transactional; import jakarta.ws.rs.GET; import jakarta.ws.rs.POST; import jakarta.ws.rs.Path; import jakarta.ws.rs.PathParam; import jakarta.ws.rs.core.Response; import java.io.IOException; import java.util.ArrayList; import java.util.List; import java.util.concurrent.Future; @Path("/poc") public class PocResource { @Inject ManagedExecutor managedExecutor; @GET @Path("/streamHPO/{maxt}") public Response streamHPO(@PathParam("maxt") int maxt) throws IOException { // 用ManagedExecutor提交任务,自动传播上下文 managedExecutor.submit(() -> { try { streamHPOThread(maxt); } catch (IOException e) { throw new RuntimeException(e); } }); return Response.ok().build(); } @POST @Transactional @Path("/copyOne") public void copyOne(long id) { ComPocTtEntity source = ComPocTtEntity.findById(id); System.out.println("Id: " + source.getId()); HighPriorityOpenEntity target = new HighPriorityOpenEntity(); target.setId(source.getId()); target.setName(source.getName()); target.setDescription(source.getDescription()); target.setType(source.getType()); target.setOwner(source.getOwner()); target.setCreated(source.getCreated()); target.setUpdated(source.getUpdated()); target.setDue(source.getDue()); target.setResolved(source.getResolved()); target.setClosed(source.getClosed()); target.setSeverity(source.getSeverity()); target.persistAndFlush(); } @Transactional // 添加事务注解,确保listAll有上下文支持 public void streamHPOThread(int maxt) throws IOException { List<ComPocTtEntity> allEntities = ComPocTtEntity.listAll(); List<Long> hpoIds = allEntities.stream() .filter(e -> e.getPriority().equals("HIGH")) .filter(e -> e.getStatus().equals("OPEN")) .map(ComPocTtEntity::getId) .toList(); System.out.println("HPO size: " + hpoIds.size()); int total = hpoIds.size(); List<Future<?>> futures = new ArrayList<>(); for (int i = 0; i < total; i++) { // 控制并发数,等待空闲线程 while (futures.stream().filter(f -> !f.isDone()).count() >= maxt) { try { Thread.sleep(100); } catch (InterruptedException e) { Thread.currentThread().interrupt(); } } long id = hpoIds.get(i); futures.add(managedExecutor.submit(() -> copyOne(id))); System.out.println("HPO: " + i + " of " + total); } // 可选:等待所有任务完成 for (Future<?> future : futures) { try { future.get(); } catch (Exception e) { throw new RuntimeException(e); } } } }
方案2:手动激活CDI请求上下文(适合必须用自定义线程的场景)
如果必须使用手动创建的线程,需要手动激活和关闭CDI请求上下文:
import jakarta.enterprise.context.RequestContext; import jakarta.inject.Inject; import jakarta.transaction.Transactional; import jakarta.ws.rs.GET; import jakarta.ws.rs.POST; import jakarta.ws.rs.Path; import jakarta.ws.rs.PathParam; import jakarta.ws.rs.core.Response; import java.io.IOException; import java.util.ArrayList; import java.util.List; @Path("/poc") public class PocResource { @Inject RequestContext requestContext; @GET @Path("/streamHPO/{maxt}") public Response streamHPO(@PathParam("maxt") int maxt) throws IOException { new Thread(() -> { requestContext.activate(); // 激活请求上下文 try { streamHPOThread(maxt); } catch (IOException e) { throw new RuntimeException(e); } finally { requestContext.deactivate(); // 务必关闭上下文 } }).start(); return Response.ok().build(); } @POST @Transactional @Path("/copyOne") public void copyOne(long id) { // 原代码不变 ComPocTtEntity source = ComPocTtEntity.findById(id); System.out.println("Id: " + source.getId()); HighPriorityOpenEntity target = new HighPriorityOpenEntity(); // 赋值逻辑... target.persistAndFlush(); } @Transactional public void streamHPOThread(int maxt) throws IOException { List<ComPocTtEntity> allEntities = ComPocTtEntity.listAll(); List<Long> hpoIds = allEntities.stream() .filter(e -> e.getPriority().equals("HIGH")) .filter(e -> e.getStatus().equals("OPEN")) .map(ComPocTtEntity::getId) .toList(); System.out.println("HPO size: " + hpoIds.size()); int total = hpoIds.size(); List<Thread> threads = new ArrayList<>(); for (int i = 0; i < total; i++) { while (threads.size() >= maxt) { for (int j = 0; j < threads.size(); j++) { if (!threads.get(j).isAlive()) { threads.remove(j); } } } long id = hpoIds.get(i); threads.add(new Thread(() -> { requestContext.activate(); try { copyOne(id); } finally { requestContext.deactivate(); } })); threads.get(i).start(); System.out.println("HPO: " + i + " of " + total); } } }
关键注意点
- 优先使用Quarkus提供的
ManagedExecutor处理异步任务,它会自动管理CDI上下文,避免手动线程的陷阱。 @Transactional注解仅在CDI上下文活跃时生效,自定义线程默认不携带该上下文,需手动处理或使用框架工具。- MSSQL服务器与该问题无关,错误本质是CDI/事务上下文缺失。
内容的提问来源于stack exchange,提问作者szg12345
相关产品推荐
相关产品推荐

