多线程异步任务触发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
相关产品推荐
相关产品推荐

