Spring Webflux WebClient处理50GB大Zip文件的内存友好方案
内存友好处理大Zip包内嵌套Gzip文件并提取数据到DynamoDB
核心解决方案:流式逐层解析压缩文件
不一次性加载整个50GB文件到内存,通过嵌套流式处理逐层解开压缩包,逐行读取JSON数据并提取目标字段,同时批量写入DynamoDB控制内存占用与性能。
具体实现步骤
- 用
ZipInputStream遍历master.zip内的文件:逐个读取Zip条目,跳过目录,只处理.json.gz后缀的子文件。 - 嵌套
GZIPInputStream解开子压缩包:直接基于当前Zip条目的流创建GZIP输入流,无需落地到本地文件。 - 逐行读取JSON并轻量解析:用
BufferedReader逐行读取解压后的文本,结合JSON流式解析库(如Jackson)只提取UserId字段,避免加载整个JSON对象。 - 批量写入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
相关产品推荐
相关产品推荐

