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

使用CompletableFuture、Executor和Semaphore时丢失JDBC连接

并发限流导致EntityManagerFactory关闭异常的排查与修复

场景说明

原本的代码可以并行处理实体,但无法限制最大并发数为2,因此引入Semaphore做限流,结果处理完前两个实体后抛出EntityManagerFactory is closed异常。

原不限流代码

CompletableFuture<Void> taskFuture = idList.stream()
    .map(id -> (Runnable) () -> processEntity(startTime, id))
    .map(task -> CompletableFuture.runAsync(task, executorService))
    .collect(Collectors.collectingAndThen(Collectors.toList(),
        tasks -> CompletableFuture.allOf(tasks.toArray(new CompletableFuture[0]))));

 taskFuture.get();

引入Semaphore后的限流代码

var semaphore = new StreamSemaphore(MAX_CONCURRENT_EXECUTIONS);
var futureTasks = itemList.stream()
    .map(semaphore::acquire)
    .map(itemId -> (Supplier<Void>) () -> { 
    var result = processItem(startTime, itemId); 
    semaphore.release(null); 
    return result;
})
    .map(task -> CompletableFuture.supplyAsync(task, taskExecutor)
        .exceptionally(e -> {
            log.error("Exception occurred while processing item", e);
            System.exit(1);
            return null;
        }))
    .toList();
futureTasks.forEach(TaskProcessor::safeGet);

StreamSemaphore实现

public class StreamSemaphore {
    private final Semaphore sema;

    public StreamSemaphore(int slots) {
        sema = new Semaphore(slots);
    }

    @SneakyThrows
    public <T> T acquire(T obj) {
        sema.acquire();
        return obj;
    }

    @SneakyThrows
    public <T> T release(T obj) {
        sema.release();
        return obj;
    }
}

异常堆栈

org.springframework.transaction.CannotCreateTransactionException: Could not open JPA EntityManager for transaction; nested exception is java.lang.IllegalStateException: EntityManagerFactory is closed
    at org.springframework.orm.jpa.JpaTransactionManager.doBegin(JpaTransactionManager.java:467)
    at org.springframework.transaction.support.AbstractPlatformTransactionManager.startTransaction(AbstractPlatformTransactionManager.java:400)
    at org.springframework.transaction.support.AbstractPlatformTransactionManager.getTransaction(AbstractPlatformTransactionManager.java:373)
    at org.springframework.transaction.interceptor.TransactionAspectSupport.createTransactionIfNecessary(TransactionAspectSupport.java:595)
    at org.springframework.transaction.interceptor.TransactionAspectSupport.invokeWithinTransaction(TransactionAspectSupport.java:382)
    at org.springframework.transaction.interceptor.TransactionInterceptor.invoke(TransactionInterceptor.java:119)
    at org.springframework.aop.framework.ReflectiveMethodInvocation.proceed(ReflectiveMethodInvocation.java:186)
    at org.springframework.aop.framework.CglibAopProxy$CglibMethodInvocation.proceed(CglibAopProxy.java:763)
    at org.springframework.aop.framework.CglibAopProxy$DynamicAdvisedInterceptor.intercept(CglibAopProxy.java:708)
    at com.xegara.process.PlanAccountRelationshipMaintenanceService$$EnhancerBySpringCGLIB$$dea4561a.handleInScopeCustomer(<generated>)
    at com.xegara.process.PlanAccountMaintenanceProcess.handleCustomerList(PlanAccountMaintenanceProcess.java:271)
    at com.xegara.process.PlanAccountMaintenanceProcess.handleCustomerListTransaction(PlanAccountMaintenanceProcess.java:247)
    at com.xegara.process.PlanAccountMaintenanceProcess.handlePlanAccountRelationshipForOwners(PlanAccountMaintenanceProcess.java:208)
    at com.xegara.process.PlanAccountMaintenanceProcess.handlePlanAccountRelationshipForDepartment(PlanAccountMaintenanceProcess.java:156)
    at com.xegara.process.PlanAccountMaintenanceProcess.handleDepartment(PlanAccountMaintenanceProcess.java:119)
    at com.xegara.process.PlanAccountMaintenanceProcess.lambda$process$0(PlanAccountMaintenanceProcess.java:92)
    at java.base/java.util.concurrent.CompletableFuture$AsyncSupply.run$$$capture(CompletableFuture.java:1768)
    at java.base/java.util.concurrent.CompletableFuture$AsyncSupply.run(CompletableFuture.java)
    at com.bytedance.cg.gcrm.overseas.common.util.executor.ContextPassThreadPoolTaskExecutor.lambda$null$1(ContextPassThreadPoolTaskExecutor.java:73)
    at java.base/java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1136)
    at java.base/java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:635)
    at java.base/java.lang.Thread.run(Thread.java:833)
Caused by: java.lang.IllegalStateException: EntityManagerFactory is closed
    at org.hibernate.internal.SessionFactoryImpl.validateNotClosed(SessionFactoryImpl.java:547)
    at org.hibernate.internal.SessionFactoryImpl.createEntityManager(SessionFactoryImpl.java:636)
    at org.hibernate.internal.SessionFactoryImpl.createEntityManager(SessionFactoryImpl.java:158)
    at org.springframework.orm.jpa.AbstractEntityManagerFactoryBean.createNativeEntityManager(AbstractEntityManagerFactoryBean.java:585)
    at jdk.internal.reflect.GeneratedMethodAccessor228.invoke(Unknown Source)
    at java.base/jdk.internal.reflect.DelegatingMethodAccessorImpl.invoke(DelegatingMethodAccessorImpl.java:43)
    at java.base/java.lang.reflect.Method.invoke(Method.java:568)
    at org.springframework.orm.jpa.AbstractEntityManagerFactoryBean.invokeProxyMethod(AbstractEntityManagerFactoryBean.java:487)
    at org.springframework.orm.jpa.AbstractEntityManagerFactoryBean$ManagedEntityManagerFactoryInvocationHandler.invoke(AbstractEntityManagerFactoryBean.java:734)
    at jdk.proxy2/jdk.proxy2.$Proxy191.createNativeEntityManager(Unknown Source)
    at org.springframework.orm.jpa.JpaTransactionManager.createEntityManagerForTransaction(JpaTransactionManager.java:485)
    at org.springframework.orm.jpa.JpaTransactionManager.doBegin(JpaTransactionManager.java:410)
    ... 21 common frames omitted

问题原因

  1. Semaphore获取时机错误:在Stream的map(semaphore::acquire)阶段就获取许可,主线程遍历所有item时会直接占满所有许可(这里是2个),主线程后续如果触发Spring容器关闭(比如当前方法是批量任务的收尾阶段),会直接关闭EntityManagerFactory,而异步任务此时还在执行,自然无法创建EntityManager。
  2. 异常处理不当:exceptionally中调用System.exit(1),一旦有异常直接退出JVM,强制关闭所有资源,导致后续任务失败。

修复方案

1. 调整Semaphore获取时机到异步任务内部

把acquire放到异步执行的Supplier里,确保许可在异步线程中获取,主线程不会提前占用许可导致后续资源被回收:

var semaphore = new Semaphore(MAX_CONCURRENT_EXECUTIONS);
var futureTasks = itemList.stream()
    .map(itemId -> (Supplier<Void>) () -> { 
        semaphore.acquire();
        try {
            return processItem(startTime, itemId); 
        } finally {
            semaphore.release(); // 确保无论成功失败都释放许可,避免泄漏
        }
    })
    .map(task -> CompletableFuture.supplyAsync(task, taskExecutor)
        .exceptionally(e -> {
            log.error("Exception occurred while processing item", e);
            // 移除System.exit(1),避免直接退出导致资源关闭
            return null;
        }))
    .toList();

// 等待所有任务完成后再结束主线程
CompletableFuture.allOf(futureTasks.toArray(new CompletableFuture[0])).join();

2. 移除自定义StreamSemaphore

原生Semaphore已经满足需求,无需额外封装,减少不必要的复杂度。

3. 检查主线程生命周期

确认当前代码所在方法是否在Spring容器关闭前执行,比如是否是@PreDestroy回调或者批量任务的最后一步。如果是,必须确保所有异步任务完成后再让主线程结束,防止容器提前回收EntityManagerFactory等资源。

注意事项

  • 使用Spring事务的异步任务,必须使用支持上下文传递的线程池(比如ContextPassThreadPoolTaskExecutor),否则无法获取正确的事务上下文。
  • 避免在异步任务中直接调用System.exit(1),应通过异常处理机制优雅处理错误,或者标记任务失败后统一处理。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.15 04:24:49