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

