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

Spring Batch复用Hibernate会话时会话被关闭的原因及解决方法

问题描述

我有一个Spring Batch任务,流程是通过Hibernate从数据库检索一批数据(任务仅执行一次),再用这批数据通过Hibernate向库中插入其他数据。

为避免org.springframework.orm.hibernate4.HibernateSystemException: Illegal attempt to associate a collection with two open sessions;异常,我尝试在Spring Batch任务里,通过CustomJobExecutionListener把Session对象传递给Writer类,想用同一个会话完成读写操作,但现在会话会被提前关闭,报错:

ERROR spi.SqlExceptionHelper - PooledConnection has already been closed.

如果不调用session.flush(),不会报错,但数据根本写不进数据库。

任务配置

<batch:partition step="customLoadFile" partitioner="customFilePartitioner">
    <batch:handler grid-size="8" task-executor="customJobTaskExecutor"/>
</batch:partition>
</batch:step>

<batch:step id="customLoadFile">
    <batch:tasklet transaction-manager="customTransactionManager">
        <batch:chunk reader="customFileReader" writer="customFileWriter" commit-interval="5"/>
    </batch:tasklet>
</batch:step>

<bean id="customFileWriter" class="com.example.batch.CustomItemWriter">
    <property name="customService" ref="customService"/>
</bean>

<bean id="customJobExecutionListener" class="com.example.batch.CustomJobExecutionListener">
    <constructor-arg ref="customFileWriter"/>
    <constructor-arg ref="customSessionFactory"/>
</bean>

相关类代码

CustomJobExecutionListener.java

public class CustomJobExecutionListener implements JobExecutionListener {

    private final CustomItemWriter customItemWriter;
    private final SessionFactory customSessionFactory;

    public CustomJobExecutionListener(CustomItemWriter customItemWriter, SessionFactory customSessionFactory) {
        this.customItemWriter = customItemWriter;
        this.customSessionFactory = customSessionFactory;
    }

    @Override
    public void beforeJob(JobExecution jobExecution) {
        Session session = customSessionFactory.openSession();
        customItemWriter.populateDataFromDb(session);
    }

    @Override
    public void afterJob(JobExecution jobExecution) {}
}

CustomItemWriter.java

public class CustomItemWriter implements ItemWriter<CustomLogModel> {

    private CustomService customService;
    private Session session;

    private void populateDataFromDb(Session session) {
        this.session = session;

        StopWatch stopWatch = new StopWatch();
        stopWatch.start();
        System.out.println("Retrieving all db info.");

        List<Object[]> data = customService.getAllData(session);

        data.forEach(record -> {
            ArticleDomain articleSupplement = (ArticleDomain) record[2];
            // 处理articleSupplement的逻辑
        });

        stopWatch.stop();
        System.out.println("Finished retrieving data.");
    }

    @Override
    public void write(List<? extends CustomLogModel> items) throws Exception {
        StopWatch stopWatch = new StopWatch();
        stopWatch.start();
        for (CustomLogModel item : items) {
            customService.insertData(item, findArticleDomainByFileName(fileName), session);
        }
        stopWatch.stop();
        System.out.println("Finished inserting data for a batch of " + items.size() + " items.");
    }

    private ArticleDomain findArticleDomainByFileName(String fileName) {
        // 实现逻辑
    }

    public void setCustomService(CustomService customService) {
        this.customService = customService;
    }
}

CustomService.java

public class CustomService {

    @Transactional(readOnly = true)
    public List<Object[]> getAllData(Session session) {
        List<Object[]> dbData = session.createQuery("select i.fileName, i.date, a " + "from ArticleDomain a " + "left join fetch a.ipList i ")
                .list();
        return dbData;
    }

    @Transactional
    public void insertData(CustomLogModel logModel, ArticleDomain articleDomain, Session session) {
        try {
            if (articleDomain.getIpList() != null && !articleDomain.getIpList().isEmpty()) {
                articleDomain.getIpList().getLast().setIsLast(false);
            }

            IpRecord ipRecord = new IpRecord();
            ipRecord.setArticleId(logModel.getArticleId());
            ipRecord.setFileName(logModel.getFileName());
            ipRecord.setFailReason(logModel.getFailReason());
            // 其他字段赋值

            articleDomain.addToIpList(ipRecord);
            session.save(articleDomain);
            session.flush();
        } catch (Exception ex) {
            log.error("Error", ex);
        }
    }
}

说明:ArticleDomain对象包含一个IpRecords列表,同一个fileName可以对应多个不同处理日期的IpRecord。

咨询问题

  1. 是谁关闭了我的会话?
  2. 如何在读写领域对象时复用会话,避免上述异常?
问题解答

1. 谁关闭了会话?

问题出在@Transactional注解和Spring的事务管理机制上:

  • 你在CustomService.getAllData和insertData方法上加了@Transactional,Spring会为每个被@Transactional修饰的方法创建独立的事务上下文。
  • 当你手动传入自己打开的Session时,Spring的事务管理器会认为这个Session不属于当前事务上下文,在方法执行完毕后(事务提交/回滚时),会自动关闭这个Session。
  • 等到后续Writer的write方法再使用这个Session时,它已经被之前的事务关闭了,所以会报PooledConnection has already been closed错误。

另外,你在beforeJob里打开Session后,没有任何代码负责维护它的生命周期,也没有绑定到Spring的事务上下文,这也是会话被提前关闭的原因之一。

2. 如何复用会话避免问题?

方案一:利用Spring Batch的事务上下文,不要手动管理Session

Spring Batch的Chunk模式本身会由transaction-manager管理事务,每个Chunk对应一个事务,会话会被Spring自动维护,无需手动传递Session:

  1. 移除CustomJobExecutionListener中手动打开Session的逻辑,以及Writer中保存Session的字段。
  2. 修改CustomService,不用手动传入Session,而是通过SessionFactory获取当前事务绑定的Session:
public class CustomService {
    private SessionFactory sessionFactory;

    // 注入SessionFactory
    public void setSessionFactory(SessionFactory sessionFactory) {
        this.sessionFactory = sessionFactory;
    }

    @Transactional(readOnly = true)
    public List<Object[]> getAllData() {
        Session session = sessionFactory.getCurrentSession();
        List<Object[]> dbData = session.createQuery("select i.fileName, i.date, a from ArticleDomain a left join fetch a.ipList i ")
                .list();
        return dbData;
    }

    @Transactional
    public void insertData(CustomLogModel logModel, ArticleDomain articleDomain) {
        try {
            if (articleDomain.getIpList() != null && !articleDomain.getIpList().isEmpty()) {
                articleDomain.getIpList().getLast().setIsLast(false);
            }

            IpRecord ipRecord = new IpRecord();
            ipRecord.setArticleId(logModel.getArticleId());
            ipRecord.setFileName(logModel.getFileName());
            ipRecord.setFailReason(logModel.getFailReason());
            // 其他字段赋值

            articleDomain.addToIpList(ipRecord);
            sessionFactory.getCurrentSession().saveOrUpdate(articleDomain);
            // 不需要手动flush,Spring会在事务提交时自动处理
        } catch (Exception ex) {
            log.error("Error", ex);
            throw ex; // 抛出异常让Spring Batch处理回滚
        }
    }
}
  1. 修改CustomItemWriter,在populateDataFromDb中直接调用customService.getAllData(),不用传递Session:
private void populateDataFromDb() {
    StopWatch stopWatch = new StopWatch();
    stopWatch.start();
    System.out.println("Retrieving all db info.");

    List<Object[]> data = customService.getAllData();

    data.forEach(record -> {
        ArticleDomain articleSupplement = (ArticleDomain) record[2];
        // 处理articleSupplement的逻辑,比如缓存起来供write方法使用
    });

    stopWatch.stop();
    System.out.println("Finished retrieving data.");
}
  1. 调整CustomJobExecutionListener,调用populateDataFromDb时不用传Session:
@Override
public void beforeJob(JobExecution jobExecution) {
    customItemWriter.populateDataFromDb();
}

这样所有操作都会使用Spring事务上下文绑定的同一个Session,既不会出现多会话关联集合的异常,也不会有会话被提前关闭的问题。

方案二:手动绑定Session到当前线程(不推荐,除非特殊场景)

如果必须手动管理Session,需要把Session绑定到Spring的TransactionSynchronizationManager,让Spring的事务管理器识别它:

  1. 在CustomJobExecutionListener的beforeJob中打开Session并绑定:
@Override
public void beforeJob(JobExecution jobExecution) {
    Session session = customSessionFactory.openSession();
    TransactionSynchronizationManager.bindResource(customSessionFactory, new SessionHolder(session));
    customItemWriter.populateDataFromDb(session);
}
  1. 在afterJob中解绑并关闭Session:
@Override
public void afterJob(JobExecution jobExecution) {
    SessionHolder sessionHolder = (SessionHolder) TransactionSynchronizationManager.unbindResource(customSessionFactory);
    Session session = sessionHolder.getSession();
    if (session.isOpen()) {
        session.close();
    }
}
  1. 移除CustomService方法上的@Transactional注解(因为我们手动管理事务和Session),所有事务操作自己控制:
public class CustomService {
    public List<Object[]> getAllData(Session session) {
        Transaction tx = null;
        try {
            tx = session.beginTransaction();
            List<Object[]> dbData = session.createQuery("select i.fileName, i.date, a from ArticleDomain a left join fetch a.ipList i ")
                    .list();
            tx.commit();
            return dbData;
        } catch (Exception e) {
            if (tx != null) tx.rollback();
            throw e;
        }
    }

    public void insertData(CustomLogModel logModel, ArticleDomain articleDomain, Session session) {
        Transaction tx = null;
        try {
            tx = session.beginTransaction();
            if (articleDomain.getIpList() != null && !articleDomain.getIpList().isEmpty()) {
                articleDomain.getIpList().getLast().setIsLast(false);
            }

            IpRecord ipRecord = new IpRecord();
            ipRecord.setArticleId(logModel.getArticleId());
            ipRecord.setFileName(logModel.getFileName());
            ipRecord.setFailReason(logModel.getFailReason());
            // 其他字段赋值

            articleDomain.addToIpList(ipRecord);
            session.saveOrUpdate(articleDomain);
            tx.commit();
        } catch (Exception ex) {
            if (tx != null) tx.rollback();
            log.error("Error", ex);
            throw ex;
        }
    }
}

这个方案需要自己处理事务的提交和回滚,容易出错,所以优先推荐方案一。


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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.26 13:42:09