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

Quarkus多线程运行时事务错误排查求助

问题分析与修复方案

问题根源

  1. POST接口的copyOne方法带有@Transactional注解,且在请求线程内执行,CDI容器会自动激活事务上下文,因此数据库操作正常。
  2. 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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.06 10:10:56