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

如何使用Camel流式处理大型JSON列表避免内存溢出(OOM)

解决Camel中流式处理超大JSON列表的内存问题

要处理超大JSON列表且不加载全量数据到内存,核心是流式解析输入+实时流式写入输出,同时维护JSON列表的合法结构。针对你的两次尝试的问题,给出以下解决方案:

问题根源分析

  • 第一次尝试:聚合器会把所有处理后的条目存储在内存中,最后一次性序列化,必然导致OOM。
  • 第二次尝试:直接追加单个JSON对象,没有处理列表的[开头、条目间的逗号和]结尾,生成的不是合法JSON数组。

方案1:基于Camel DSL的流式聚合实现

关键要点

  1. 确保JSONPath流式解析:给jsonpath添加streaming=true参数,避免加载整个JSON到内存。
  2. 自定义流式聚合策略:不缓存结果,实时写入文件,同时维护JSON数组的结构。
  3. 单个条目序列化时不要包装成列表: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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.31 12:05:33