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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.29 04:15:38