如何使用Camel流式处理大型JSON列表避免内存溢出(OOM)
解决Camel中流式处理超大JSON列表的内存问题
要处理超大JSON列表且不加载全量数据到内存,核心是流式解析输入+实时流式写入输出,同时维护JSON列表的合法结构。针对你的两次尝试的问题,给出以下解决方案:
问题根源分析
- 第一次尝试:聚合器会把所有处理后的条目存储在内存中,最后一次性序列化,必然导致OOM。
- 第二次尝试:直接追加单个JSON对象,没有处理列表的
[开头、条目间的逗号和]结尾,生成的不是合法JSON数组。
方案1:基于Camel DSL的流式聚合实现
关键要点
- 确保JSONPath流式解析:给
jsonpath添加streaming=true参数,避免加载整个JSON到内存。 - 自定义流式聚合策略:不缓存结果,实时写入文件,同时维护JSON数组的结构。
- 单个条目序列化时不要包装成列表:
marshal时设置useList=false,避免每个条目生成[{}]格式。
步骤1:实现流式聚合策略
public class StreamingJsonListAggregator implements AggregationStrategy { private boolean isFirstEntry = true; private final BufferedWriter writer; public StreamingJsonListAggregator() throws IOException { this.writer = new BufferedWriter(new FileWriter("/temp/to.json")); writer.write("["); // 先写入数组开头 } @Override public Exchange aggregate(Exchange oldExchange, Exchange newExchange) { try { String processedJson = newExchange.getIn().getBody(String.class); if (!isFirstEntry) { writer.write(","); // 非第一个条目前加逗号 } writer.write(processedJson); isFirstEntry = false; writer.flush(); // 实时写入磁盘,避免内存缓存 return newExchange; } catch (IOException e) { throw new RuntimeException("Failed to write processed entry", e); } } // 拆分完成后写入数组结尾 public void finish() throws IOException { writer.write("]"); writer.close(); } }
步骤2:配置Camel路由
<bean id="streamingJsonListAggregator" class="com.something.StreamingJsonListAggregator"/> <route> <from uri="file:/mytestdirectory/?fileName=input.json"/> <!-- 流式拆分JSON列表,streaming=true确保不加载全量数据 --> <split streaming="true"> <jsonpath expression="$[*]" streaming="true" writeAsString="true"/> <unmarshal> <json unmarshalType="com.something.InputItem" namingStrategy="UPPER_CAMEL_CASE"/> </unmarshal> <bean beanType="com.something.Processor" method="doSomething"/> <marshal> <json namingStrategy="LOWER_CAMEL_CASE" useList="false"/> <!-- 单个对象序列化,不包装列表 --> </marshal> <!-- 聚合时实时写入文件,completionFromBatchConsumer确保拆分完成后触发finish --> <aggregate strategy="#streamingJsonListAggregator" completionFromBatchConsumer="true"> <constant>true</constant> <!-- 所有条目归为同一聚合组 --> <onCompletion> <method ref="streamingJsonListAggregator" method="finish"/> </onCompletion> </aggregate> </split> </route>
方案2:基于Jackson流式API的底层实现
如果追求极致的内存控制,直接用Jackson的流式API写自定义Processor,完全绕开Camel组件可能的内存开销:
自定义Processor
public class StreamingJsonProcessor implements Processor { @Override public void process(Exchange exchange) throws Exception { File inputFile = exchange.getIn().getBody(File.class); File outputFile = new File("/temp/to.json"); // 流式读取输入,流式写入输出 try (JsonParser parser = new JsonFactory().createParser(inputFile); JsonGenerator generator = new JsonFactory().createGenerator(outputFile, JsonEncoding.UTF8)) { generator.writeStartArray(); // 写入输出数组开头 parser.nextToken(); // 跳过输入的START_ARRAY标记 // 逐个读取并处理对象 while (parser.nextToken() == JsonToken.START_OBJECT) { InputItem item = parser.readValueAs(InputItem.class); OutputItem processedItem = new com.something.Processor().doSomething(item); generator.writeObject(processedItem); // 实时写入处理后的对象 } generator.writeEndArray(); // 写入输出数组结尾 } } }
路由配置
<route> <from uri="file:/mytestdirectory/?fileName=input.json"/> <bean beanType="com.something.StreamingJsonProcessor" method="process"/> </route>
核心注意事项
- 无论哪种方案,都要确保使用流式JSON解析器,避免将整个JSON文件加载到内存。
- 输出时必须维护JSON数组的结构,手动处理开头、逗号分隔和结尾。
- 实时写入磁盘并刷新缓冲区,避免中间结果在内存中积累。
内容的提问来源于stack exchange,提问作者Marcanpilami
相关产品推荐
相关产品推荐

