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

使用CopyManager导出百万级PostgreSQL记录时遭遇Java堆内存溢出问题

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.04.23 14:12:45