Apache Camel 3.10.0与3.11.0流式处理完成时机差异问题
Camel 3.10.0→3.11.0 版本间并行路由执行顺序问题解决方案
问题现象
使用camel-spring-boot-dependencies 3.10.0时,route_stage中的"This is the end"日志会在route_test的所有拆分、聚合及富化逻辑完成后打印;升级到3.11.0后,该日志提前输出,富化路由(route_enrich_A/route_enrich_B)的日志反而在其后出现。需要恢复3.10.0的行为,确保所有流式处理环节完成后再触发结束日志。
相关代码
路由定义:
from(quartz("quartz") // Prevent concurrent calls .stateful(true) .cron("0/5 * * * * ?") .autoStartScheduler(true) ).id("quartz") .autoStartup(true) .startupOrder(1) .process(utils::addElement) .log("Quartz lance") .to(direct("route_intermediaire")) .delay(300000); from(direct("route_stage")).id("route_stage").to(direct("route_test")).log("This is the end"); from(direct("route_enrich_A")) .id("route_enrich_A") .delay(1000).syncDelayed() .to(log("route_enrich_A")); from(direct("route_enrich_B")) .id("route_enrich_B") .delay(1000).syncDelayed() .to(log("route_enrich_B")); from(direct("route_test")) .id("route_test") .split(body()) .streaming().parallelProcessing() .threads(5) .aggregate(AggregationStrategies.flexible().accumulateInCollection(ArrayList.class)) .constant(true) .completionSize(2) .completionInterval(1000) .enrich(direct("route_enrich_A")) .enrich(direct("route_enrich_B")) .end();
数据生成代码:
public static void addElement(Exchange exchange) { List<Integer> liste = new ArrayList<>(); for (Integer i=0;i<10;i++){ liste.add(i); } exchange.getIn().setBody(liste); }
问题原因
Camel 3.11.x版本对并行拆分(parallelProcessing)后的聚合(aggregate)与富化(enrich)的同步逻辑做了优化调整:3.10.0中聚合完成后会同步等待所有富化路由执行完毕再返回上游;3.11.0中默认将富化逻辑转为异步执行,导致上游route_stage无需等待即可继续打印结束日志。
解决方案
方案1:显式指定富化为同步执行
在route_test的enrich方法中添加第三个参数false,强制使用同步模式(默认enrich为同步,但并行场景下可能被自动优化为异步):
from(direct("route_test")) .id("route_test") .split(body()) .streaming().parallelProcessing() .threads(5) .aggregate(AggregationStrategies.flexible().accumulateInCollection(ArrayList.class)) .constant(true) .completionSize(2) .completionInterval(1000) // 显式指定同步富化 .enrich(direct("route_enrich_A"), AggregationStrategies.passThrough(), false) .enrich(direct("route_enrich_B"), AggregationStrategies.passThrough(), false) .end();
方案2:在route_stage中添加同步等待
修改route_stage,在调用route_test后添加.sync(),确保等待所有下游逻辑完成:
from(direct("route_stage")) .id("route_stage") .to(direct("route_test")) // 强制同步等待下游路由完成 .sync() .log("This is the end");
额外优化:修正聚合完成条件
当前completionSize(2)结合completionInterval(1000)会导致聚合在收集到2个元素或1秒超时后提前完成,与addElement生成的10个元素不符。若需求是等待所有拆分项处理完成,建议调整为:
.completionSize(10) // 对应生成的10个元素 // 或使用completionFromBatchConsumer()自动匹配拆分的总数量 .completionFromBatchConsumer()
内容的提问来源于stack exchange,提问作者Samshay
相关产品推荐
相关产品推荐

