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
相关产品推荐
相关产品推荐

