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

Java多线程跨库批量数据迁移及内存溢出问题解决方案咨询

海量数据跨库同步的内存问题与解决方案咨询

我需要开发一个程序,从数据库查询数据并插入到另一个数据库的历史表中,该历史表每小时新增约4000万条记录。测试时处理100万条数据,采用分页查询(每页10万条)可行,但处理海量数据时出现Java堆内存溢出,调整内存参数后仍未解决。现咨询:

  • 是否可通过线程实现?具体方案是什么?
  • 还有哪些替代方案?

以下是测试可用但海量数据下会触发内存溢出的代码:

public void leerDatos() throws IOException {
    String user = "postgres";
    String pss = "admin123";
    
    empleadoList = new ArrayList<>();
    Empleados empleado = null;
    ResultSet result = null;
    PreparedStatement stm = null;
    Connection connect = null;
    
    final int pageSize = 100000;
    int offset = 0;
    try {
        connect = DriverManager.getConnection("jdbc:postgresql://localhost:5432/bd_test", user, pss);
        PreparedStatement stmCount = connect.prepareStatement("SELECT COUNT(*) from empleados");
        ResultSet resultCount = stmCount.executeQuery();
        int records = 0;
        
        while (resultCount.next()) {
            records = resultCount.getInt(1);
        }
    
        while (offset < records) {
            String query = "SELECT * FROM empleados ORDER BY id LIMIT ? OFFSET ?";
            stm = connect.prepareStatement(query);
            stm.setInt(1, pageSize);
            stm.setInt(2, offset);
            result = stm.executeQuery();
            while (result.next()) {
                empleado = new Empleados();
                empleado.setNombre(result.getString("nombre"));
                empleado.setApellido(result.getString("apellido"));
                empleado.setSalario(result.getString("salario"));
                empleadoList.add(empleado);
            }
            offset += pageSize;
           
        }
        if (!empleadoList.isEmpty()) {
            
            insertarDatos(empleadoList);
        }
        
    } catch (SQLException ex) {
        System.out.println("SQL Error: " + ex.getMessage());
    } finally {
        try {
            if (result != null) {
                result.close();
            }
            if (stm != null) {
                stm.close();
            }
            if (connect != null) {
                connect.close();
            }
            
        } catch (SQLException ex) {
            System.out.println("SQL Error en cierre: " + ex.getMessage());
        }
    }
}

一、线程实现的可行性与具体方案

线程方案完全可行,但核心是避免内存累积,不能像原代码那样把所有数据加载到全局List后再处理。具体方案如下:

1. 生产者-消费者模式

  • 生产者线程:负责分页查询数据,每次查询一页(10万条)后,将该页数据封装成List提交到阻塞队列,然后立即释放当前页的ResultSet、Statement资源,避免内存占用。查询逻辑改用主键范围查询替代OFFSET(原OFFSET在数据量极大时会越来越慢,且数据库需要扫描大量前置数据):
    // 替换原LIMIT OFFSET的查询语句
    String query = "SELECT * FROM empleados WHERE id > ? ORDER BY id LIMIT ?";
    // 每次用上一页的最大id作为下一页的起始条件,不需要OFFSET
    
  • 消费者线程:从阻塞队列中取出数据块,调用insertarDatos批量插入。根据目标数据库的写入能力设置线程数(比如2-4个,避免连接过载)。
  • 关键注意点:
    • 阻塞队列设置合理容量(比如5-10个数据块),防止生产者过快导致内存溢出。
    • 每个数据块插入完成后,主动清空List并触发GC(调用System.gc(),或让List变量超出作用域)。

2. 线程池管理

用ThreadPoolExecutor创建固定大小的线程池,统一管理生产者和消费者线程,避免手动创建线程带来的资源泄漏问题。


二、替代方案

1. 流式分页处理(最直接的内存优化)

不需要线程,直接修改原代码,每页查询完成后立即插入,不存储全局List:

public void leerDatos() throws IOException {
    String user = "postgres";
    String pss = "admin123";
    
    ResultSet result = null;
    PreparedStatement stm = null;
    Connection connect = null;
    
    final int pageSize = 100000;
    int offset = 0;
    try {
        connect = DriverManager.getConnection("jdbc:postgresql://localhost:5432/bd_test", user, pss);
        PreparedStatement stmCount = connect.prepareStatement("SELECT COUNT(*) from empleados");
        ResultSet resultCount = stmCount.executeQuery();
        int records = 0;
        
        while (resultCount.next()) {
            records = resultCount.getInt(1);
        }
    
        while (offset < records) {
            String query = "SELECT * FROM empleados ORDER BY id LIMIT ? OFFSET ?";
            stm = connect.prepareStatement(query);
            stm.setInt(1, pageSize);
            stm.setInt(2, offset);
            // 设置fetchSize,让驱动流式获取数据,避免一次性加载整页到内存
            stm.setFetchSize(1000);
            result = stm.executeQuery();
            
            List<Empleados> pageList = new ArrayList<>(pageSize);
            while (result.next()) {
                Empleados empleado = new Empleados();
                empleado.setNombre(result.getString("nombre"));
                empleado.setApellido(result.getString("apellido"));
                empleado.setSalario(result.getString("salario"));
                pageList.add(empleado);
            }
            // 每页查询完立即插入
            if (!pageList.isEmpty()) {
                insertarDatos(pageList);
                pageList.clear(); // 主动清空,帮助GC回收
            }
            offset += pageSize;
           
        }
        
    } catch (SQLException ex) {
        System.out.println("SQL Error: " + ex.getMessage());
    } finally {
        try {
            if (result != null) {
                result.close();
            }
            if (stm != null) {
                stm.close();
            }
            if (connect != null) {
                connect.close();
            }
            
        } catch (SQLException ex) {
            System.out.println("SQL Error en cierre: " + ex.getMessage());
        }
    }
}

这种方式内存中始终只保留一页的数据,彻底解决内存溢出问题,实现成本最低。

2. 数据库原生同步工具

如果两个数据库都是PostgreSQL,优先用原生工具,效率远高于Java程序:

  • COPY命令:将源表数据导出为CSV文件,再用COPY命令导入目标表,适合一次性全量同步。
  • 逻辑复制:通过pg_logical插件实现增量同步,实时将源库的变更同步到目标历史表,适合持续同步场景。

3. JDBC驱动参数优化

PostgreSQL JDBC驱动默认会将整个ResultSet加载到内存,设置fetchSize参数让驱动流式获取数据:

stm.setFetchSize(1000); // 每次从数据库拉取1000条记录,而不是整页

结合流式处理,进一步降低内存占用。


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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.24 14:40:33