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

Spring Webflux WebClient处理50GB大Zip文件的内存友好方案

内存友好处理大Zip包内嵌套Gzip文件并提取数据到DynamoDB

核心解决方案:流式逐层解析压缩文件

不一次性加载整个50GB文件到内存,通过嵌套流式处理逐层解开压缩包,逐行读取JSON数据并提取目标字段,同时批量写入DynamoDB控制内存占用与性能。

具体实现步骤

  1. 用ZipInputStream遍历master.zip内的文件:逐个读取Zip条目,跳过目录,只处理.json.gz后缀的子文件。
  2. 嵌套GZIPInputStream解开子压缩包:直接基于当前Zip条目的流创建GZIP输入流,无需落地到本地文件。
  3. 逐行读取JSON并轻量解析:用BufferedReader逐行读取解压后的文本,结合JSON流式解析库(如Jackson)只提取UserId字段,避免加载整个JSON对象。
  4. 批量写入DynamoDB:攒够一定数量的UserId后批量写入,减少网络请求开销,同时控制内存占用。

代码示例

import java.io.*;
import java.util.zip.*;
import com.amazonaws.services.dynamodbv2.AmazonDynamoDB;
import com.amazonaws.services.dynamodbv2.AmazonDynamoDBClientBuilder;
import com.amazonaws.services.dynamodbv2.document.DynamoDB;
import com.amazonaws.services.dynamodbv2.document.Table;
import com.fasterxml.jackson.core.JsonFactory;
import com.fasterxml.jackson.core.JsonParser;
import com.fasterxml.jackson.core.JsonToken;

public class LargeZipProcessor {
    // 自定义缓冲区大小,控制内存占用(8KB适合多数场景)
    private static final int BUFFER_SIZE = 8192;
    // 批量写入DynamoDB的批次大小,平衡内存和性能
    private static final int BATCH_SIZE = 1000;

    public static void main(String[] args) throws IOException {
        // 初始化DynamoDB客户端(根据ECS环境配置调整)
        AmazonDynamoDB client = AmazonDynamoDBClientBuilder.defaultClient();
        DynamoDB dynamoDB = new DynamoDB(client);
        Table targetTable = dynamoDB.getTable("YourDynamoDBTableName");

        // 外层Zip流,用BufferedInputStream优化读取性能
        try (ZipInputStream zipInputStream = new ZipInputStream(
                new BufferedInputStream(new FileInputStream("master.zip"), BUFFER_SIZE))) {

            ZipEntry entry;
            while ((entry = zipInputStream.getNextEntry()) != null) {
                // 跳过目录,只处理目标子压缩文件
                if (!entry.isDirectory() && entry.getName().endsWith("child.json.gz")) {
                    processGzipChildFile(zipInputStream, targetTable);
                    zipInputStream.closeEntry(); // 关闭当前条目,释放资源
                }
            }
        }
    }

    private static void processGzipChildFile(InputStream inputStream, Table targetTable) throws IOException {
        // 嵌套GZIP流和解压后的字符流
        try (GZIPInputStream gzipInputStream = new GZIPInputStream(
                new BufferedInputStream(inputStream, BUFFER_SIZE));
             BufferedReader lineReader = new BufferedReader(new InputStreamReader(gzipInputStream))) {

            JsonFactory jsonFactory = new JsonFactory();
            String jsonLine;
            int batchCounter = 0;

            while ((jsonLine = lineReader.readLine()) != null) {
                // 流式解析JSON,只提取UserId字段
                try (JsonParser jsonParser = jsonFactory.createParser(jsonLine)) {
                    String userId = null;
                    while (jsonParser.nextToken() != JsonToken.END_OBJECT) {
                        String fieldName = jsonParser.getCurrentName();
                        if ("UserId".equals(fieldName)) {
                            jsonParser.nextToken();
                            userId = jsonParser.getText();
                            break; // 找到目标字段后立即停止解析,节省资源
                        }
                    }

                    if (userId != null) {
                        // 写入DynamoDB,达到批次大小后可做批量优化(此处为简化示例用单条写入)
                        targetTable.putItem(new com.amazonaws.services.dynamodbv2.document.Item()
                                .withString("UserId", userId));

                        if (++batchCounter >= BATCH_SIZE) {
                            // 可选:添加批次完成日志或等待异步写入完成
                            batchCounter = 0;
                        }
                    }
                }
            }
        }
    }
}

关键优化点说明

  • 内存控制:所有流都使用固定大小的缓冲区,不会加载整个文件到内存;JSON流式解析只提取所需字段,避免解析整个JSON对象。
  • 解决EOF错误:ZipInputStream会自动处理每个Zip条目的边界,不会出现分块截断导致的ZipException: unexpected EOF问题——之前的DataBuffer方案可能因为没有识别Zip条目边界,导致分块截断了压缩数据。
  • 性能优化:批量写入DynamoDB减少网络请求次数;BufferedInputStream和BufferedReader减少磁盘IO次数。

额外建议

  • 如果ECS容器内存仍紧张,可以进一步缩小缓冲区大小(如4KB),或降低DynamoDB批量写入的批次大小。
  • 加入错误处理逻辑:比如捕获JSON解析错误、DynamoDB写入失败,将异常数据暂存到SQS或本地文件后续重试。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.18 15:07:24