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

JDBC ResultSet新增数据被误标记为已导出的原因排查

问题描述

我通过变更指示器从AS/400数据库的多张表查询数据,导出至文件后更新变更指示器标记数据为已导出。但流程执行期间,若用户新增数据,这些新数据会被意外加入ResultSet并标记为已导出。调试发现,在markResultSetAsExported方法中调用extractedData.next()时,ResultSet的内容会因表中新增数据而增多。

相关代码片段

创建Statement与ResultSet

stmt = conn.createStatement(ResultSet.TYPE_SCROLL_INSENSITIVE, ResultSet.CONCUR_UPDATABLE);
extractedData = stmt.executeQuery(sql);

数据导出核心方法

private void writeDataToFile(TableIdentifier table, String tableHeader, ResultSet extractedData,
        TableMetaData tableMetaData, Date currentTimestamp) {
    BufferedFileWriter bfw = null;
    boolean shouldUpdate = true;
    try {
        logger.info(Message.COLLECTOR_EXPORT_START, table.getName());
        ResultSetMetaData rsMetaData = extractedData.getMetaData();
        int numberOfColumns = rsMetaData.getColumnCount();
        for (ExportConfig config : tableQueueMap.get(table)) {
            expDir = WorkDir.getWorkDir(context, applicationDomain + File.separator + EXPDIR + File.separator
                    + (config.isSendComplete() ? "complete_" : "") + config.getSendQueue());
            String[] qualTable = table.getName().split("\\.");
            String tableName = table.getName();
            if (qualTable.length > 1) {
                tableName = qualTable[1];
            }
            File exportFile = new File(expDir, tableName);

            if (!exportFile.exists()) {
                bfw = new BufferedFileWriter(exportFile);
                writeHeader(tableHeader, bfw);
            } else {
                bfw = new BufferedFileWriter(exportFile, true);
            }
            writeLines(tableName, extractedData, bfw, numberOfColumns);
            setTimestamp(table.getName(), currentTimestamp, config);
            bfw.close();
            if (config.isSendComplete()) {
                shouldUpdate = false;
            }
        }
        if (tableMetaData.getUpdate() != null && shouldUpdate) {
            extractedData.beforeFirst();
            if (tableMetaData.getVersion() != null) {
                markDataAsExported(extractedData, tableMetaData);
            } else {
                markResultSetAsExported(extractedData, tableMetaData);
            }
        }
        tableQueueMap.remove(table);
        dataExported = true;
    } catch (IOException ex0) {
        logger.error(Message.COLLECTOR_EXPORT_FILE_ERROR, table.getName(), ex0.getMessage());
    } catch (SQLException ex1) {
        logger.error(Message.COLLECTOR_RESULTSET_ERROR, table.getName(), ex1.getMessage());
    } catch (WorkDirException ex2) {
        logger.error(Message.COLLECTOR_WORKDIR_ERROR, ex2.getMessage());
    } finally {
        if (bfw != null)
            try {
                bfw.close();
            } catch (IOException e) {
                // ignore
            }
    }
}

标记数据为已导出的方法

private void markResultSetAsExported(ResultSet extractedData, TableMetaData tableMetaData) throws SQLException {
    while (extractedData.next()) {
        if (tableMetaData.getUpdateColumnType().equals("N")) {
            extractedData.updateInt(tableMetaData.getUpdateColumn(),
                    Integer.parseInt(tableMetaData.getUpdateValue()));
        } else {
            extractedData.updateString(tableMetaData.getUpdateColumn(), tableMetaData.getUpdateValue());
        }
        extractedData.updateRow();
    }
}

问题原因

虽然使用了ResultSet.TYPE_SCROLL_INSENSITIVE,但JDBC规范中该类型仅保证对数据库的更新/删除操作不敏感,并不保证对新增数据完全隔离。部分数据库驱动(包括AS/400的JDBC驱动)实现的ResultSet可能是动态游标,会在每次调用next()时从数据库拉取最新数据,导致导出过程中新增的数据被加入结果集。

解决方案

方案1:将ResultSet一次性加载到内存

查询完成后,立即将所有数据读取到内存集合中,后续导出和标记操作都基于内存数据,彻底隔离数据库变更:

// 执行查询后,先把ResultSet数据转存到内存列表
List<Map<String, Object>> dataSnapshot = new ArrayList<>();
ResultSetMetaData metaData = extractedData.getMetaData();
int columnCount = metaData.getColumnCount();
while (extractedData.next()) {
    Map<String, Object> row = new HashMap<>();
    for (int i = 1; i <= columnCount; i++) {
        row.put(metaData.getColumnName(i), extractedData.getObject(i));
    }
    dataSnapshot.add(row);
}
// 后续导出时遍历dataSnapshot,而不是直接操作ResultSet
// 标记已导出时,可根据主键批量更新数据库,或者单独查询每条记录更新

方案2:强制驱动一次性获取所有结果

通过设置Statement的fetchSize为Integer.MIN_VALUE(部分驱动支持该参数强制一次性拉取所有数据):

stmt = conn.createStatement(ResultSet.TYPE_SCROLL_INSENSITIVE, ResultSet.CONCUR_UPDATABLE);
stmt.setFetchSize(Integer.MIN_VALUE); // 强制一次性加载所有数据到客户端内存
extractedData = stmt.executeQuery(sql);

方案3:查询时锁定时间范围

在查询SQL中添加时间条件,仅查询执行查询时刻之前的未导出数据,确保结果集是固定快照:

SELECT * FROM YOUR_TABLE 
WHERE CHANGE_INDICATOR = '未导出' 
AND CREATE_TIMESTAMP <= CURRENT_TIMESTAMP -- 或者查询开始前记录的时间戳

内容的提问来源于stack exchange,提问作者Moh-Aw

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.17 07:20:09