Camel 3.11迁移:如何用Split、Aggregate配置单线程工作流?
解决Camel 3.2→3.21迁移中文件串行处理失效的问题
核心问题定位
Camel 3.11版本后,Aggregate EIP默认启用专用工作线程池(即使parallelProcessing=false),导致Split/Aggregate的线程与后续块处理线程分离,文件消费者线程未被阻塞,因此提前读取下一个文件。而使用SynchronousExecutorService时,因线程上下文绑定或UnitOfWork传递问题导致挂起。
具体解决思路
1. 从文件组件层面强制串行
直接通过文件端点配置确保只有当前文件完全处理完毕才会读取下一个:
- 设置
readLock=markerFile:处理文件时创建标记文件,处理完成后删除,未删除时不会读取下一个文件;配合readLockTimeout=0(无限等待锁)。 - 添加
maxMessagesPerPoll=1:每次轮询仅处理一个文件。 - 示例配置(Java DSL):
from("file:/input?readLock=markerFile&readLockTimeout=0&maxMessagesPerPoll=1") .split(xpath("//fragment")) // ...后续Split/Aggregate逻辑
2. 修正Aggregate线程池配置,避免挂起
放弃SynchronousExecutorService,改用可控的单线程池,同时确保上下文传递:
- 注册单线程Executor Bean(以Spring为例):
@Bean(name = "singleAggregatePool") public ExecutorService singleAggregatePool() { return Executors.newSingleThreadExecutor(); } - 在Aggregate中指定该线程池,同时开启
shareUnitOfWork:.aggregate(header("CamelFileName"), new MyAggregationStrategy()) .completionSize(100) // 你的块大小 .executorService("#singleAggregatePool") .shareUnitOfWork(true) .parallelProcessing(false) .to("direct:processBlock")
3. 用全局路由策略限制并发
通过ThrottlingInflightRoutePolicy限制路由内同时处理的Exchange数量为1,从全局层面保证串行:
- 配置路由策略:
这样整个路由同一时间只会处理一个文件的所有流程,彻底避免并发读取。ThrottlingInflightRoutePolicy policy = new ThrottlingInflightRoutePolicy(); policy.setMaxInflightExchanges(1); route.setRoutePolicy(policy);
4. 检查Split与Aggregate的上下文关联
确保Split后的片段正确关联到原文件,避免聚合逻辑混乱:
- Split时设置
shareUnitOfWork=true,保证UnitOfWork在拆分片段间传递; - Aggregate的
correlationExpression必须准确关联同一个文件的片段(比如用文件名字段),防止跨文件聚合; - 开启Split的
stopOnException=true,若文件处理失败则终止,避免后续文件提前读取。
内容的提问来源于stack exchange,提问作者Max Vasileuski
相关产品推荐
相关产品推荐

