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

多线程异步任务触发Arjuna事务中止告警及数据未插入问题

问题原因分析

1. 事务跨线程上下文冲突

executeJob方法标记了@Transactional,但被放在CompletableFuture.supplyAsync中执行,默认使用ForkJoinPool.commonPool共享线程池。并发请求时线程池复用线程,导致不同事务上下文混在同一线程中,Arjuna事务管理器检测到同一事务ID下存在多个活跃线程,触发ARJUNA012095警告并中止事务,最终数据回滚。

2. 事务超时被自动中止

远程服务循环拉取数据耗时可能超过事务默认超时时间,Arjuna的Transaction Reaper线程会主动中止超时事务,这也是日志中出现Transaction Reaper Worker 0线程的原因。事务中止后所有数据库操作回滚,因此看不到数据插入。

3. JAX-RS Client资源泄漏

sendProxyPostRequest每次调用都创建新的Client实例,但仅关闭了Response未关闭Client。并发场景下会导致大量未释放的网络资源,引发请求阻塞,进一步加剧事务超时概率。


解决方案

1. 为异步任务使用独立线程池

放弃默认共享线程池,创建专属线程池隔离每个异步任务的事务上下文:

// 定义专属线程池(可配置为Spring/Quarkus Bean)
private static final ExecutorService JOB_EXECUTOR = Executors.newFixedThreadPool(10);

@Override
public CompletableFuture<String> initiateJob(String derivativeType, String minDate, String maxDate) {
    String jobId = jobService.generateJobDerivativeId(derivativeType, minDate, maxDate);
    // 使用自定义线程池执行异步任务
    return CompletableFuture.supplyAsync(() -> {
        executeJob(jobId, derivativeType);
        return jobId;
    }, JOB_EXECUTOR);
}

2. 拆分事务边界,避免远程调用占用事务

将远程拉取数据逻辑移出事务范围,仅在数据库操作时开启事务,减少事务持有时间:

// 去掉executeJob的@Transactional,拆分数据库操作的事务
void executeJob(String jobId, String derivativeType) {
    // 单独事务创建Job
    JobPSBOMEntity job = jobService.createJob(jobId, derivativeType);
    LOG.info("Job is created with id: " + job.getId());
    
    // 远程拉取数据(无事务)
    ResponseDTO responseDTO = remoteService.sendProxyPostRequest(derivativeType);
    
    // 单独开启事务执行数据入库
    dataImportService.saveDataInTransaction(job, responseDTO);
}

// 数据入库的事务方法
@Transactional
public void saveDataInTransaction(JobPSBOMEntity job, ResponseDTO responseDTO) {
    // 后续数据入库逻辑
}

3. 延长事务超时时间

在Quarkus配置文件中调整事务超时,适配远程调用耗时:

# application.properties
quarkus.transaction-manager.default-timeout=300s

4. 正确管理JAX-RS Client资源

复用线程安全的Client实例,或使用try-with-resources确保Client被关闭:

// 修改RemoteServiceClientProvider为单例Client
@Singleton
public class RemoteServiceClientProvider {
    private final Client client;

    public RemoteServiceClientProvider() {
        try {
            TrustManager[] trustAllCerts = new TrustManager[]{
                    new X509TrustManager() {
                        public X509Certificate[] getAcceptedIssuers() {
                            return null;
                        }
                        public void checkClientTrusted(X509Certificate[] certs, String authType) {}
                        public void checkServerTrusted(X509Certificate[] certs, String authType) {}
                    }
            };
            SSLContext sslContext = SSLContext.getInstance("SSL");
            sslContext.init(null, trustAllCerts, new java.security.SecureRandom());
            HostnameVerifier allHostsValid = (hostname, session) -> true;
            this.client = jakarta.ws.rs.client.ClientBuilder.newBuilder()
                    .sslContext(sslContext)
                    .hostnameVerifier(allHostsValid)
                    .build();
        } catch (Exception e) {
            throw new RuntimeException("Failed to create a client with disabled SSL validation", e);
        }
    }

    public Client getClient() {
        return client;
    }

    // 应用关闭时销毁Client
    @PreDestroy
    public void destroy() {
        client.close();
    }
}

// 修改sendProxyPostRequest使用复用的Client
@Override
public ResponseDTO sendProxyPostRequest(String derivativeType) {
    Client client = remoteServiceClientProvider.getClient();
    String authHeader = Base64Encoder.getEncodedBasicAuth(username, password);
    String dataVersion = "";
    String cursor = "";
    ResponseDTO responseDTO;
    int rows = 10000;
    Set<PartUsageDTO> partUsageDTOSet = new HashSet<>();
    Set<MaterialVersionDTO> materialVersionDTOS = new HashSet<>();
    Set<MoboRuleDTO> moboRuleDTOS = new HashSet<>();
    do {
        LOG.info("CURRENT ONE IS: " + Thread.currentThread().getName());
        String requestBody = TestData.getBody(dataVersion, cursor, derivativeType, rows);
        try (Response response = client.target(url)
                .request(MediaType.APPLICATION_JSON)
                .header(HttpHeaders.AUTHORIZATION, authHeader)
                .post(Entity.json(requestBody))) {
            responseDTO = response.readEntity(ResponseDTO.class);
            partUsageDTOSet.addAll(responseDTO.partUsages());
            if (responseDTO.referencedEntities() != null) {
                materialVersionDTOS.addAll(responseDTO.referencedEntities().materialVersions());
            }
            if (responseDTO.referencedEntities() != null && responseDTO.referencedEntities().moboRules() != null && !responseDTO.referencedEntities().moboRules().isEmpty()) {
                moboRuleDTOS.addAll(responseDTO.referencedEntities().moboRules());
            }
            dataVersion = responseDTO.metadata().dataVersion();
            cursor = responseDTO.metadata().nextPageCursor();
        }
    } while (cursor != null && !cursor.isEmpty());
    // 组装返回逻辑...
}

5. 明确事务传播行为

确保jobService.createJob的REQUIRES_NEW传播行为正确生效,避免事务上下文混乱:

@Transactional(propagation = Propagation.REQUIRES_NEW)
public JobPSBOMEntity createJob(String jobId, String derivativeType) {
    // 创建Job的逻辑
}

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.23 02:15:54