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

Java从S3处理大CSV文件报错,求高效无错处理方案

问题描述

尝试从Amazon S3获取并处理130,000 KB的大型CSV文件时,流式读取处理耗时超1小时,且中途频繁出现连接断开错误;改用TransferManager下载文件时,又触发400 Bad Request错误。

原流式读取代码

try (Reader reader = new InputStreamReader(s3ServiceStock.getObject(key).getObjectContent(), StandardCharsets.UTF_8);
     CSVReader csvReader = new CSVReaderBuilder(reader).withCSVParser(parser).build()) {
    importContent(csvReader);
} catch (Exception e) {
    log.error("Error", e);
}

原CSV处理逻辑

private void importContent(CSVReader csvReader) throws IOException {
    String[] nextRecord;
    while ((nextRecord = csvReader.readNext()) != null) {
        // 记录转换逻辑
        repository.save(entity);
    }
}

触发的连接错误

org.apache.http.ConnectionClosedException: Premature end of Content-Length delimited message body (expected: 130,272,542; received: 22,061,056)
    at org.apache.http.impl.io.ContentLengthInputStream.read(ContentLengthInputStream.java:178) ~[httpcore-4.4.11.jar:4.4.11]
    // 省略后续栈追踪信息

原TransferManager下载代码(触发400错误)

@Override
public void downloadFileWithTransferManager(String key, String downloadFilePath){
    TransferManager transferManager = TransferManagerBuilder.standard()
            .withS3Client(s3)
            .withExecutorFactory(() -> Executors.newFixedThreadPool(5))
            .withMinimumUploadPartSize(Long.valueOf(5L * 1024 * 1024))
            .withMultipartUploadThreshold(Long.valueOf(10L * 1024 * 1024))
            .build();
    Download download = transferManager.download(bucketName, key, new File(downloadFilePath));
    try{
        download.waitForCompletion();
    } catch (Exception e) {
        LOG.log(Level.FINER, e.getMessage());
    }finally {
        transferManager.shutdownNow();
    }
}

解决方案

1. 修复TransferManager的400错误

原配置错误使用了上传分片参数,而非下载相关参数,同时补充S3客户端的基础容错配置:

@Override
public void downloadFileWithTransferManager(String key, String downloadFilePath){
    // 配置S3客户端的超时与重试策略
    ClientConfiguration clientConfig = new ClientConfiguration()
            .withConnectionTimeout(5000)    // 连接超时5秒
            .withSocketTimeout(30000)       // 读取超时30秒
            .withMaxErrorRetry(3);          // 失败重试3次

    AmazonS3 s3Client = AmazonS3ClientBuilder.standard()
            .withClientConfiguration(clientConfig)
            .withRegion(Regions.YOUR_REGION) // 替换为你的S3桶所在区域
            .build();

    TransferManager transferManager = TransferManagerBuilder.standard()
            .withS3Client(s3Client)
            .withExecutorFactory(() -> Executors.newFixedThreadPool(5))
            .withMultipartDownloadThreshold(10L * 1024 * 1024) // 超过10MB启用分片下载
            .withMinimumDownloadPartSize(5L * 1024 * 1024)     // 每个分片大小5MB
            .build();

    Download download = transferManager.download(bucketName, key, new File(downloadFilePath));
    try{
        download.waitForCompletion();
    } catch (Exception e) {
        LOG.log(Level.SEVERE, "S3文件下载失败", e); // 调高日志级别,便于排查具体错误
    }finally {
        transferManager.shutdownNow();
    }
}

核心修复点:

  • 替换上传参数为下载专用的withMultipartDownloadThreshold和withMinimumDownloadPartSize
  • 增加客户端超时与重试配置,提升网络容错性
  • 调整日志级别,获取400错误的完整详情(如区域不匹配、签名错误等)

2. 优化CSV处理效率

即使文件下载到本地,单条插入数据库仍会导致处理过慢,改用批量插入优化:

private void importContentFromLocalFile(String localFilePath) throws IOException {
    int batchSize = 1000; // 根据数据库性能调整批量大小
    List<YourEntity> entityBatch = new ArrayList<>(batchSize);

    try (Reader reader = Files.newBufferedReader(Paths.get(localFilePath), StandardCharsets.UTF_8);
         CSVReader csvReader = new CSVReaderBuilder(reader).withCSVParser(parser).build()) {

        String[] nextRecord;
        while ((nextRecord = csvReader.readNext()) != null) {
            // 转换CSV记录为实体
            YourEntity entity = convertToEntity(nextRecord);
            entityBatch.add(entity);

            // 达到批量阈值时执行批量插入
            if (entityBatch.size() >= batchSize) {
                repository.saveAll(entityBatch);
                entityBatch.clear();
            }
        }

        // 插入剩余未批量的记录
        if (!entityBatch.isEmpty()) {
            repository.saveAll(entityBatch);
        }
    }
}

优化点:

  • 使用saveAll批量插入替代单条save,大幅减少数据库IO次数
  • 本地文件读取避免了S3连接长时间保持的超时风险

3. 额外优化建议

  • 若仅需CSV部分列,可使用S3 Select筛选数据,减少传输量
  • 开启数据库批量插入优化(如MySQL的rewriteBatchedStatements=true参数)
  • 处理过程中增加进度日志,便于监控处理状态

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.15 18:15:08