存储过程多结果集传ExecutorService报Closed Resultset错误,如何并行处理?
问题根因
你遇到的Closed Resultset报错核心原因是Groovy的Sql.call方法会在闭包执行完成后自动关闭所有关联的ResultSet、数据库连接资源。你在闭包内仅向线程池提交了异步任务,没有等待任务执行完成,闭包就执行结束,资源被提前回收,等线程池中的任务真正开始读取ResultSet时,资源已经被关闭,就触发了该异常。
另外你代码里的resultSetCount是多线程共享变量,直接做自增操作存在线程安全问题,计数会不准确。
可行方案
方案1:提前读取结果集数据到内存(推荐,实现简单且稳定)
适合结果集数据量不大,不会占用过多内存的场景,你可以在Sql.call的闭包内先把所有ResultSet的数据读取出来,转成Map或者自定义POJO的集合,再把这些内存中的集合提交给线程池处理,完全不依赖数据库连接资源。
示例代码:
// 先定义存储结果的列表 def resultSetsData = [] sql.call("{call MY_PROCEDURE(?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?)}", [caseNumber, version, followUpNumber, caseType, Sql.resultSet(OracleTypes.CURSOR), Sql.resultSet(OracleTypes.CURSOR), // 省略其他参数 Sql.NUMERIC]) { po_case_info, po_prod_info, po_event_info, /* 省略其他结果集 */ po_version_num -> // 读取第一个结果集到内存 if (po_case_info) { def caseData = [] while (po_case_info.next()) { // 把行数据转成Map存到内存,按需取字段 caseData << [ id: po_case_info.getString("id"), // 其他字段按需添加 ] } resultSetsData << caseData } // 依次读取其他所有结果集到内存 if (po_prod_info) { def prodData = [] while (po_prod_info.next()) { prodData << [/* 存对应字段 */] } resultSetsData << prodData } // 其他结果集同理处理 } // 现在所有数据都在内存里,提交给线程池处理 def atomicCount = new AtomicInteger(0) def futures = resultSetsData.collect { dataList -> executorService.submit({ -> dataList.each { row -> atomicCount.incrementAndGet() someMethod(row) } if (dataList.isEmpty()) { someMethod() } } as Runnable) } // 等待所有任务执行完成 futures.each { it.get() } // 后续处理计数逻辑
方案2:手动管控资源生命周期(适合大数据量结果集)
如果结果集数据量太大,无法全部加载到内存,你需要放弃Groovy Sql的自动资源管理,手动获取连接、调用存储过程,等所有异步任务全部执行完成后再手动关闭ResultSet、Statement、连接。
注意事项:
- 单个ResultSet只能绑定给一个线程处理,JDBC的ResultSet实现不是线程安全的,不要多线程操作同一个ResultSet
- 需要捕获异常确保资源最终被关闭,避免数据库连接泄漏
- 计数必须使用原子类保证线程安全
示例核心逻辑:
def conn = sql.dataSource.connection def cs = null try { cs = conn.prepareCall("{call MY_PROCEDURE(?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?)}") // 设置入参 cs.setString(1, caseNumber) cs.setString(2, version) // 省略其他入参设置 // 注册出参 cs.registerOutParameter(5, OracleTypes.CURSOR) cs.registerOutParameter(6, OracleTypes.CURSOR) // 省略其他出参注册 cs.execute() // 获取所有结果集 def resultSets = [ cs.getObject(5), cs.getObject(6), // 其他结果集依次添加 ] def atomicCount = new AtomicInteger(0) def futures = resultSets.collect { rs -> executorService.submit({ -> if (rs == null) return try { int count = 0 while (rs.next()) { count++ atomicCount.incrementAndGet() someMethod(rs) } if (count == 0) { someMethod() } } finally { try { rs.close() } catch (Exception ignore) {} } } as Runnable) } // 等待所有任务执行完成 futures.each { it.get() } } finally { // 关闭资源 try { cs?.close() } catch (Exception ignore) {} try { conn.close() } catch (Exception ignore) {} }
内容的提问来源于stack exchange,提问作者Nitesh kumar Rai
相关产品推荐
相关产品推荐

