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. 谁关闭了会话?
问题出在@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:
- 移除
CustomJobExecutionListener中手动打开Session的逻辑,以及Writer中保存Session的字段。 - 修改
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处理回滚 } } }
- 修改
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."); }
- 调整
CustomJobExecutionListener,调用populateDataFromDb时不用传Session:
@Override public void beforeJob(JobExecution jobExecution) { customItemWriter.populateDataFromDb(); }
这样所有操作都会使用Spring事务上下文绑定的同一个Session,既不会出现多会话关联集合的异常,也不会有会话被提前关闭的问题。
方案二:手动绑定Session到当前线程(不推荐,除非特殊场景)
如果必须手动管理Session,需要把Session绑定到Spring的TransactionSynchronizationManager,让Spring的事务管理器识别它:
- 在
CustomJobExecutionListener的beforeJob中打开Session并绑定:
@Override public void beforeJob(JobExecution jobExecution) { Session session = customSessionFactory.openSession(); TransactionSynchronizationManager.bindResource(customSessionFactory, new SessionHolder(session)); customItemWriter.populateDataFromDb(session); }
- 在
afterJob中解绑并关闭Session:
@Override public void afterJob(JobExecution jobExecution) { SessionHolder sessionHolder = (SessionHolder) TransactionSynchronizationManager.unbindResource(customSessionFactory); Session session = sessionHolder.getSession(); if (session.isOpen()) { session.close(); } }
- 移除
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

