如何拆分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,同时保留并行处理能力。具体改动如下:
- 替换单行拆分逻辑:将
split(body())改为按批量大小拆分List,减少Exchange实例数量。 - 简化中间聚合步骤:拆分阶段直接生成批量数据,无需再聚合单行,降低内存中转开销。
修改后的路由代码
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
相关产品推荐
相关产品推荐

