Spring Boot读取S3 Athena中40GB+CSV数据失败问题排查
问题场景
在Spring Boot项目中读取S3 Athena生成的40GB以上CSV文件时,系统抛出如下警告:
2024-06-25 18:16:25.447 WARN 509917 --- [pool-6-thread-1] c.a.s.s.internal.S3AbortableInputStream : Not all bytes were read from the S3ObjectInputStream, aborting HTTP connection. This is likely an error and may result in sub-optimal behavior. Request only the bytes you need via a ranged GET or drain the input stream after use.
核心处理代码如下:
private static int saveResultFile(S3Object s3Object, String outputFile, List<ColumnInfo> columnInfoList) { System.out.println("This is the save file method."); System.out.println(outputFile + "this is output"); String lastTwoDigits = outputFile.substring(outputFile.length() - 2); boolean hasData = false; // Initialize to false int totalSum = 0; // Initialize totalSum before the loop // Set the maximum records per page int maxRecordsPerPage = 1000000; try (BufferedReader reader = new BufferedReader( new InputStreamReader(s3Object.getObjectContent(), StandardCharsets.UTF_8)); ZipOutputStream zipOutputStream = new ZipOutputStream(new FileOutputStream(outputFile + ".zip"))) { CSVReader csvReader = new CSVReader(reader); String[] header = csvReader.readNext(); if (header == null) { // No header, return false System.out.println("No header found in the CSV file."); return NO_HEADERS_FOUND; } int recordsProcessed = 0; int pageNumber = 1; CSVWriter writer = null; try { // Create a ZipOutputStream to write to the zip file writer = new CSVWriter(new FileWriter(outputFile + "_" + pageNumber + ".csv"), ',', CSVWriter.DEFAULT_QUOTE_CHARACTER); // Write the header for the current page writer.writeNext(header); String[] line; while ((line = csvReader.readNext()) != null) { hasData = true; // Set to true if any data is processed if (lastTwoDigits.equals("02")) { // Calculate sum for sms_split column if (line.length > 8) { // Assuming index 7 is the 8th column try { totalSum += Integer.parseInt(line[8]); // Assuming sms_split is an integer } catch (NumberFormatException e) { // Handle parsing error if necessary e.printStackTrace(); } } } int lastColumnIndex = line.length - 1; String lastColumnName = header[lastColumnIndex]; if (isTextColumn(lastColumnName)) { // Perform encryption and decryption operations only if it's a text column String encryptedValue = line[lastColumnIndex]; String decryptedText = decrypt(encryptedValue); if (isHexString(decryptedText)) { decryptedText = hexToAscii(decryptedText); } line[lastColumnIndex] = decryptedText; } if (recordsProcessed >= maxRecordsPerPage) { // Close the current writer and create a new one for the next page writer.close(); pageNumber++; recordsProcessed = 0; writer = new CSVWriter(new FileWriter(outputFile + "_" + pageNumber + ".csv"), ',', CSVWriter.DEFAULT_QUOTE_CHARACTER); // Write the header for the new page writer.writeNext(header); } writer.writeNext(line); recordsProcessed++; } } catch (IOException e) { e.printStackTrace(); System.out.println("Error creating ZIP file: " + outputFile + ".zip"); } finally { if (writer != null) { writer.close(); } } if (hasData) { for (int i = 1; i <= pageNumber; i++) { String pageFileName = outputFile + "_" + i + ".csv"; try (FileInputStream fis = new FileInputStream(new File(pageFileName))) { String zipEntryName = pageFileName.substring(pageFileName.lastIndexOf("/") + 1); ZipEntry zipEntry = new ZipEntry(zipEntryName); zipOutputStream.putNextEntry(zipEntry); byte[] bytes = new byte[1024]; int length; while ((length = fis.read(bytes)) >= 0) { zipOutputStream.write(bytes, 0, length); } zipOutputStream.closeEntry(); if (new File(pageFileName).delete()) { // System.out.println("Deleted file: " + pageFileName); } else { System.out.println("Failed to delete file: " + pageFileName); } } catch (IOException e) { e.printStackTrace(); System.out.println("Error processing file: " + pageFileName); } } System.out.println("Result files saved successfully with decryption."); System.out.println("Total sum of sms_split column: " + totalSum); } else { System.out.println("No data found to save."); return NO_HEADERS_FOUND; } } catch (IOException e) { e.printStackTrace(); } finally { try { if (s3Object != null) { s3Object.close(); } } catch (IOException e) { e.printStackTrace(); } } return totalSum; }
问题原因
警告的核心原因是S3ObjectInputStream未被完全读取就被关闭,导致AWS SDK强制中断HTTP连接。即使使用了try-with-resources自动管理流,若CSV解析器未消费完底层流的所有字节(比如循环提前退出、解析异常未处理),就会触发该警告,可能引发连接池资源浪费或后续请求异常。
修复方案
1. 强制排空剩余流字节
在关闭S3Object之前,手动读取流中剩余的所有字节,确保流被完全消费。修改原代码的finally块:
finally { try { if (s3Object != null) { InputStream objectContent = s3Object.getObjectContent(); // 用缓冲区排空剩余字节,避免残留数据 byte[] drainBuffer = new byte[4096]; while (objectContent.read(drainBuffer) != -1) { // 仅读取不处理,确保流被完全消费 } s3Object.close(); } } catch (IOException e) { e.printStackTrace(); } }
2. 优化资源管理逻辑
将CSVReader也纳入try-with-resources,避免因解析器未关闭导致的流泄漏:
try (BufferedReader reader = new BufferedReader( new InputStreamReader(s3Object.getObjectContent(), StandardCharsets.UTF_8)); ZipOutputStream zipOutputStream = new ZipOutputStream(new FileOutputStream(outputFile + ".zip")); CSVReader csvReader = new CSVReader(reader)) { // 原有业务逻辑... }
3. 分段读取(可选,超大规模文件)
对于40GB以上的超大文件,可采用AWS推荐的范围请求(Ranged GET),按字节分段读取文件,分批处理。这种方式不仅能避免流未完全消费的问题,还能降低单次请求的内存负载,适合超大规模数据场景。
验证
修改后重新运行任务,若警告不再出现,且数据处理正常,则说明修复有效。
内容的提问来源于stack exchange,提问作者prakhar1807

