如何提前触发JDBC executeQuery()并返回未消费的JPA结果流?
问题:如何强制提前触发JDBC executeQuery()并返回未消费的JPA结果流?
我使用Hibernate 6对接一个存在故障的只读PostgreSQL数据库,由于需要获取大规模结果集,因此选择调用Query.getResultStream()而非getResultList(),以获取无缓存的仅向前结果流。
数据库故障仅发生在JDBC的executeQuery()调用阶段,处理已打开的ResultSet时从未出现失败。我希望在代码中添加重试机制,但错误仅在将结果流返回给第三方代码后才会出现——因为结果流是延迟执行的,仅当调用终端操作时才会触发JDBC调用,对应的错误栈如下:
org.hibernate.exception.LockAcquisitionException caused by org.postgresql.util.PSQLException ... at org.postgresql.jdbc.PgPreparedStatement.executeQuery ... at java.base/java.util.Iterator.forEachRemaining(Unknown Source) at java.base/java.util.Spliterators$IteratorSpliterator.forEachRemaining(Unknown Source) ... at java.base/java.util.stream.ReferencePipeline.count(Unknown Source) at com.thirdparty.Program.main
我的原始代码如下:
LOGGER.debug("JPQL (" + offset + "," + pagesize + "): " + jpqlQuery); TypedQuery<Map<String, Object>> builder = queryBuilder(jpqlQuery); return builder.setFirstResult(offset).setMaxResults(pagesize).getResultStream();
解决方案
核心思路是提前触发executeQuery()执行,同时保留结果流的仅向前、无缓存特性。可以通过手动操作结果流的迭代器,强制触发JDBC查询,再将迭代器重新包装为流返回,这样就能在当前方法内捕获executeQuery()阶段的错误并执行重试。
修改后的代码示例(带重试逻辑)
LOGGER.debug("JPQL (" + offset + "," + pagesize + "):\n" + jpqlQuery); int retryCount = 3; // 设置重试次数 TypedQuery<Map<String, Object>> query = queryBuilder(jpqlQuery) .setFirstResult(offset) .setMaxResults(pagesize); for (int i = 0; i < retryCount; i++) { try { Stream<Map<String, Object>> resultStream = query.getResultStream(); Iterator<Map<String, Object>> iterator = resultStream.iterator(); // 调用hasNext()强制触发JDBC executeQuery()执行 boolean hasElements = iterator.hasNext(); // 将迭代器重新包装为流,保留仅向前特性 return StreamSupport.stream( Spliterators.spliteratorUnknownSize(iterator, Spliterator.ORDERED), false); } catch (LockAcquisitionException | PSQLException e) { LOGGER.warn("第{}次查询失败,准备重试", i + 1, e); if (i == retryCount - 1) { // 重试次数耗尽,抛出最终异常 throw e; } // 可选:添加重试间隔,避免频繁重试 try { Thread.sleep(1000); } catch (InterruptedException ie) { Thread.currentThread().interrupt(); throw new RuntimeException("重试被中断", ie); } } } // 理论上不会走到这里,返回空流作为兜底 return Stream.empty();
关键逻辑说明
- 提前触发查询:调用
iterator.hasNext()会强制Hibernate执行底层的JDBCexecuteQuery(),并初始化ResultSet。这一步如果出现数据库故障,就能在当前方法内捕获,而不是等到第三方代码消费流时才报错。 - 保留流特性:通过
StreamSupport.stream()将已触发查询的迭代器重新包装为流,依然保持仅向前、无缓存的特性,不会像getResultList()那样把所有数据加载到内存中。 - 重试逻辑:在循环中捕获
executeQuery()阶段的异常(LockAcquisitionException、PSQLException),达到重试次数上限后再抛出最终异常。
内容的提问来源于stack exchange,提问作者basin
相关产品推荐
相关产品推荐

