使用CopyManager导出百万级PostgreSQL记录时遭遇Java堆内存溢出问题
看起来你遇到了典型的内存缓冲过载问题——当处理百万级数据时,把所有CSV内容先存在StringWriter里(本质是内存中的字符串缓冲区),必然会耗尽Java堆内存,这也是堆栈里OutOfMemoryError: Java heap space指向StringWriter.write的核心原因。
先还原你的场景:
- 处理几千行数据时内存足够支撑,百万行量级直接触发OOM
- 当前用
CopyManager.copyOut把查询结果写入PrintWriter(包装了StringWriter),最后一次性转成ByteArrayInputStream上传S3 - 分页参数:OFFSET从0开始,每次取300万条
问题根源分析
你的代码逻辑是把所有查询到的CSV数据先全部加载到内存中,百万级的CSV数据(按每条1KB估算就是3GB)远远超过了JVM默认堆内存大小,直接导致堆内存被耗尽。
解决方案建议
1. 流式写入本地文件,避免全量内存缓冲
不要用StringWriter,直接用文件输出流写本地文件,写完后再上传S3。这样内存里只会保留当前批次的少量数据,不会累积全部内容:
Connection connection = null; try { String now = LocalDate.now().format(DateTimeFormatter.ofPattern("dd/MM/yyyy")).replace("/", ""); String localFileName = table.getTableName() + "_" + now + ".csv"; File localFile = new File(localFileName); try (BufferedWriter writer = new BufferedWriter(new FileWriter(localFile))) { connection = dataSource.getConnection(); long offset = 0; final long limit = 3_000_000; while (offset <= totalRecords && totalRecords > 0) { if (connection != null && connection.isWrapperFor(PGConnection.class)) { PGConnection pgConnection = connection.unwrap(PGConnection.class); CopyManager copyManager = pgConnection.getCopyAPI(); String sql = "SELECT " + table.getAttributeList() + " FROM " + table.getSchemaName().trim() + "." + table.getTableName().trim() + " WHERE modified_ts > '" + dateTime + "'" + " OFFSET " + offset + " LIMIT " + limit; LOGGER.info(sql); long rowsCopied = copyManager.copyOut("COPY (" + sql + " ) TO STDOUT WITH (FORMAT CSV)", writer); LOGGER.info("Total no of records in {} : {}", table.getTableName(), rowsCopied); writer.flush(); offset += limit; } } } // 上传本地文件到S3 try (InputStream inputStream = new FileInputStream(localFile)) { postObjectToS3.uploadFile(saveFilePath + localFileName, inputStream); } // 可选:删除本地临时文件 localFile.delete(); } catch (Exception e) { LOGGER.error("导出失败", e); throw e; } finally { if (connection != null) { try { connection.close(); } catch (SQLException e) { LOGGER.error("关闭连接失败", e); } } }
2. 优化分页逻辑:替换OFFSET为键值分页
虽然你用了OFFSET分页,但PostgreSQL处理大OFFSET时性能很差——它需要扫描前面所有的行才能定位到OFFSET的位置。建议用**主键或有序字段(比如modified_ts)**来做分页,比如:
// 初始条件 long lastModifiedTs = 0; // 或者你的起始dateTime对应的时间戳 final long limit = 3_000_000; while (true) { String sql = "SELECT " + table.getAttributeList() + ", modified_ts FROM " + table.getSchemaName().trim() + "." + table.getTableName().trim() + " WHERE modified_ts > '" + dateTime + "'" + " AND modified_ts > " + lastModifiedTs + " ORDER BY modified_ts LIMIT " + limit; // 执行COPY导出,同时解析当前批次的最大modified_ts // 当返回的rowsCopied < limit时,说明没有更多数据,退出循环 long rowsCopied = copyManager.copyOut("COPY (" + sql + " ) TO STDOUT WITH (FORMAT CSV)", writer); if (rowsCopied < limit) { break; } // 这里需要额外查询当前批次的最大modified_ts,用于下一轮分页 // 可以通过单独的SELECT语句获取,或者在COPY的结果中解析最后一行的modified_ts lastModifiedTs = getLastModifiedTsFromCurrentBatch(); }
这种方式避免了大OFFSET的性能损耗,同时也让分页逻辑更稳定。
3. 直接流式上传到S3(无需本地文件)
如果不想写本地文件,可以用AWS S3的TransferManager或者直接用PutObjectRequest接受输出流,把CopyManager的输出直接写到S3的上传流中:
Connection connection = null; try { String now = LocalDate.now().format(DateTimeFormatter.ofPattern("dd/MM/yyyy")).replace("/", ""); String s3FileName = saveFilePath + table.getTableName() + "_" + now + ".csv"; // 创建S3上传的输出流 TransferManager transferManager = TransferManagerBuilder.defaultTransferManager(); // 构建合并分页结果的输入流 List<InputStream> batchStreams = new ArrayList<>(); connection = dataSource.getConnection(); PGConnection pgConnection = connection.unwrap(PGConnection.class); CopyManager copyManager = pgConnection.getCopyAPI(); long offset = 0; final long limit = 3_000_000; while (offset <= totalRecords && totalRecords > 0) { String sql = "SELECT " + table.getAttributeList() + " FROM " + table.getSchemaName().trim() + "." + table.getTableName().trim() + " WHERE modified_ts > '" + dateTime + "'" + " OFFSET " + offset + " LIMIT " + limit; String copySql = "COPY (" + sql + " ) TO STDOUT WITH (FORMAT CSV)"; batchStreams.add(copyManager.copyOut(copySql)); offset += limit; } // 合并所有批次的输入流 InputStream combinedStream = new SequenceInputStream(Collections.enumeration(batchStreams)); // 上传到S3 PutObjectRequest putObjectRequest = new PutObjectRequest("your-bucket-name", s3FileName, combinedStream, new ObjectMetadata()); Upload upload = transferManager.upload(putObjectRequest); upload.waitForCompletion(); } catch (Exception e) { LOGGER.error("导出上传失败", e); throw e; } finally { if (connection != null) { try { connection.close(); } catch (SQLException e) { LOGGER.error("关闭连接失败", e); } } }
这种方式完全跳过内存和本地文件,直接把PostgreSQL的导出流传到S3,内存占用极低。
4. 临时应急:调整JVM堆内存
如果只是临时解决,可以增大JVM堆内存,比如启动参数加上-Xmx4G(分配4GB堆内存),但这只是治标不治本,当数据量继续增长(比如10亿行)还是会OOM,优先推荐前面的流式方案。
你的原始代码和堆栈信息
原始代码
Connection connection = null; try(StringWriter out = new StringWriter(); Writer printWriter = new PrintWriter(out, true)) { connection = dataSource.getConnection(); OFFSET = 0; // totalRecords is the total rows in the table while (OFFSET <= totalRecords && totalRecords > 0) { if (connection != null && connection.isWrapperFor(PGConnection.class)) { PGConnection pgConnection = connection.unwrap(PGConnection.class); CopyManager copyManager = null; copyManager = pgConnection.getCopyAPI(); String sql = null; sql = "SELECT " + table.getAttributeList() + " FROM " + table.getSchemaName().trim() + "." + table.getTableName().trim() + " WHERE " + "modified_ts" + " > " + "'" + dateTime + "'" + " OFFSET " + OFFSET + " LIMIT " + LIMIT; LOGGER.info(sql); long i; // Here I am trying to copy the Result into the printwriter i = copyManager.copyOut("COPY (" + sql + " ) TO STDOUT WITH (FORMAT CSV)", printWriter); LOGGER.info("Total no of records in {} : {}", table.getTableName(), i); printWriter.flush(); OFFSET = OFFSET + LIMIT; } } String now = LocalDate.now().format(DateTimeFormatter.ofPattern("dd/MM/yyyy")).replace("/", ""); String localFileName = table.getTableName() + "_" + now + ".csv"; InputStream inputStream = new ByteArrayInputStream(out.toString().getBytes(UTF8)); postObjectToS3.uploadFile(saveFilePath + localFileName, inputStream); inputStream.close(); }
堆栈信息
Exception in thread "main" java.lang.reflect.InvocationTargetException at java.base/jdk.internal.reflect.NativeMethodAccessorImpl.invoke0(Native Method) at java.base/jdk.internal.reflect.NativeMethodAccessorImpl.invoke(NativeMethodAccessorImpl.java:62) at java.base/jdk.internal.reflect.DelegatingMethodAccessorImpl.invoke(DelegatingMethodAccessorImpl.java:43) at java.base/java.lang.reflect.Method.invoke(Method.java:566) at org.springframework.boot.loader.MainMethodRunner.run(MainMethodRunner.java:49) at org.springframework.boot.loader.Launcher.launch(Launcher.java:108) at org.springframework.boot.loader.Launcher.launch(Launcher.java:58) at org.springframework.boot.loader.JarLauncher.main(JarLauncher.java:88) Caused by: java.lang.OutOfMemoryError: Java heap space at java.base/java.util.Arrays.copyOf(Arrays.java:3745) at java.base/java.lang.AbstractStringBuilder.ensureCapacityInternal(AbstractStringBuilder.java:172) at java.base/java.lang.AbstractStringBuilder.append(AbstractStringBuilder.java:633) at java.base/java.lang.StringBuffer.append(StringBuffer.java:397) at java.base/java.io.StringWriter.write(StringWriter.java:122) at java.base/java.io.PrintWriter.write(PrintWriter.java:542) at java.base/java.io.PrintWriter.write(PrintWriter.java:559) at org.postgresql.copy.CopyManager.copyOut(CopyManager.java:92)
备注:内容来源于stack exchange,提问作者Himanshu Arya

