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

Spring Boot读取S3 Athena中40GB+CSV数据失败问题排查

解决S3 Athena读取大CSV文件时的S3AbortableInputStream警告

问题场景

在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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.21 00:24:51