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

微服务获取百万级数据致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框架。它原生支持分片处理、批量读写、内存自动管理,无需手动实现分批逻辑。

核心配置思路:

  1. 配置JdbcCursorItemReader:流式读取源库数据,通过fetchSize控制每次从数据库拉取的行数
  2. 配置JdbcBatchItemWriter:批量写入目标库
  3. 设置Step的chunk size为你的批次大小,框架自动处理分批、事务、内存回收

额外优化建议

  • 避免SELECT *,只查询需要归档的字段,减少单条记录的内存占用
  • 归档操作放在业务低峰期执行,降低对线上服务的影响
  • 调整JVM参数(如-XX:+UseG1GC)优化垃圾回收效率,配合分批处理进一步降低内存压力

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.14 18:50:27