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

如何拆分Exchange的List类型Body为非单元素批量片段以优化内存?

Apache Camel处理大CSV内存优化:批量拆分替代单行拆分

问题背景

当前使用Apache Camel处理大CSV文件时,路由通过split(body())将Exchange的List类型Body拆分为单元素Exchange,产生大量单行Exchange,导致内存占用过高——处理60MB CSV文件时内存占用超过1GB。目标是将Body拆分为更大的批量片段,把内存占用控制在500MB以内。

现有路由代码:

from(url.concat("RAW(" + filesToInclude + ")"))
        .log("Start processing file: ${header.CamelFileName}")
        .routeId(CamelRoutes.LIST_ROUTE_ID + endpointConfig.getConfigId())
        .setHeader(CamelHeaders.ENDPOINT_CONFIG, constant(endpointConfig))

        .unmarshal(getCsvDataFormat(endpointConfig))

        .split(body())
            .streaming().parallelProcessing().executorService(Executors.newFixedThreadPool(CONCURRENT_THREADS_NUMBER))
            .stopOnException().stopOnAggregateException()
            .marshal().json(JsonLibrary.Jackson)
            .log("Starting aggregating for batch processing")
            .aggregate(header(Exchange.FILE_NAME), new ListAggregationStrategy())
                .completionPredicate(new BatchSizePredicate(endpointConfig.getMaxBatchSize()))
                .completionTimeout(endpointConfig.getBatchTimeout())
                .log("Start one batch processing")
                .bean(rowCamelService, "saveInDB")
                .log("Finished one batch processing")
            .end()
            .log("Finished aggregating for batch processing")
            .log("Finished processing file: ${header.CamelFileName}")
        .end();

优化方案

核心思路是直接将List类型的Body拆分为指定大小的子列表批量,避免生成大量单行Exchange,同时保留并行处理能力。具体改动如下:

  1. 替换单行拆分逻辑:将split(body())改为按批量大小拆分List,减少Exchange实例数量。
  2. 简化中间聚合步骤:拆分阶段直接生成批量数据,无需再聚合单行,降低内存中转开销。

修改后的路由代码

int batchSize = endpointConfig.getMaxBatchSize();

from(url.concat("RAW(" + filesToInclude + ")"))
        .log("Start processing file: ${header.CamelFileName}")
        .routeId(CamelRoutes.LIST_ROUTE_ID + endpointConfig.getConfigId())
        .setHeader(CamelHeaders.ENDPOINT_CONFIG, constant(endpointConfig))

        .unmarshal(getCsvDataFormat(endpointConfig))
        // 关键:将List拆分为指定大小的子列表批量
        .split().method(ListBatchSplitter.class, "splitList(" + batchSize + ")")
            .streaming()
            .parallelProcessing()
            .executorService(Executors.newFixedThreadPool(CONCURRENT_THREADS_NUMBER))
            .stopOnException()
            .marshal().json(JsonLibrary.Jackson)
            .log("Start processing batch of size ${body.size()}")
            // 直接调用批量保存方法(需确保saveInDB支持List参数)
            .bean(rowCamelService, "saveInDB")
            .log("Finished processing batch of size ${body.size()}")
        .end()
        .log("Finished processing file: ${header.CamelFileName}");

关键细节说明

  • 批量拆分工具类实现:上述代码用到的ListBatchSplitter用于将大List拆分为指定大小的子列表,示例代码如下:
    public class ListBatchSplitter {
        public <T> List<List<T>> splitList(List<T> input, int batchSize) {
            return IntStream.range(0, (input.size() + batchSize - 1) / batchSize)
                            .mapToObj(i -> input.subList(Math.min(i * batchSize, input.size()), 
                                                        Math.min((i + 1) * batchSize, input.size())))
                            .collect(Collectors.toList());
        }
    }
    
  • 内存优化原理:拆分后每个Exchange携带batchSize行数据,Exchange数量从原来的数万/数十万级降至数百/数千级,大幅减少内存中同时存在的实例数与序列化开销。
  • 并行处理适配:保留parallelProcessing和线程池配置,批量处理的并行度更可控,避免过多线程抢占资源。
  • 方法适配:确保rowCamelService.saveInDB支持接收List类型参数,若无法修改原有方法,可保留原聚合逻辑,但拆分阶段攒批的内存效率更高。

额外优化建议

  • 开启CSV反序列化流模式:检查getCsvDataFormat(endpointConfig)是否设置csvDataFormat.setStreaming(true),避免一次性加载整个CSV到内存。
  • 微调线程池大小:根据批量大小和系统资源调整CONCURRENT_THREADS_NUMBER,减少线程上下文切换开销。
  • 内存监控调参:通过JConsole等工具跟踪内存变化,微调batchSize直到内存占用稳定在500MB以内。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.22 12:43:20