微服务获取百万级数据致JVM堆内存耗尽,求非扩容解决方案
无需扩容内存处理海量数据归档的解决方案
当前代码通过queryForList一次性将百万级记录加载到JVM堆内存,直接触发内存耗尽报错。以下是不用扩容内存就能解决问题的可行方案:
方案一:流式读取+逐批处理(基于JdbcTemplate原生能力)
放弃一次性加载全量数据,改用流式方式逐行读取结果集,每积累指定批次量就写入目标库,同时清理内存中的临时数据,避免内存堆积。
修改核心代码如下:
@Value("${batch-size}") private int batchSize; public void archiveTableRecords(JdbcTemplate sourceDbTemplate, JdbcTemplate targetDbTemplate, ArchiveConfigDTO archiveObj) { try { String sourceTable = archiveObj.getSourceTable(); String archive_months = archiveObj.getArchiveCriteriaMonths(); String compareDate1 = getCSTDateNew(archive_months); logger.info("Archive criteria date: {}", compareDate1); String fetchSql = ArchiveSQLQueries.buildSQLQueryToFetchSourceRecords(sourceTable); // 流式读取结果集,逐批处理 sourceDbTemplate.query(fetchSql, new Object[]{compareDate1}, rs -> { List<Map<String, Object>> batchRecords = new ArrayList<>(batchSize); List<Object> batchPrimaryKeys = new ArrayList<>(batchSize); int count = 0; while (rs.next()) { Map<String, Object> record = new HashMap<>(); ResultSetMetaData metaData = rs.getMetaData(); int columnCount = metaData.getColumnCount(); for (int i = 1; i <= columnCount; i++) { String columnName = metaData.getColumnName(i); record.put(columnName, rs.getObject(i)); if (columnName.equals(archiveObj.getPrimaryKeyColumn())) { batchPrimaryKeys.add(rs.getObject(i)); } } batchRecords.add(record); count++; // 达到批次量执行插入+删除 if (count >= batchSize) { copyBatchRecords(targetDbTemplate, archiveObj.getTargetTable(), archiveObj.getPrimaryKeyColumn(), batchRecords); deleteBatchRecords(sourceDbTemplate, sourceTable, archiveObj.getPrimaryKeyColumn(), batchPrimaryKeys); // 清空批次集合释放内存 batchRecords.clear(); batchPrimaryKeys.clear(); count = 0; } } // 处理剩余不足一批的记录 if (!batchRecords.isEmpty()) { copyBatchRecords(targetDbTemplate, archiveObj.getTargetTable(), archiveObj.getPrimaryKeyColumn(), batchRecords); deleteBatchRecords(sourceDbTemplate, sourceTable, archiveObj.getPrimaryKeyColumn(), batchPrimaryKeys); } }); } catch (Exception e) { logger.error("Exception in archiveTableRecords: {} {}", e.getMessage(), e); } } // 批量插入单批次数据 private int copyBatchRecords(JdbcTemplate targetDbTemplate, String targetTable, String primaryKeyColumn, List<Map<String, Object>> batchRecords) { if (batchRecords.isEmpty()) return 0; int[][] insertResult = targetDbTemplate.batchUpdate( ArchiveSQLQueries.buildSQLTargetRecordInsertionQuery(targetTable, batchRecords.get(0), primaryKeyColumn), batchRecords, batchSize, (ps, argument) -> { int index = 1; for (Map.Entry<String, Object> obj : argument.entrySet()) { if (!obj.getKey().equals(primaryKeyColumn)) { ps.setObject(index++, obj.getValue()); } } }); int result = getSumOfArray(insertResult); logger.info("Inserted {} record(s) in {}", result, targetTable); return result; } // 批量删除单批次数据 private void deleteBatchRecords(JdbcTemplate sourceDbTemplate, String sourceTable, String primaryKeyColumn, List<Object> batchPrimaryKeys) { if (batchPrimaryKeys.isEmpty()) return; String deleteSql = "DELETE FROM " + sourceTable + " WHERE " + primaryKeyColumn + " IN (" + String.join(",", Collections.nCopies(batchPrimaryKeys.size(), "?")) + ")"; sourceDbTemplate.update(deleteSql, batchPrimaryKeys.toArray()); logger.info("Deleted {} record(s) from {}", batchPrimaryKeys.size(), sourceTable); }
方案二:数据库分页分段读取
基于主键或时间字段做分页,每次查询一段范围的数据,处理完后再查询下一段,避免一次性加载全量数据。适合主键自增或时间字段连续的场景。
核心逻辑示例:
public void archiveTableRecords(JdbcTemplate sourceDbTemplate, JdbcTemplate targetDbTemplate, ArchiveConfigDTO archiveObj) { try { String sourceTable = archiveObj.getSourceTable(); String primaryKey = archiveObj.getPrimaryKeyColumn(); String archive_months = archiveObj.getArchiveCriteriaMonths(); String compareDate1 = getCSTDateNew(archive_months); logger.info("Archive criteria date: {}", compareDate1); long lastId = 0; while (true) { // 分页查询:基于主键分段,每次取batchSize条 String fetchSql = String.format("SELECT * FROM %s WHERE update_dts <= ? AND %s > ? ORDER BY %s LIMIT %d", sourceTable, primaryKey, primaryKey, batchSize); List<Map<String, Object>> batchRecords = sourceDbTemplate.queryForList(fetchSql, compareDate1, lastId); if (batchRecords.isEmpty()) break; // 更新下一批的起始主键 lastId = (long) batchRecords.get(batchRecords.size()-1).get(primaryKey); // 插入目标库 List<Object> primaryKeyValueList = new ArrayList<>(); int recordsInserted = copySourceRecords(targetDbTemplate, archiveObj.getTargetTable(), primaryKey, batchRecords, primaryKeyValueList); // 删除源库数据 if (recordsInserted > 0) { deleteSourceRecords(sourceDbTemplate, sourceTable, primaryKey, primaryKeyValueList); } // 手动触发GC(可选,根据内存情况调整) System.gc(); } } catch (Exception e) { logger.error("Exception in archiveTableRecords: {} {}", e.getMessage(), e); } }
方案三:使用Spring Batch框架(适合复杂归档场景)
如果归档逻辑涉及事务保障、错误重试或监控需求,直接用Spring Batch框架。它原生支持分片处理、批量读写、内存自动管理,无需手动实现分批逻辑。
核心配置思路:
- 配置
JdbcCursorItemReader:流式读取源库数据,通过fetchSize控制每次从数据库拉取的行数 - 配置
JdbcBatchItemWriter:批量写入目标库 - 设置
Step的chunk size为你的批次大小,框架自动处理分批、事务、内存回收
额外优化建议
- 避免
SELECT *,只查询需要归档的字段,减少单条记录的内存占用 - 归档操作放在业务低峰期执行,降低对线上服务的影响
- 调整JVM参数(如
-XX:+UseG1GC)优化垃圾回收效率,配合分批处理进一步降低内存压力
内容的提问来源于stack exchange,提问作者QueenOfCode
相关产品推荐
相关产品推荐

