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

Camel如何实现direct端点的真正同步,解决aggregate异步执行问题

Camel Direct路由搭配Aggregate组件消息无序问题

环境信息

  • 框架版本:Camel 3.11 ~ 3.12.0-SNAPSHOT
  • JDK版本:Java 16

需求说明

本示例中timer会发送两条消息,当前第一条消息尚未完成处理时,第二条消息就已开始执行,期望单条DIRECT路由中的消息可以按顺序逐个处理。

代码实现

@Override
public void configure() throws Exception {
    from("timer:cTimer?period=1&repeatCount=2")
            .process(e -> {
                final Integer oldBody = timerMessageId.getAndIncrement();
                log.info("Send message: {}", oldBody);
                e.getIn().setHeader(KEY, oldBody);
                final List<Integer> list = IntStream.range(1, 10).boxed().collect(Collectors.toList());
                e.getIn().setBody(list);
            })
            .to("direct:direct-order-file-process-route:" + psId);

    from(direct("direct-order-file-process-route:" + psId))
            .routeId("routeId:" + psId)
            .process(exchange -> {
                log.info("Started process message: {}", exchange.getIn().getHeader(KEY));
            })
            .split()
            .body()
            .aggregate(header(KEY), this::aggregatorExchanges)
            .completionSize(5)
            .completionTimeout(1000L)
            .process(exchange -> {
                Thread.sleep(100);
                log.info("DONE: {} / {}", exchange.getIn().getHeader(KEY), exchange.getIn().getBody());
            })
            .end();
}

private Exchange aggregatorExchanges(Exchange oldExchange, Exchange newExchange) {
    List<Integer> drinks;
    if (oldExchange == null) {
        drinks = new ArrayList<>();
    } else {
        drinks = (List<Integer>) oldExchange.getIn().getBody();
    }
    drinks.add((Integer) newExchange.getIn().getBody());
    newExchange.getIn().setBody(drinks);
    return newExchange;
}

运行日志对比

实际运行日志

Send message: 1
Started process message: 1
Send message: 2
Started process message: 2
DONE: 1 / [1, 2, 3, 4, 5]
DONE: 2 / [1, 2, 3, 4, 5]
DONE: 1 / [6, 7, 8, 9]
DONE: 2 / [6, 7, 8, 9]

期望运行日志

Send message: 1
Started process message: 1
DONE: 1 / [1, 2, 3, 4, 5]
DONE: 1 / [6, 7, 8, 9]

Send message: 2
Started process message: 2
DONE: 2 / [1, 2, 3, 4, 5]
DONE: 2 / [6, 7, 8, 9]

问题原因

当前问题与aggregate组件相关,未使用aggregate时direct运行逻辑符合预期。换言之,direct与aggregate未同步(二者运行在不同线程中),direct不会等待聚合操作完成。


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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.10.07 06:06:03