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

