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

执行消息迁移至Elasticsearch时遇Closed Resultset : next错误求助

问题描述
  • 执行Oracle消息迁移至Elasticsearch操作,运行一段时间后触发Closed ResultSet错误
  • 已确认Oracle数据库端直接执行查询无异常
  • 触发错误场景:调用fetchMessageIds方法获取50条ID时,无法读取查询结果

相关代码

迁移主方法

public static void migrateMessagesToElastic(Long firstMsgId, Long lastMsgId, Integer size) {
    int offset = 0;
    size = size != null ? size : 100;
    List<Long> messageIdList;
    List<Message> messageList = new ArrayList<>();
    do {
        try (HibernateSession session = SessionManager.openSession()) {
            messageIdList = fetchMessageIds(firstMsgId, lastMsgId, offset, size, session);
            messageList = messageIdList.stream().map(e -> MessageCRUD.loadMessage(e, session)).collect(Collectors.toList());
            bulkAddToElastic(messageList, session);
            offset += size;
            if (messageList.size() >= 1) {
                CacheUtil.putLastMigratedMessageId(messageList.get(messageList.size() - 1).getId());
            }
            System.out.println("number of migrated messages : " + offset);
        } catch (Exception e) {
            e.printStackTrace();
        }
    }
    while (size <= messageList.size());
}

获取消息ID方法

private static List<Long> fetchMessageIds(Long firstMsgId, Long lastMsgId, int offset, int size, HibernateSession session) throws Exception {
    try {
        String hql = "select m.id from Message m left outer join m.excludedUsers where m.id >= :firstMsgId";
        if (lastMsgId != null) {
            hql += "  and m.id <= :lastMsgId ";
        }
        hql += " order by m.id";
        Query query = session.createQuery(hql)
                .setParameter("firstMsgId", firstMsgId);
        if (lastMsgId != null) {
            query.setParameter("lastMsgId", lastMsgId);
        }
        List<Long> ids = (List<Long>) query
                .setFirstResult(offset)
                .setMaxResults(size)
                .list();
        return ids;
    } catch (Exception ex) {
        ex.printStackTrace();
        throw ex;
    }
}

错误信息

java.sql.SQLRecoverableException: Closed Resultset: next
    at oracle.jdbc.driver.OracleResultSetImpl.next(OracleResultSetImpl.java:277)
    ...(省略后续栈帧)

问题分析与修复方案

核心原因

  1. 结果集膨胀:left outer join m.excludedUsers会让同一条Message生成多条重复记录(对应多个关联的excludedUser),setMaxResults(size)限制的是膨胀后的记录数,而非实际唯一Message ID的数量,导致分页逻辑混乱,触发驱动层面的结果集关闭异常。
  2. 分页逻辑错误:循环条件while (size <= messageList.size())判断逻辑错误,当实际返回的有效数据量小于size时,仍会继续循环,引发空查询或会话资源异常。
  3. 会话生命周期与延迟加载冲突:流处理stream().map(...)可能延迟读取数据,若会话提前关闭,会导致尝试访问已关闭的结果集。

修复步骤

1. 修正HQL,避免结果集重复

添加distinct去重,确保每个Message ID只返回一次:

// 修改后的HQL
String hql = "select distinct m.id from Message m left outer join m.excludedUsers where m.id >= :firstMsgId";

如果关联excludedUsers没有实际过滤需求,直接简化查询(左连接不影响主表数据):

String hql = "select m.id from Message m where m.id >= :firstMsgId";

2. 调整分页循环逻辑

基于实际返回的ID数量判断是否继续循环,同时用实际返回数累加偏移量:

do {
    try (HibernateSession session = SessionManager.openSession()) {
        messageIdList = fetchMessageIds(firstMsgId, lastMsgId, offset, size, session);
        if (messageIdList.isEmpty()) {
            break;
        }
        // 提前转为普通List,避免延迟加载导致会话关闭后无法读取
        List<Long> ids = new ArrayList<>(messageIdList);
        messageList = ids.stream().map(e -> MessageCRUD.loadMessage(e, session)).collect(Collectors.toList());
        
        bulkAddToElastic(messageList, session);
        offset += ids.size();
        CacheUtil.putLastMigratedMessageId(ids.get(ids.size() - 1));
        System.out.println("number of migrated messages : " + offset);
    } catch (Exception e) {
        e.printStackTrace();
        // 可选:添加单次失败重试逻辑,避免中断整个迁移
    }
} while (messageIdList.size() == size); // 仅当返回满size条数据时继续循环

3. 检查并升级Oracle JDBC驱动

确保驱动版本与Oracle数据库版本兼容,老版本驱动在游标管理上存在bug,建议升级至ojdbc8或对应数据库版本的官方驱动。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.29 08:33:25